Skip to content
Draft
Show file tree
Hide file tree
Changes from 1 commit
Commits
Show all changes
22 commits
Select commit Hold shift + click to select a range
d1c4a9b
WIP using rerun specific table options
timsaucer Jun 10, 2026
baa9ed3
WIP on making table providers spark data sources with a working example
timsaucer Jun 10, 2026
e7b5848
feat(spark): add Spark DataSource V2 connector + pyspark FFI demo
timsaucer Jun 10, 2026
7ee7c9d
refactor(build): consolidate Rust crates into a Cargo workspace
timsaucer Jun 10, 2026
8db9d4a
feat(examples): pass user options through FFI table provider demo
timsaucer Jun 10, 2026
e8c70a9
update examples to build after last commit
timsaucer Jun 10, 2026
088474d
feat(spark): per-partition payload + preferred locations in FFI factory
timsaucer Jun 10, 2026
f7d3972
docs(examples): update SPARK_INTEGRATION for PartitionInfo + per-slic…
timsaucer Jun 10, 2026
daa3ba5
feat(spark): SupportsReportPartitioning via optional reportPartitioni…
timsaucer Jun 10, 2026
c926d3d
feat(spark): shared-scan mode with per-executor provider cache
timsaucer Jun 11, 2026
1cffd93
refactor(spark): fold scan planning/execution into connector cdylib
timsaucer Jun 11, 2026
1f73a6f
docs: rewrite examples README and move Spark guide into spark/
timsaucer Jun 11, 2026
e9f3f61
feat(spark): add datafusion-spark-bridge SDK for static bridges
timsaucer Jun 11, 2026
dc909ce
feat(spark): ScanBackend dispatch + one-method-minimum factory
timsaucer Jun 11, 2026
45c9613
feat(spark): reusable native loader + bridge packaging recipe
timsaucer Jun 11, 2026
cc35958
feat(spark): bridge scaffold generator
timsaucer Jun 11, 2026
9e47f0e
refactor(spark)!: remove the FFI provider path
timsaucer Jun 11, 2026
13bca92
docs: scrub dual-path language after FFI removal
timsaucer Jun 11, 2026
a1815a8
refactor(spark): rename optionsProtoBytes, fix stale createProvider docs
timsaucer Jun 11, 2026
4b99734
refactor(spark): rename LegacyMode to PerPartitionMode
timsaucer Jun 11, 2026
827068c
refactor(spark): move bridge scaffold from dev/ to spark/scaffold/
timsaucer Jun 12, 2026
0cc474a
add support for fixed sized list widening
timsaucer Jun 12, 2026
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Prev Previous commit
Next Next commit
refactor(spark)!: remove the FFI provider path
Every known bridge owns its provider's Rust source, so the
FFI_TableProvider handover was speculative generality with real
costs: a second cdylib bundled in the connector jar, datafusion-ffi
ABI lockstep between artifacts, and pointer-ownership rules across
JNI. Static export_bridge! bridges are now the only path.

- delete the datafusion-spark-helper cdylib, FfiHelperNative,
  FfiScanBackend, and the bridge SDK's ffi module + feature;
  datafusion-ffi leaves the dependency tree entirely
- the connector jar goes pure JVM: no cargo prerequisite, no
  per-platform builds; native code ships inside each bridge's jar
- convert examples/native to an export_bridge! bridge
  (datafusion-java-example-bridge) — the committed, runnable
  concrete-provider example; demo renamed to bridge_demo.py

BREAKING CHANGE: FfiProviderFactory is renamed BridgeProviderFactory;
createProvider is gone and scanBackend() is the single required
method. Re-adding an FFI path later is mechanical: the ScanBackend
seam and scan.rs's provider-source closure are unchanged, and the
deleted code sits intact in this branch's history.

Verified: cargo + mvn suites, spotless/RAT verify, pyspark demo (both
scan modes), and a regenerated scaffold built + smoke-tested.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
  • Loading branch information
timsaucer and claude committed Jun 11, 2026
commit 9e47f0e192e2ad345124bd859025820742df36f0
306 changes: 38 additions & 268 deletions Cargo.lock

Large diffs are not rendered by default.

2 changes: 0 additions & 2 deletions Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -22,7 +22,6 @@ members = [
"native-common",
"examples/native",
"spark/bridge",
"spark/native",
]

# Every dependency used by any workspace member is declared here so version
Expand All @@ -34,7 +33,6 @@ members = [
arrow = { version = "58", features = ["ffi"] }
async-trait = "0.1"
datafusion = { version = "53.1.0" }
datafusion-ffi = "53.1.0"
datafusion-proto = "53.1.0"
datafusion-substrait = "53.1.0"
futures = "0.3"
Expand Down
4 changes: 1 addition & 3 deletions dev/bridge-template/native/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -13,11 +13,9 @@ name = "__LIB__"
crate-type = ["cdylib"]

[dependencies]
# default-features = false drops the datafusion-ffi import path — a static
# bridge never crosses an FFI_TableProvider boundary.
# TODO: replace the path with a git or crates.io dependency once you build
# outside a local datafusion-java checkout.
datafusion-spark-bridge = { path = "__BRIDGE_SDK_PATH__", default-features = false }
datafusion-spark-bridge = { path = "__BRIDGE_SDK_PATH__" }

[profile.release]
strip = "debuginfo"
Original file line number Diff line number Diff line change
@@ -1,13 +1,13 @@
package __PKG__;

import io.datafusion.spark.FfiProviderFactory;
import io.datafusion.spark.BridgeProviderFactory;
import io.datafusion.spark.ScanBackend;

/**
* The bridge's contract with the Spark connector. This is a STATIC bridge — the provider is built
* inside this bridge's own cdylib — so the only required override is {@link #scanBackend()}.
* The bridge's contract with the Spark connector: the provider is built inside this bridge's own
* cdylib, and {@link #scanBackend()} is the only required method.
*
* <p>Useful optional overrides (see their javadoc on {@link FfiProviderFactory}):
* <p>Useful optional overrides (see their javadoc on {@link BridgeProviderFactory}):
*
* <ul>
* <li>{@code encodeOptions} — only if you have your own options schema; the default ships the
Expand All @@ -19,7 +19,7 @@
* task per DataFusion output partition. Mind the determinism contract.
* </ul>
*/
public final class __PREFIX__ProviderFactory implements FfiProviderFactory {
public final class __PREFIX__ProviderFactory implements BridgeProviderFactory {

@Override
public ScanBackend scanBackend() {
Expand Down
9 changes: 5 additions & 4 deletions docs/source/contributor-guide/development.md
Original file line number Diff line number Diff line change
Expand Up @@ -75,7 +75,7 @@ disk space.
The repository is a multi-module Maven build:

- `Cargo.toml` — Rust workspace root declaring the three crate members
(`native`, `examples/native`, `spark/native`) and `[workspace.dependencies]`
(`native`, `native-common`, `examples/native`, `spark/bridge`) and `[workspace.dependencies]`
that pin shared versions in one place. Cargo writes artifacts to
`rust-target/` (overridden in `.cargo/config.toml`) so `mvn clean` at the
repo root does not nuke the Rust build cache.
Expand All @@ -84,12 +84,13 @@ The repository is a multi-module Maven build:
- `core/` — `datafusion-java` library module (Java sources, tests, and
generated protobuf classes).
- `spark/` — `datafusion-java-spark` Spark DataSource V2 connector
(Scala + Java) and its `spark/native/` widening cdylib crate.
(Scala + Java, pure JVM) and its `spark/bridge/` Rust SDK crate
(`datafusion-spark-bridge`: widening, scan machinery, `export_bridge!`).
- `examples/` — `datafusion-java-examples` module containing runnable
examples that depend on the library; built alongside the library so they
cannot fall out of sync with the API. Includes `examples/native/`, a
small FFI table-provider cdylib used by the Spark connector demo
(`ExampleFfiProviderFactory` + the pyspark script under
small `export_bridge!` cdylib used by the Spark connector demo
(`ExampleBridgeProviderFactory` + the pyspark script under
`examples/python/`).
- `native/` — `datafusion-jni` Rust crate (JNI + Arrow C Data Interface).
- `proto/` — Protobuf definitions shared between Java and Rust.
Expand Down
19 changes: 10 additions & 9 deletions examples/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -64,17 +64,18 @@ the result rows. Swap `SqlQueryExample` for any class in the table below.
## The Spark connector example

One example is not a standalone `main`:
`ExampleFfiProviderFactory` implements the Spark connector's
`FfiProviderFactory` interface over a tiny Rust-built in-memory table (the
cdylib under [`native/`](native/)). It exists to be loaded *by Spark* — the
runnable end-to-end version is the PySpark demo under
[`python/`](python/), and the guide to building your own connector is
`ExampleBridgeProviderFactory` implements the Spark connector's
`BridgeProviderFactory` interface over a tiny in-memory table built inside
the example bridge cdylib (the `export_bridge!` crate under
[`native/`](native/)). It exists to be loaded *by Spark* — the runnable
end-to-end version is the PySpark demo under [`python/`](python/), and the
guide to building your own connector is
[`../spark/README.md`](../spark/README.md).

To build its cdylib (workspace member, buildable from anywhere in the tree):

```bash
cargo build -p datafusion-java-ffi-example --release
cargo build -p datafusion-java-example-bridge --release
```

Building the examples jar then bundles the cdylib inside it (under
Expand All @@ -83,7 +84,7 @@ there at runtime via the connector's `NativeLibraryLoader` — the same
packaging recipe a real bridge uses (see "Packaging your bridge" in
[`../spark/README.md`](../spark/README.md)). To run against an unpackaged
local build instead, pass
`-Dexample.ffi.lib.path=/abs/path/to/libdatafusion_java_ffi_example.{so,dylib}`.
`-Dexample.bridge.lib.path=/abs/path/to/libdatafusion_example_bridge.{so,dylib}`.

## Troubleshooting

Expand All @@ -94,5 +95,5 @@ local build instead, pass
built in a different profile than Maven expects. Re-run build step 1 and
keep `-Ddatafusion.native.profile=release` consistent between the cargo
profile (`--release`) and the Maven flag.
- **`UnsatisfiedLinkError ... datafusion_java_ffi_example`** — only the FFI
example's cdylib is missing; see "The Spark connector example" above.
- **`UnsatisfiedLinkError ... datafusion_example_bridge`** — only the example
bridge cdylib is missing; see "The Spark connector example" above.
12 changes: 7 additions & 5 deletions examples/native/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -9,19 +9,21 @@
# http://www.apache.org/licenses/LICENSE-2.0

[package]
name = "datafusion-java-ffi-example"
name = "datafusion-java-example-bridge"
version = "0.1.0"
edition = "2021"
publish = false

[lib]
# Built as a cdylib so the JVM-side example can System.load() the artifact.
# `rlib` lets us add Rust-level unit tests if needed.
name = "datafusion_example_bridge"
# Built as a cdylib so the JVM loads it via NativeLibraryLoader; `rlib` keeps
# the Rust-level unit tests (options decoding, partition layout) runnable.
crate-type = ["cdylib", "rlib"]

[dependencies]
arrow = { workspace = true }
datafusion = { workspace = true }
datafusion-ffi = { workspace = true }
jni = { workspace = true }
datafusion-spark-bridge = { path = "../../spark/bridge" }

[dev-dependencies]
tokio = { workspace = true }
133 changes: 30 additions & 103 deletions examples/native/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -15,20 +15,19 @@
// specific language governing permissions and limitations
// under the License.

//! Example cdylib that produces a small DataFusion `MemTable` wrapped as an
//! `FFI_TableProvider`, returned to the JVM as a `jlong` (the raw boxed
//! pointer). The Spark connector consumes the pointer via
//! `FfiHelperNative.createScan` / `providerSchemaIpc`, which widen the
//! provider and plan/execute the scan inside the connector cdylib.
//! Example bridge cdylib: a small DataFusion `MemTable` exposed to Spark
//! through the `datafusion-spark-bridge` SDK. `export_bridge!` generates the
//! whole JNI surface for `org.apache.datafusion.examples.ExampleBridgeNative`;
//! this crate only decodes the options blob and builds the provider.
//!
//! The same pattern is what domain bridges (HDF5, custom Iceberg, in-house formats) use
//! to expose their TableProviders to Spark via the connector-core DataSource
//! V2 plumbing.
//! The same pattern is what domain bridges (HDF5, custom Iceberg, in-house
//! formats) use to expose their TableProviders to Spark via the connector's
//! DataSource V2 plumbing.
//!
//! ## Options wire format
//!
//! `createMemTableProvider` accepts an opaque `byte[]` that the JVM-side
//! `ExampleFfiProviderFactory.encodeOptions` produces. Layout (little-endian):
//! The provider builder accepts an opaque `byte[]` that the JVM-side
//! `ExampleBridgeProviderFactory.encodeOptions` produces. Layout (little-endian):
//!
//! ```text
//! [u32 name_prefix_len][name_prefix UTF-8 bytes][u32 num_rows][u32 num_batches]
Expand All @@ -38,45 +37,19 @@
//! Empty/`null` bytes decode as all defaults: `name_prefix="row"`, `num_rows=4`,
//! `num_batches=1`, `num_partitions=1`, `shared_scan=false`. The trailing
//! fields are optional so blobs from older encoders keep decoding. The
//! `shared_scan` flag is consumed JVM-side (`ExampleFfiProviderFactory.sharedScan`);
//! `shared_scan` flag is consumed JVM-side (`ExampleBridgeProviderFactory.sharedScan`);
//! this decoder carries it only so one blob format serves both sides. Real
//! bridges use a real proto schema here; this example hand-rolls the encoding
//! to keep the wire layer obvious.
//! bridges can use the connector's default `OptionsCodec` instead (decoded via
//! `datafusion_spark_bridge::options`); this example hand-rolls the encoding
//! to show a custom wire layer.

use std::sync::Arc;

use arrow::array::{Float64Array, Int64Array, RecordBatch, StringArray};
use arrow::datatypes::{DataType, Field, Schema as ArrowSchema};
use datafusion::catalog::TableProvider;
use datafusion::datasource::MemTable;
use datafusion::execution::TaskContextProvider;
use datafusion::prelude::SessionContext;
use datafusion_ffi::execution::FFI_TaskContextProvider;
use datafusion_ffi::table_provider::FFI_TableProvider;
use jni::objects::{JByteArray, JClass};
use jni::sys::jlong;
use jni::JNIEnv;
use tokio::runtime::{Handle, Runtime};

/// Tokio runtime that the FFI provider is anchored to. Shared across calls
/// for the lifetime of the cdylib so successive `createMemTableProvider`
/// invocations don't spawn fresh runtimes.
fn runtime() -> &'static Handle {
use std::sync::OnceLock;
static RT: OnceLock<Runtime> = OnceLock::new();
RT.get_or_init(|| Runtime::new().expect("tokio runtime init failed"))
.handle()
}

/// Host `SessionContext` used only to obtain a `TaskContextProvider` for
/// `FFI_TableProvider::new`. Static on purpose: the `FFI_TaskContextProvider`
/// holds a non-owning reference, so this context must outlive every provider
/// built from it. Nothing is ever registered on it.
fn host_session_context() -> &'static Arc<SessionContext> {
use std::sync::OnceLock;
static CTX: OnceLock<Arc<SessionContext>> = OnceLock::new();
CTX.get_or_init(|| Arc::new(SessionContext::new()))
}
use datafusion_spark_bridge::{export_bridge, BridgeContext, JniResult};

#[derive(Debug)]
struct Options {
Expand Down Expand Up @@ -186,70 +159,22 @@ fn build_mem_table(
Ok(Arc::new(MemTable::try_new(schema, partitions)?))
}

/// JNI entry point: decode the options blob, build a `MemTable` accordingly,
/// wrap it in an `FFI_TableProvider`, return the raw boxed pointer as a `jlong`.
/// Ownership of the boxed FFI transfers to the caller — the matching
/// `Box::from_raw` is performed by the consumer (the Spark connector's
/// `FfiHelperNative.createScan` / `providerSchemaIpc`).
#[no_mangle]
pub extern "system" fn Java_org_apache_datafusion_examples_FfiTableProviderExampleNative_createMemTableProvider<
'local,
>(
mut env: JNIEnv<'local>,
_class: JClass<'local>,
options_bytes: JByteArray<'local>,
) -> jlong {
let result: Result<jlong, Box<dyn std::error::Error + Send + Sync>> = (|| {
let bytes: Vec<u8> = if options_bytes.is_null() {
Vec::new()
} else {
env.convert_byte_array(&options_bytes)
.map_err(|e| format!("failed to read options byte[] from JVM: {e}"))?
};
let opts = decode_options(&bytes)?;

let mem_table = build_mem_table(&opts)?;
let provider: Arc<dyn TableProvider> = mem_table;

let ctx_provider: Arc<dyn TaskContextProvider> =
Arc::clone(host_session_context()) as Arc<dyn TaskContextProvider>;
let ffi_task_ctx = FFI_TaskContextProvider::from(&ctx_provider);
let ffi = FFI_TableProvider::new(
provider,
/*can_support_pushdown_filters=*/ true,
Some(runtime().clone()),
ffi_task_ctx,
/*logical_codec=*/ None,
);
Ok(Box::into_raw(Box::new(ffi)) as jlong)
})();

match result {
Ok(ptr) => ptr,
Err(err) => {
let _ = env.throw_new("java/lang/RuntimeException", err.to_string());
0
}
}
/// Build the example provider for one scan: decode the options blob, build
/// the `MemTable` accordingly. `partition` is unused — the example reports a
/// single partition (or relies on shared-scan mode), so there is no per-task
/// payload to interpret.
fn build_provider(
_ctx: &BridgeContext,
options: &[u8],
_partition: &[u8],
) -> JniResult<Arc<dyn TableProvider>> {
let opts = decode_options(options)?;
Ok(build_mem_table(&opts)?)
}

/// Drop a previously-created FFI_TableProvider whose pointer was NOT handed
/// off to a consumer. Exposed for the error path — callers that pass the
/// pointer to `createScan` / `providerSchemaIpc` must NOT also call this;
/// ownership has already transferred.
#[no_mangle]
pub extern "system" fn Java_org_apache_datafusion_examples_FfiTableProviderExampleNative_dropProvider<
'local,
>(
_env: JNIEnv<'local>,
_class: JClass<'local>,
ffi_ptr: jlong,
) {
if ffi_ptr != 0 {
unsafe {
drop(Box::from_raw(ffi_ptr as *mut FFI_TableProvider));
}
}
export_bridge! {
jni_class: "org_apache_datafusion_examples_ExampleBridgeNative",
build_provider: build_provider,
}

#[cfg(test)]
Expand Down Expand Up @@ -316,6 +241,8 @@ mod tests {
let table = build_mem_table(&opts).unwrap();
// MemTable has no partition accessor; verify via scan output partitioning.
use datafusion::catalog::TableProvider;
use datafusion::prelude::SessionContext;
use tokio::runtime::Runtime;
let ctx = SessionContext::new();
let rt = Runtime::new().unwrap();
let plan = rt
Expand Down
16 changes: 8 additions & 8 deletions examples/pom.xml
Original file line number Diff line number Diff line change
Expand Up @@ -47,7 +47,7 @@ under the License.
<groupId>org.apache.datafusion</groupId>
<artifactId>datafusion-java</artifactId>
</dependency>
<!-- Provides io.datafusion.spark.FfiProviderFactory for ExampleFfiProviderFactory.
<!-- Provides io.datafusion.spark.BridgeProviderFactory for ExampleBridgeProviderFactory.
Compile-only: the example factory does not pull Spark itself onto the classpath;
the pyspark demo supplies Spark at runtime. -->
<dependency>
Expand Down Expand Up @@ -105,14 +105,14 @@ under the License.
<version>3.1.0</version>
<executions>
<execution>
<id>copy-ffi-example-cdylib</id>
<id>copy-example-bridge-cdylib</id>
<phase>process-classes</phase>
<goals><goal>run</goal></goals>
<configuration>
<target>
<property name="example.ffi.lib.source"
value="${maven.multiModuleProjectDirectory}/rust-target/${datafusion.native.profile}/${example.ffi.filename}"/>
<fail message="Example bridge cdylib not found at ${example.ffi.lib.source}. Run 'cargo build -p datafusion-java-ffi-example' (or '--release' for the release profile) before building the JAR.">
<fail message="Example bridge cdylib not found at ${example.ffi.lib.source}. Run 'cargo build -p datafusion-java-example-bridge' (or '--release' for the release profile) before building the JAR.">
<condition><not><available file="${example.ffi.lib.source}"/></not></condition>
</fail>
<mkdir dir="${project.build.outputDirectory}/org/apache/datafusion/examples/${example.ffi.os}/${example.ffi.arch}"/>
Expand All @@ -139,7 +139,7 @@ under the License.
<properties>
<example.ffi.os>linux</example.ffi.os>
<example.ffi.arch>x86_64</example.ffi.arch>
<example.ffi.filename>libdatafusion_java_ffi_example.so</example.ffi.filename>
<example.ffi.filename>libdatafusion_example_bridge.so</example.ffi.filename>
</properties>
</profile>
<profile>
Expand All @@ -150,7 +150,7 @@ under the License.
<properties>
<example.ffi.os>linux</example.ffi.os>
<example.ffi.arch>x86_64</example.ffi.arch>
<example.ffi.filename>libdatafusion_java_ffi_example.so</example.ffi.filename>
<example.ffi.filename>libdatafusion_example_bridge.so</example.ffi.filename>
</properties>
</profile>
<profile>
Expand All @@ -161,7 +161,7 @@ under the License.
<properties>
<example.ffi.os>linux</example.ffi.os>
<example.ffi.arch>aarch64</example.ffi.arch>
<example.ffi.filename>libdatafusion_java_ffi_example.so</example.ffi.filename>
<example.ffi.filename>libdatafusion_example_bridge.so</example.ffi.filename>
</properties>
</profile>
<profile>
Expand All @@ -172,7 +172,7 @@ under the License.
<properties>
<example.ffi.os>darwin</example.ffi.os>
<example.ffi.arch>x86_64</example.ffi.arch>
<example.ffi.filename>libdatafusion_java_ffi_example.dylib</example.ffi.filename>
<example.ffi.filename>libdatafusion_example_bridge.dylib</example.ffi.filename>
</properties>
</profile>
<profile>
Expand All @@ -183,7 +183,7 @@ under the License.
<properties>
<example.ffi.os>darwin</example.ffi.os>
<example.ffi.arch>aarch64</example.ffi.arch>
<example.ffi.filename>libdatafusion_java_ffi_example.dylib</example.ffi.filename>
<example.ffi.filename>libdatafusion_example_bridge.dylib</example.ffi.filename>
</properties>
</profile>
</profiles>
Expand Down
Loading