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 optionsProtoBytes, fix stale createProvider docs
The default options encoding is OptionsCodec key/value strings, not
protobuf, so optionsProtoBytes was misleading. Rename to optionsBytes
across main, test, and examples. filterProtoBytes is left as-is — those
genuinely are LogicalExprNode proto bytes.

Also fix five doc references to the removed createProvider method (now
ScanBackend.createScan): two broken {@link}s in PartitionInfo.java that
would fail -Xdoclint, two comments in DatafusionInputPartition.scala,
and the bridge_demo.py note (which also claimed a non-existent stdout
line — reworded to the native build_provider).

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
  • Loading branch information
timsaucer and claude committed Jun 11, 2026
commit a1815a8d97fdc026d38eaeabb89b9285793a51bd
4 changes: 2 additions & 2 deletions examples/python/bridge_demo.py
Original file line number Diff line number Diff line change
Expand Up @@ -222,8 +222,8 @@ def main() -> None:

# Note on cache scope: the executor cache is keyed by a per-query scanId,
# so sharing happens across the TASKS of one query (4 tasks above -> one
# provider build per executor JVM, observable via the factory's
# createProvider stdout line), not across separate actions. Each new
# provider build per executor JVM, in the bridge's native build_provider),
# not across separate actions. Each new
# action plans a new scan with a fresh scanId; its entry simply joins the
# cache until the idle TTL evicts it.
count_again = shared.count()
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -60,7 +60,7 @@
* </ul>
*
* <p>Real bridges (HDF5, custom Iceberg, in-house formats) use a protobuf schema for {@code
* optionsProtoBytes}; this example uses a hand-rolled length-prefixed binary format to keep the
* optionsBytes}; this example uses a hand-rolled length-prefixed binary format to keep the
* wire layer obvious:
*
* <pre>
Expand Down Expand Up @@ -112,31 +112,31 @@ public byte[] encodeOptions(Map<String, String> sparkOptions) {
}

@Override
public PartitionInfo[] listPartitions(byte[] optionsProtoBytes) {
public PartitionInfo[] listPartitions(byte[] optionsBytes) {
// Single partition; the example MemTable is not actually sliced. A real bridge would
// populate `partitionBytes` per slice and `preferredLocations` with the hosts holding it.
return new PartitionInfo[] {new PartitionInfo("p0", new byte[0], new String[0])};
}

@Override
public PartitionInfo[] listPartitions(byte[] optionsProtoBytes, byte[][] filterProtoBytes) {
public PartitionInfo[] listPartitions(byte[] optionsBytes, byte[][] filterProtoBytes) {
// The example cannot prune its single partition, but a real bridge would inspect the
// pushed predicates here and drop partitions that cannot match.
System.out.println(
"ExampleBridgeProviderFactory.listPartitions received "
+ filterProtoBytes.length
+ " pushed filter(s)");
return listPartitions(optionsProtoBytes);
return listPartitions(optionsBytes);
}

@Override
public boolean sharedScan(byte[] optionsProtoBytes) {
public boolean sharedScan(byte[] optionsBytes) {
// The flag is the final byte of the options blob (present only when the encoder wrote the
// trailing fields). The bridge owns its wire format, so decoding it here is fair game.
return optionsProtoBytes != null
&& optionsProtoBytes.length >= 1
&& hasTrailingFields(optionsProtoBytes)
&& optionsProtoBytes[optionsProtoBytes.length - 1] == 1;
return optionsBytes != null
&& optionsBytes.length >= 1
&& hasTrailingFields(optionsBytes)
&& optionsBytes[optionsBytes.length - 1] == 1;
}

private static boolean hasTrailingFields(byte[] bytes) {
Expand Down
8 changes: 4 additions & 4 deletions spark/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -149,8 +149,8 @@ into more than one Spark task:

```java
@Override
public PartitionInfo[] listPartitions(byte[] optionsProtoBytes) {
MySlice[] slices = MyBridgeNative.listSlices(optionsProtoBytes);
public PartitionInfo[] listPartitions(byte[] optionsBytes) {
MySlice[] slices = MyBridgeNative.listSlices(optionsBytes);
PartitionInfo[] out = new PartitionInfo[slices.length];
for (int i = 0; i < slices.length; i++) {
out[i] = new PartitionInfo(slices[i].id(), slices[i].payload(), slices[i].hosts());
Expand Down Expand Up @@ -333,7 +333,7 @@ provider builds dominate. Opting in via

```java
@Override
public boolean sharedScan(byte[] optionsProtoBytes) { return true; }
public boolean sharedScan(byte[] optionsBytes) { return true; }
```

flips the mapping: the provider is built **once per executor JVM per query**
Expand All @@ -356,7 +356,7 @@ Choosing between the modes:

Shared-scan's price of admission is a **determinism contract**: the
provider's schema, partitioning, and per-partition contents must be a pure
function of `optionsProtoBytes`. Remote sources must pin a snapshot
function of `optionsBytes`. Remote sources must pin a snapshot
(version/timestamp) inside the options. The connector fails tasks when an
executor's partition count diverges from the driver's, but equal counts with
different contents are undetectable by construction. The provider's
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -79,7 +79,7 @@ default byte[] encodeOptions(Map<String, String> sparkOptions) {
* #sharedScan(byte[])}) before pointing it at anything large. Size guidance lives in {@code
* spark/README.md}.
*/
default PartitionInfo[] listPartitions(byte[] optionsProtoBytes) {
default PartitionInfo[] listPartitions(byte[] optionsBytes) {
return new PartitionInfo[] {new PartitionInfo("p0", new byte[0], new String[0])};
}

Expand All @@ -95,8 +95,8 @@ default PartitionInfo[] listPartitions(byte[] optionsProtoBytes) {
* conjunction of all pushed predicates. The default delegates to the filter-unaware overload (no
* pruning), which is always correct.
*/
default PartitionInfo[] listPartitions(byte[] optionsProtoBytes, byte[][] filterProtoBytes) {
return listPartitions(optionsProtoBytes);
default PartitionInfo[] listPartitions(byte[] optionsBytes, byte[][] filterProtoBytes) {
return listPartitions(optionsBytes);
}

/**
Expand All @@ -118,7 +118,7 @@ default PartitionInfo[] listPartitions(byte[] optionsProtoBytes, byte[][] filter
*
* <ul>
* <li>The provider's schema, partitioning, and per-partition row content are a pure function of
* {@code optionsProtoBytes}. Remote sources must pin a snapshot (version, timestamp) inside
* {@code optionsBytes}. Remote sources must pin a snapshot (version, timestamp) inside
* the options; data that compacts or moves between driver planning and executor execution
* otherwise yields wrong results that no runtime check can catch.
* <li>The provider's {@code ExecutionPlan} supports calling {@code execute(i)} more than once
Expand All @@ -129,7 +129,7 @@ default PartitionInfo[] listPartitions(byte[] optionsProtoBytes, byte[][] filter
* <p>The connector fails tasks with a clear error when the executor's partition count diverges
* from the driver's — but identical counts with different contents cannot be detected.
*/
default boolean sharedScan(byte[] optionsProtoBytes) {
default boolean sharedScan(byte[] optionsBytes) {
return false;
}

Expand All @@ -154,7 +154,7 @@ default boolean sharedScan(byte[] optionsProtoBytes) {
* KeyGroupedPartitioning} entirely. Storage-partitioned joins additionally require {@code
* spark.sql.sources.v2.bucketing.enabled=true}.
*/
default ReportedPartitioning reportPartitioning(byte[] optionsProtoBytes) {
default ReportedPartitioning reportPartitioning(byte[] optionsBytes) {
return null;
}
}
6 changes: 3 additions & 3 deletions spark/src/main/java/io/datafusion/spark/PartitionInfo.java
Original file line number Diff line number Diff line change
Expand Up @@ -22,7 +22,7 @@
/**
* Driver-side descriptor for a single partition produced by {@link
* BridgeProviderFactory#listPartitions(byte[])}. Carries the bridge-specific slice payload that the
* executor passes back into {@link BridgeProviderFactory#createProvider(byte[], byte[])}, plus
* executor passes back into {@link ScanBackend#createScan}, plus
* optional host hints for Spark's scheduler.
*
* <p>Fields:
Expand All @@ -32,8 +32,8 @@
* Surfaces in Spark UI, logs, and exception messages. Must be non-empty.
* <li>{@code partitionBytes} — opaque per-partition payload. Bridge encodes whatever the executor
* needs to materialise *this* slice (offsets, row ranges, sub-options, etc.). Combined with
* the global {@code optionsProtoBytes} in {@link BridgeProviderFactory#createProvider(byte[],
* byte[])}. Empty array = no per-partition state (single-partition table).
* the global {@code optionsBytes} in {@link ScanBackend#createScan}. Empty array = no
* per-partition state (single-partition table).
* <li>{@code preferredLocations} — hostnames where this partition's data lives. Returned from
* {@code InputPartition.preferredLocations()} so Spark can co-locate the task with the data.
* Empty array = no preference. Honoured subject to {@code spark.locality.wait}.
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -43,7 +43,7 @@ class DatafusionBatch(val scan: DatafusionScan) extends Batch {
partitions.iterator.map { p =>
val base = DatafusionInputPartition(
factoryFqcn = scan.factoryFqcn,
optionsProtoBytes = scan.optionsProtoBytes,
optionsBytes = scan.optionsBytes,
projectionColumnNames = projection,
filterProtoBytes = filterBytes,
partitionId = p.id,
Expand All @@ -63,7 +63,7 @@ class DatafusionBatch(val scan: DatafusionScan) extends Batch {
Array.tabulate[InputPartition](numPartitions) { i =>
DatafusionSharedScanPartition(
factoryFqcn = scan.factoryFqcn,
optionsProtoBytes = scan.optionsProtoBytes,
optionsBytes = scan.optionsBytes,
projectionColumnNames = projection,
filterProtoBytes = filterBytes,
scanId = scanId,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -49,7 +49,7 @@ class DatafusionColumnarPartitionReader(
private val scanHandle: Long =
try {
backend.createScan(
partition.optionsProtoBytes,
partition.optionsBytes,
partition.partitionBytes,
/* targetPartitions = */ -1,
/* batchSize = */ -1,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -32,22 +32,22 @@ sealed trait DatafusionPartition extends InputPartition
* Per-task payload for the per-partition payload (legacy) read path.
*
* - `factoryFqcn`: fully-qualified class name of the bridge's `BridgeProviderFactory`. The
* executor reflectively instantiates this and calls `createProvider(optionsProtoBytes,
* partitionBytes)`.
* - `optionsProtoBytes`: bridge-specific global connection options, encoded by the bridge.
* executor reflectively instantiates this and calls
* `scanBackend().createScan(optionsBytes, partitionBytes, …)`.
* - `optionsBytes`: bridge-specific global connection options, encoded by the bridge.
* Opaque to connector-core. Same bytes ride along on every partition.
* - `projectionColumnNames`: pruned column list (post-`pruneColumns`).
* - `filterProtoBytes`: V2 `Predicate` → DataFusion `LogicalExprNode` proto bytes; each one is
* applied natively via `ScanBackend.createScan`.
* - `partitionId`: stable identifier (e.g. a segment or file id) — surfaces in Spark UI/logs/errors.
* - `partitionBytes`: opaque per-partition payload from `PartitionInfo.partitionBytes`. Passed
* back into `createProvider` so the bridge materialises *this* slice.
* back into `ScanBackend.createScan` so the bridge materialises *this* slice.
* - `preferredLocs`: hostnames where this partition's data lives; returned from
* `preferredLocations()` so Spark schedules the task there subject to `spark.locality.wait`.
*/
final case class DatafusionInputPartition(
factoryFqcn: String,
optionsProtoBytes: Array[Byte],
optionsBytes: Array[Byte],
projectionColumnNames: Array[String],
filterProtoBytes: Array[Array[Byte]],
partitionId: String,
Expand Down Expand Up @@ -93,7 +93,7 @@ final case class DatafusionKeyedInputPartition(
*/
final case class DatafusionSharedScanPartition(
factoryFqcn: String,
optionsProtoBytes: Array[Byte],
optionsBytes: Array[Byte],
projectionColumnNames: Array[String],
filterProtoBytes: Array[Array[Byte]],
scanId: String,
Expand All @@ -107,7 +107,7 @@ final case class DatafusionSharedScanPartition(
SharedScanSpec(
scanId = scanId,
factoryFqcn = factoryFqcn,
optionsProtoBytes = optionsProtoBytes,
optionsBytes = optionsBytes,
projectionColumnNames = projectionColumnNames,
filterProtoBytes = filterProtoBytes,
pinnedConfig = pinnedConfig
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -69,7 +69,7 @@ final case class SharedScanMode(
*/
class DatafusionScan(
val factoryFqcn: String,
val optionsProtoBytes: Array[Byte],
val optionsBytes: Array[Byte],
val fullSchema: StructType,
val prunedSchema: StructType,
val pushedPredicates: Array[Predicate],
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -45,7 +45,7 @@ import org.apache.spark.sql.types.StructType
*/
class DatafusionScanBuilder(
factoryFqcn: String,
optionsProtoBytes: Array[Byte],
optionsBytes: Array[Byte],
fullSchema: StructType
) extends ScanBuilder
with SupportsPushDownV2Filters
Expand Down Expand Up @@ -86,11 +86,11 @@ class DatafusionScanBuilder(
override def build(): Scan = {
val factory = instantiateFactory(factoryFqcn)
val mode: DatafusionScanMode =
if (factory.sharedScan(optionsProtoBytes)) buildSharedScanMode()
if (factory.sharedScan(optionsBytes)) buildSharedScanMode()
else buildLegacyMode(factory)
new DatafusionScan(
factoryFqcn,
optionsProtoBytes,
optionsBytes,
fullSchema,
pruned,
pushed,
Expand All @@ -101,13 +101,13 @@ class DatafusionScanBuilder(

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

/**
Expand All @@ -125,7 +125,7 @@ class DatafusionScanBuilder(
val probeSpec = SharedScanSpec(
scanId = scanId,
factoryFqcn = factoryFqcn,
optionsProtoBytes = optionsProtoBytes,
optionsBytes = optionsBytes,
projectionColumnNames = pruned.fieldNames,
filterProtoBytes = pushedBytes,
pinnedConfig = pinned
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -38,7 +38,7 @@ import org.apache.spark.sql.util.CaseInsensitiveStringMap
*
* Schema discovery happens driver-side inside the bridge's native scan backend
* (`ScanBackend.providerSchemaIpc`), which widens the provider and returns its Arrow schema as
* IPC bytes. The same `optionsProtoBytes` (and the factory FQCN) is then carried verbatim through
* IPC bytes. The same `optionsBytes` (and the factory FQCN) is then carried verbatim through
* `DatafusionInputPartition`, so each executor task repeats the same factory → backend pipeline
* locally.
*/
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -31,7 +31,7 @@ import org.apache.spark.sql.util.CaseInsensitiveStringMap
*/
class DatafusionTable(
val factoryFqcn: String,
val optionsProtoBytes: Array[Byte],
val optionsBytes: Array[Byte],
val sparkSchema: StructType
) extends Table
with SupportsRead {
Expand All @@ -47,5 +47,5 @@ class DatafusionTable(
}

override def newScanBuilder(scanOpts: CaseInsensitiveStringMap): ScanBuilder =
new DatafusionScanBuilder(factoryFqcn, optionsProtoBytes, sparkSchema)
new DatafusionScanBuilder(factoryFqcn, optionsBytes, sparkSchema)
}
Original file line number Diff line number Diff line change
Expand Up @@ -80,7 +80,7 @@ private[spark] object NativeSharedScanResources extends Logging {
// Shared mode builds the dataset-wide provider: empty partitionBytes, like the
// driver-side schema probe. DataFusion-native partitioning replaces listPartitions.
val scanHandle = backend.createScan(
spec.optionsProtoBytes,
spec.optionsBytes,
Array.emptyByteArray,
spec.pinnedConfig.targetPartitions,
spec.pinnedConfig.batchSize,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -31,7 +31,7 @@ import org.apache.arrow.vector.ipc.ArrowReader
final case class SharedScanSpec(
scanId: String,
factoryFqcn: String,
optionsProtoBytes: Array[Byte],
optionsBytes: Array[Byte],
projectionColumnNames: Array[String],
filterProtoBytes: Array[Array[Byte]],
pinnedConfig: PinnedSessionConfig
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -47,7 +47,7 @@ class SharedScanPartitionReader(
throw new IllegalStateException(
s"shared-scan determinism violation for scanId=${partition.scanId}: driver planned " +
s"${partition.numPartitions} partition(s) but this executor planned $executorCount. " +
"The provider's partitioning must be a pure function of optionsProtoBytes; pin your " +
"The provider's partitioning must be a pure function of optionsBytes; pin your " +
"source snapshot (see BridgeProviderFactory.sharedScan).")
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -50,8 +50,8 @@ class BridgeProviderFactoryDefaultsTest extends AnyFunSuite {

override def scanBackend(): ScanBackend = StubBackend

override def listPartitions(optionsProtoBytes: Array[Byte]): Array[PartitionInfo] = {
lastListPartitionsOpts = optionsProtoBytes
override def listPartitions(optionsBytes: Array[Byte]): Array[PartitionInfo] = {
lastListPartitionsOpts = optionsBytes
Array(new PartitionInfo("p0", Array.emptyByteArray, Array.empty[String]))
}
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -32,7 +32,7 @@ class SharedScanCacheTest extends AnyFunSuite {
SharedScanSpec(
scanId = scanId,
factoryFqcn = "test.Factory",
optionsProtoBytes = Array.emptyByteArray,
optionsBytes = Array.emptyByteArray,
projectionColumnNames = Array.empty,
filterProtoBytes = Array.empty,
pinnedConfig = PinnedSessionConfig(8, 8192, Vector.empty)
Expand Down