This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
The connector probes the schema once (on the driver) when the `DataFrame` is
created, then plans the scan. `df.filter(...)` and `df.select(...)` are pushed
into the DataFusion scan where possible (see below); the rest Spark evaluates
itself.
## Two partitioning models that must be reconciled
The heart of the connector is mapping DataFusion's idea of a partitioned scan
onto Spark's idea of tasks. The two engines size scans on **different
principles**, and a provider author who ignores the difference will get a job
that runs but performs badly.
### DataFusion: parallelism-bound
DataFusion partitions a scan to **saturate the cores of one machine**. Its
built-in `ListingTable`, for example, packs input files into roughly
`target_partitions` groups — where `target_partitions` defaults to the local
core count — using a minimum-size floor and optional single-file splitting only
to *reach* that count. The partition count is therefore approximately
```text
N ≈ target_partitions ≈ cores
```
independent of how much data there is. More data means **bigger** partitions,
not more of them.
### Spark: one task per partition, byte-bound
Spark turns **each input partition into one task**, and a native Spark file
source sizes partitions by **bytes**: it targets roughly one partition per
`spark.sql.files.maxPartitionBytes` (default 128 MB), cutting files into splits
and bin-packing them to that size. (Spark shrinks the target below 128 MB when
the data is small, so that there are at least enough partitions to keep every
core busy.)
So large data means **many** partitions — typically far more than the cluster
has cores — which run in waves. Each task processes a bounded slice (≈128 MB),
which is what keeps memory bounded and lets Spark reschedule stragglers.
### How the connector joins them
For a scan that only pushes projection / filter / limit (no DataFusion-side
join or aggregation), the physical plan's output partitioning **is** the
provider's `scan()` partitioning, and the mapping to Spark is **1:1**:
```text
provider.scan() output partitions
→ DataFusion physical plan partition count N
→ N ADBC partition descriptors
→ N Spark input partitions = N Spark tasks (scan stage)
```
Spark never splits or merges a partition afterward — `df.rdd().getNumPartitions()
== N`, and the cluster core count only bounds how many of the `N` run at once.
**The consequence you must plan for:** suppose a provider keeps DataFusion's
default of `N = target_partitions ≈ cores` partitions. On a large dataset that
produces only a handful of partitions, each holding about `total_bytes / cores`
of data — so Spark runs a few very large tasks. Because a task is the unit of
work, those oversized tasks bring memory pressure, stragglers that hold up the
whole stage, and no way for Spark to rebalance. To feed Spark well, a provider
should size partitions by **bytes** — the way a Spark file source does — rather
than by the `target_partitions` count.
## Sizing the scan
Aim for the Spark sweet spot:
- **Floor:** `N` ≥ total executor cores, or cores sit idle. `~2–4× cores` gives
slack for skew and stragglers.
- **Per-partition size:** big enough to amortize per-task overhead — the
canonical target is **128–256 MB** of data per partition (matching Spark's
`spark.sql.files.maxPartitionBytes` default of 128 MB).
- **Ceiling — and ADBC lowers it.** Beyond Spark's generic per-task cost, this
path adds two per-partition costs: each of the `N` ADBC descriptors carries a
copy of the **whole serialized physical plan**, and each task **deserializes**
that plan when it reads its partition. (The driver caches a deserialized plan
per connection, but that does not help here: each Spark task uses its own
connection and reads a single partition, so the plan is deserialized once per
task and never reused.) Tens of thousands of tiny partitions is far costlier
here than for a native file source. Keep `N` in the hundreds; prefer bigger
partitions over a huge count.
Rule of thumb, balanced by **bytes** rather than split count:
```text
N ≈ clamp(total_bytes / target_bytes, floor=cores, ceiling≈hundreds)
```
## Writing a provider that feeds Spark well
Your provider's `scan()` receives the session, so it can read the connector's
parallelism hint and then decide based on what it knows about the data:
```text
T = state.config().target_partitions() // connector hint ≈ desired parallelism
target_bytes = ~128–256 MB // bias larger for ADBC to keep N down
min_bytes = floor so partitions never get tiny // ~ tens of MB
if total_bytes known AND splittable (files / row groups / key ranges):
N = clamp(ceil(total_bytes / target_bytes), 1, CEILING) // byte-bound
bin-pack splits into N groups BALANCED BY BYTES // not by count
do not split below min_bytes; merge small splits
elif split_count known but not bytes (shards / external partitions):
N = min(split_count, CEILING) // natural partitioning
if split_count >> T: coalesce shards into ~T balanced groups
else (size and splits unknown; opaque / streaming):
N = T // fall back to the hint
lean slightly high — stragglers hurt more than overhead
```
- **Balance by bytes, not by split count.** A task is atomic; Spark cannot
rebalance a fat partition. One 10 GB partition among 100 small ones stalls the
whole stage. Bin-pack so partition bytes are even.
- **`target_partitions` is a hint, cap, and fallback** — not the data-bound
count. Use it when bytes are unknown or as a parallelism ceiling; when you know
the data size, let **bytes** drive `N`.
### Planning happens once, on the driver
Your provider's `scan()` runs **once**, on the driver, while the multi-partition
descriptors are built. The driver plans the query, fixes the partitioning, and
**serializes the whole physical plan** into each descriptor. Each executor task
then deserializes that plan and executes its partition index. So partition
`index i` always means the same slice: there is no second planning pass that
could disagree with the first, and no need to make `N` stable across a
driver/executor boundary.
The one requirement this places on a provider is that its plan be
**serializable**: built-in DataFusion nodes round-trip through the default codec,
but a custom `ExecutionPlan` node needs a `PhysicalExtensionCodec` registered
with the driver so it can be encoded on the driver and decoded on the executor.
## Connector options
| Option | Required | Meaning |
| --- | --- | --- |
| `driver` | yes | Path to (or manifest name of) the native ADBC driver library. |
| `entrypoint` | depends on driver | Name of the driver's C init symbol (e.g. `AdbcDatafusionExampleInit`). |
| `table` | yes | Table the provider exposes; its schema is probed on the driver. |
| `target_partitions` | no | Parallelism hint. Defaults to `SparkContext.defaultParallelism` (total executor cores). The connector issues `SET datafusion.execution.target_partitions = N` on the planning session, so it influences the partition count when the driver builds the physical plan; that plan (with its partitioning fixed) is what gets serialized into the descriptors. Pass `k × cores` to raise parallelism. |
| `manifest.path` | no | Extra search path for ADBC driver manifests when `driver` is a name rather than an absolute path. |
| *(anything else)* | no | Forwarded verbatim as native ADBC database options, so provider-specific knobs pass straight through. |
Two things are intentionally **not** connector options:
- **`maxPartitionBytes`** — byte-sizing lives in the provider's `scan()`; the
connector does not know the data size and cannot derive it. Pass a sizing knob
as a provider-specific option if your provider supports one.
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
feat: ADBC-backed Spark DataSource for DataFusion table providers #111
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Are you sure you want to change the base?
Uh oh!
There was an error while loading. Please reload this page.
feat: ADBC-backed Spark DataSource for DataFusion table providers #111
Filter by extension
Only manifest files
Viewed files
Uh oh!
There was an error while loading. Please reload this page.
Uh oh!
There was an error while loading. Please reload this page.
There are no files selected for viewing
Uh oh!
There was an error while loading. Please reload this page.
Uh oh!
There was an error while loading. Please reload this page.