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): rename LegacyMode to PerPartitionMode
The "legacy" naming implied a deprecated path, but this is all new,
unreleased code with no prior path. Rename the scan mode and scrub
"legacy" wording from comments to describe what the mode actually does.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
  • Loading branch information
timsaucer and claude committed Jun 11, 2026
commit 4b997345099a39025acf07269455109628370664
7 changes: 3 additions & 4 deletions spark/bridge/src/scan.rs
Original file line number Diff line number Diff line change
Expand Up @@ -38,9 +38,8 @@
//! succeeds when every operator in that partition's pipeline supports
//! repeated `execute()` — stateless scans do, `RepartitionExec`
//! pipelines do not;
//! - [`execute_stream`] — the whole plan as one stream (legacy
//! per-partition payload mode, where the provider itself is the task's
//! slice);
//! - [`execute_stream`] — the whole plan as one stream (per-partition
//! mode, where the provider itself is the task's slice);
//! - [`close_scan`] — drop the plan. The single unsafe interleaving is
//! closing a handle that still has an in-flight call; the Java consumer
//! (the shared-scan cache) prevents it with a refcount covering every
Expand Down Expand Up @@ -282,7 +281,7 @@ pub fn execute_stream_partition(
})
}

/// Whole-plan stream for legacy per-partition payload mode (the provider
/// Whole-plan stream for per-partition mode (the provider
/// itself is the task's slice, so all plan partitions merge into one reader).
pub fn execute_stream(env: &mut JNIEnv, handle: jlong, ffi_stream_addr: jlong) {
try_unwrap_or_throw(env, (), |_env| -> JniResult<()> {
Expand Down
2 changes: 1 addition & 1 deletion spark/src/main/java/io/datafusion/spark/ScanBackend.java
Original file line number Diff line number Diff line change
Expand Up @@ -70,7 +70,7 @@ long createScan(

/**
* Stream the WHOLE plan (all partitions coalesced) into the caller-allocated {@code
* FFI_ArrowArrayStream} at {@code ffiStreamAddr}. Used by legacy per-partition payload mode.
* FFI_ArrowArrayStream} at {@code ffiStreamAddr}. Used by per-partition mode.
*/
void executeStream(long scanHandle, long ffiStreamAddr);

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -26,7 +26,7 @@ import org.apache.spark.sql.connector.read.{Batch, InputPartition, PartitionRead

/**
* Spark `Batch` for a DataFusion-backed scan. Driver-side partition planning:
* - [[LegacyMode]]: one task per `PartitionInfo` (resolved by [[DatafusionScanBuilder]]); when
* - [[PerPartitionMode]]: one task per `PartitionInfo` (resolved by [[DatafusionScanBuilder]]); when
* the bridge reported a partitioning and every entry carries key values, tasks implement
* `HasPartitionKey` so Spark can actually use the `KeyGroupedPartitioning`.
* - [[SharedScanMode]]: one task per DataFusion plan partition index.
Expand All @@ -38,7 +38,7 @@ class DatafusionBatch(val scan: DatafusionScan) extends Batch {
val filterBytes: Array[Array[Byte]] = scan.pushedPredicateBytes

scan.mode match {
case LegacyMode(partitions, reported) =>
case PerPartitionMode(partitions, reported) =>
val keyed = DatafusionBatch.validateKeyedState(scan.factoryFqcn, partitions, reported)
partitions.iterator.map { p =>
val base = DatafusionInputPartition(
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -26,7 +26,7 @@ import org.apache.spark.sql.types.StructType
import org.apache.spark.sql.vectorized.ColumnarBatch

/**
* Per-task columnar reader for the per-partition payload (legacy) path. Lifecycle:
* Per-task columnar reader for the per-partition path. Lifecycle:
*
* 1. Reflectively instantiate the bridge's `BridgeProviderFactory` (no-arg) and take its
* [[ScanBackend]].
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -29,7 +29,7 @@ import org.apache.spark.sql.connector.read.{HasPartitionKey, InputPartition}
sealed trait DatafusionPartition extends InputPartition

/**
* Per-task payload for the per-partition payload (legacy) read path.
* Per-task payload for the per-partition read path.
*
* - `factoryFqcn`: fully-qualified class name of the bridge's `BridgeProviderFactory`. The
* executor reflectively instantiates this and calls
Expand Down Expand Up @@ -59,7 +59,7 @@ final case class DatafusionInputPartition(
}

/**
* Legacy-path payload that additionally carries this partition's key values, precomputed
* Per-partition payload that additionally carries this partition's key values, precomputed
* driver-side into an [[InternalRow]]. Emitted by [[DatafusionBatch]] when the bridge reported a
* partitioning AND every `PartitionInfo` carries `partitionKeyValues` — implementing
* [[HasPartitionKey]] is what makes the reported `KeyGroupedPartitioning` visible to Spark 3.3+
Expand Down
8 changes: 4 additions & 4 deletions spark/src/main/scala/io/datafusion/spark/DatafusionScan.scala
Original file line number Diff line number Diff line change
Expand Up @@ -39,7 +39,7 @@ sealed trait DatafusionScanMode extends Serializable
* from that entry's `partitionBytes`. `reported` is the bridge's optional partitioning
* declaration (may be null).
*/
final case class LegacyMode(
final case class PerPartitionMode(
partitions: Array[PartitionInfo],
reported: ReportedPartitioning
) extends DatafusionScanMode
Expand All @@ -62,7 +62,7 @@ final case class SharedScanMode(
* executor applies natively via `ScanBackend.createScan`, and the driver-resolved
* [[DatafusionScanMode]].
*
* Legacy mode with a bridge-declared [[ReportedPartitioning]] surfaces `KeyGroupedPartitioning`
* Per-partition mode with a bridge-declared [[ReportedPartitioning]] surfaces `KeyGroupedPartitioning`
* via `SupportsReportPartitioning`; note Spark 3.3+ only consumes it when the input partitions
* also implement `HasPartitionKey` (see [[DatafusionBatch]]). Shared-scan mode always reports
* `UnknownPartitioning` — DataFusion-native partitions carry no key contract.
Expand All @@ -82,7 +82,7 @@ class DatafusionScan(

override def description(): String = {
val modeDesc = mode match {
case LegacyMode(partitions, reported) =>
case PerPartitionMode(partitions, reported) =>
s"mode=per-partition, partitions=${partitions.length}," +
s" reportedPartitioning=${if (reported == null) "unknown" else "key-grouped"}"
case SharedScanMode(scanId, n, _, _) =>
Expand All @@ -95,7 +95,7 @@ class DatafusionScan(
override def toBatch: Batch = new DatafusionBatch(this)

override def outputPartitioning(): Partitioning = mode match {
case LegacyMode(partitions, reported) =>
case PerPartitionMode(partitions, reported) =>
if (reported == null) new UnknownPartitioning(partitions.length)
else new KeyGroupedPartitioning(reported.keys().toArray, partitions.length)
case SharedScanMode(_, numPartitions, _, _) =>
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -87,7 +87,7 @@ class DatafusionScanBuilder(
val factory = instantiateFactory(factoryFqcn)
val mode: DatafusionScanMode =
if (factory.sharedScan(optionsBytes)) buildSharedScanMode()
else buildLegacyMode(factory)
else buildPerPartitionMode(factory)
new DatafusionScan(
factoryFqcn,
optionsBytes,
Expand All @@ -99,15 +99,15 @@ class DatafusionScanBuilder(
)
}

private def buildLegacyMode(factory: BridgeProviderFactory): LegacyMode = {
private def buildPerPartitionMode(factory: BridgeProviderFactory): PerPartitionMode = {
val partitions: Array[PartitionInfo] =
factory.listPartitions(optionsBytes, pushedBytes)
if (partitions == null || partitions.isEmpty) {
throw new IllegalStateException(
s"BridgeProviderFactory '$factoryFqcn' returned no partitions to scan"
)
}
LegacyMode(partitions, factory.reportPartitioning(optionsBytes))
PerPartitionMode(partitions, factory.reportPartitioning(optionsBytes))
}

/**
Expand Down