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
docs: scrub dual-path language after FFI removal
Doc comments still described two binding styles, a connector cdylib,
and a separate widening library — none of which exist. Module docs,
javadoc, error messages, and READMEs now describe the single path:
every bridge cdylib is an export_bridge! expansion over the
datafusion-spark-bridge SDK, widening included.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
  • Loading branch information
timsaucer and claude committed Jun 11, 2026
commit 13bca92dc1aa8398b892f148f4a2953fd4a89c66
12 changes: 6 additions & 6 deletions examples/python/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -110,16 +110,16 @@ root
+---+-----+
```

Filter row count drops from 4 → 2 because the predicate is pushed across the
FFI boundary as a `LogicalExprNode` proto and applied inside DataFusion before
Arrow batches cross back to Spark.
Filter row count drops from 4 → 2 because the predicate is pushed into the
bridge cdylib as a `LogicalExprNode` proto and applied inside DataFusion
before Arrow batches cross back to Spark.

## Notes

- `master("local[2]")` keeps driver + executor in one JVM so the example
cdylib loads once. Cluster mode would need the cdylib pre-staged on every
worker (the widening lib is bundled in `datafusion-java-spark`; only the
per-bridge example lib is not).
cdylib loads once. In cluster mode nothing extra is needed: the bridge
cdylib travels inside the examples jar and `NativeLibraryLoader` extracts
it on every worker.
- `extraClassPath` (not `--packages` / `userClassPathFirst`) is used because
the Spark distro ships Arrow 12, flatbuffers 1.12, and protobuf 2.5, all
of which we need to override; userClassPathFirst splits Netty across two
Expand Down
9 changes: 5 additions & 4 deletions native-common/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -15,10 +15,11 @@
// specific language governing permissions and limitations
// under the License.

//! JNI plumbing shared by this workspace's cdylibs (`datafusion-jni` and the
//! Spark connector helper): the error-to-Java-exception mapping, the
//! per-cdylib Tokio runtime singleton, and the async-stream-to-
//! `FFI_ArrowArrayStream` bridge.
//! JNI plumbing shared by this workspace's native crates (`datafusion-jni`
//! and `datafusion-spark-bridge`, and through the latter every bridge
//! cdylib): the error-to-Java-exception mapping, the per-cdylib Tokio
//! runtime singleton, and the async-stream-to-`FFI_ArrowArrayStream`
//! bridge.
//!
//! Each cdylib statically links its own copy of this rlib, so [`runtime`] is
//! a per-cdylib singleton -- exactly the behaviour each crate had when this
Expand Down
4 changes: 2 additions & 2 deletions native/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -82,8 +82,8 @@ pub(crate) fn jvm() -> &'static JavaVM {

pub(crate) fn runtime() -> &'static Runtime {
// The singleton itself lives in datafusion-jni-common (shared with the
// Spark helper cdylib; each cdylib statically links its own copy, so the
// runtime stays per-library). The init hook eagerly installs the
// datafusion-spark-bridge SDK; each cdylib statically links its own
// copy, so the runtime stays per-library). The init hook eagerly installs the
// runtime-metrics accumulator (no-op when the `runtime-metrics` Cargo
// feature is off). Initialising here -- not lazily on the first
// `runtimeStats()` call -- means the RuntimeMonitor's sampling baseline
Expand Down
3 changes: 2 additions & 1 deletion spark/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -380,7 +380,8 @@ Shared-scan operational details:
- **Schema inference** — your provider's Arrow schema, widened, becomes the
Spark schema. Driver-side, one probe build with empty `partitionBytes`.
- **Type widening** — Spark's columnar readers reject several Arrow types
DataFusion happily produces. The connector cdylib transparently casts
DataFusion happily produces. The SDK (inside your bridge's cdylib)
transparently casts
unsigned ints → wider signed, `Float16` → `Float32`, `Time*` → wider ints,
any-unit/tz `Timestamp` → microsecond, recursively through
`List`/`LargeList`/`FixedSizeList` (see
Expand Down
2 changes: 1 addition & 1 deletion spark/bridge/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -75,7 +75,7 @@ pub(crate) fn runtime_handle() -> &'static Handle {
datafusion_jni_common::runtime().handle()
}

/// Generate the JNI entry points for a static bridge cdylib.
/// Generate the JNI entry points for a bridge cdylib.
///
/// `jni_class` is the **underscore-mangled** binary name of the Java class
/// declaring the matching `native` methods: dots become underscores
Expand Down
12 changes: 6 additions & 6 deletions spark/bridge/src/scan.rs
Original file line number Diff line number Diff line change
Expand Up @@ -15,13 +15,13 @@
// specific language governing permissions and limitations
// under the License.

//! Planning and execution of a Spark scan, provider-source-agnostic.
//! Planning and execution of a Spark scan.
//!
//! Every function here is the body of one JNI entry point; the caller (the
//! generic FFI cdylib, or a static bridge's `export_bridge!` expansion)
//! supplies only how the provider is obtained, as a `make` closure. The
//! provider is wrapped in a [`WideningTableProvider`] here, so both binding
//! styles get identical Spark-compatible Arrow types.
//! Every function here is the body of one JNI entry point generated by a
//! bridge's `export_bridge!` expansion, which supplies only how the provider
//! is obtained, as a `make` closure. The provider is wrapped in a
//! [`WideningTableProvider`] here, so every bridge gets identical
//! Spark-compatible Arrow types.
//!
//! [`create_scan`] registers the widened provider on a private
//! `SessionContext` built from the caller-pinned config, applies the pruned
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -35,9 +35,10 @@ import org.apache.spark.sql.types._
* Null. No unsigned-int or Time accessor exists; we surface a clear error at schema discovery
* for those — the alternative is silent corruption.
*
* The widening cdylib (connector-core/native/) inserts a `WideningTableProvider` upstream of the
* Spark reader that casts unsupported types kernel-side (UInt*→signed wider, Float16→Float32,
* non-µs Timestamp→µs Timestamp, Time→Int) so Spark only ever sees compatible Arrow types.
* The widening layer (datafusion-spark-bridge, compiled into every bridge cdylib) inserts a
* `WideningTableProvider` upstream of the Spark reader that casts unsupported types kernel-side
* (UInt*→signed wider, Float16→Float32, non-µs Timestamp→µs Timestamp, Time→Int) so Spark only
* ever sees compatible Arrow types.
*/
object ArrowToSparkSchema {

Expand Down Expand Up @@ -67,7 +68,7 @@ object ArrowToSparkSchema {
unsupported(
f,
s"unsigned integer UInt$bits (Spark ArrowColumnVector has no unsigned accessor; " +
"widening cdylib casts these before Spark sees them — this branch indicates the " +
"widening layer casts these before Spark sees them — this branch indicates the " +
"WideningTableProvider was bypassed)"
)
case (bits, signed) => unsupported(f, s"Int(bits=$bits, signed=$signed)")
Expand All @@ -76,7 +77,7 @@ object ArrowToSparkSchema {
case t: ArrowType.FloatingPoint =>
t.getPrecision match {
case FloatingPointPrecision.HALF =>
unsupported(f, "Float16 (widening cdylib must cast to Float32 before Spark)")
unsupported(f, "Float16 (widening layer must cast to Float32 before Spark)")
case FloatingPointPrecision.SINGLE => FloatType
case FloatingPointPrecision.DOUBLE => DoubleType
case other => unsupported(f, s"FloatingPoint($other)")
Expand Down