You signed in with another tab or window. Reload to refresh your session.You signed out in another tab or window. Reload to refresh your session.You switched accounts on another tab or window. Reload to refresh your session.Dismiss alert
Run Iceberg reads and writes in parallel on Ballista executors: apache/datafusion-ballista#2217 adds an iceberg-ballista crate whose
logical and physical extension codecs send Iceberg plans to the scheduler and
executors.
A codec in another process has to name each node's type, read what it was built
from, and rebuild an equivalent node from those parts alone. This epic tracks the
changes that make that possible, split into what blocks apache/datafusion-ballista#2217 and what improves the integration once it has
landed.
This likely is all relevant to datafusion-distributed as well, but written from the perspective of Ballista, so it may not be consistent to their needs.
I'm likely missing to-do's as well, so please let me know if someone notices anything and I can update this epic.
Required (blocking)
apache/datafusion-ballista#2217 can't merge until these are on
datafusion-iceberg main. Today it depends on a fork branch that carries all of
them except partitioned scans.
feat: make the plan nodes inspectable and rebuildable #20 Make the plan nodes inspectable and rebuildable: public IcebergCommitExec, IcebergWriteExec, IcebergMetadataScan, accessors,
and constructors that take only what the accessors return
IcebergTableProvider::new_with_schema (PR to be opened; draft, stacked on feat: make the plan nodes inspectable and rebuildable #20). Lets the scheduler rebuild a catalog-backed provider with the schema the
client planned against, without loading the table.
Opt-in idempotent commits (PR to be opened; draft, stacked on feat: make the plan nodes inspectable and rebuildable #20). IcebergCommitExec records a commit id in the snapshot summary and skips
committing when an earlier execution already did, since Ballista reruns tasks
after losing an executor. Two attempts running at the same time can still
both commit; see the idempotency key under iceberg-rust below.
The node exposes its task groups and has a constructor that takes them,
like feat: make the plan nodes inspectable and rebuildable #20's nodes. The files must be planned once, on the scheduler, and
shipped to executors: an executor that planned again would read a
different grouping, since file planning order isn't deterministic. FileScanTask is public and serializable.
No dependency on TableScan::arrow_reader_builder(), which #2671 adds to
iceberg-rust but neither iceberg-rust main nor the revision chore(deps): align with current iceberg-rust revision #21 moves
to has.
Build the ArrowReaderBuilder here instead (its setters are public, and Runtime::try_current() works at execute time), or upstream the
accessor first (see iceberg-rust below).
If eager planning stays opt-in (#2671 uses iceberg.enable_eager_scan_planning), the setting reaches the Ballista
scheduler, or iceberg-ballista turns it on.
A scan with no files reports one partition, not zero.
Nice to have
Fixes and improvements that don't block apache/datafusion-ballista#2217. Most
of the iceberg-rust items remove a workaround in iceberg-ballista or close a
known gap.
datafusion-iceberg
fix: accept NOT NULL input columns in partitioned INSERT #19 Accept NOT NULL input columns in partitioned INSERT. Without it, a
partitioned INSERT into optional columns from a NOT NULL source (such as a VALUES list with no NULLs, whose columns DataFusion infers as
non-nullable) fails at plan time.
Hash-partition INSERT input for every partition transform. repartition.rs
hashes on _partition only when the spec has an identity or bucket field,
and otherwise round-robins, which Ballista drops. Each writer task can then
write to every table partition, and the INSERT produces many small files. _partition already holds the transformed value, so hashing on it works
for any transform.
Idempotency key on fast_append, checked inside the commit retry loop.
Closes the concurrent-commit gap above: today a conflicting commit is retried
and re-applied on the new head without re-checking the commit id.
A scan on a runtime that has shut down returns no rows instead of failing
The REST catalog never refreshes its OAuth token, so a long-lived scheduler
fails to plan once its token expires
Debug output of FileIO/StorageConfig prints storage credentials
Table::with_runtime that keeps the table's manifest cache
Goal
Run Iceberg reads and writes in parallel on
Ballista executors: apache/datafusion-ballista#2217 adds an
iceberg-ballistacrate whoselogical and physical extension codecs send Iceberg plans to the scheduler and
executors.
A codec in another process has to name each node's type, read what it was built
from, and rebuild an equivalent node from those parts alone. This epic tracks the
changes that make that possible, split into what blocks
apache/datafusion-ballista#2217 and what improves the integration once it has
landed.
This likely is all relevant to datafusion-distributed as well, but written from the perspective of Ballista, so it may not be consistent to their needs.
I'm likely missing to-do's as well, so please let me know if someone notices anything and I can update this epic.
Required (blocking)
apache/datafusion-ballista#2217 can't merge until these are on
datafusion-iceberg
main. Today it depends on a fork branch that carries all ofthem except partitioned scans.
PartitionExprreconstructible (closes PartitionExpr cannot be reconstructed from its public API #13)IcebergCommitExec,IcebergWriteExec,IcebergMetadataScan, accessors,and constructors that take only what the accessors return
IcebergTableProvider::new_with_schema(PR to be opened; draft, stacked onfeat: make the plan nodes inspectable and rebuildable #20). Lets the scheduler rebuild a catalog-backed provider with the schema the
client planned against, without loading the table.
IcebergCommitExecrecords a commit id in the snapshot summary and skipscommitting when an earlier execution already did, since Ballista reruns tasks
after losing an executor. Two attempts running at the same time can still
both commit; see the idempotency key under iceberg-rust below.
IcebergTableScanreportsUnknownPartitioning(1), soevery scan runs as one task on one executor, which lists and reads every
file. It also limits writes: Ballista drops round-robin repartitions, so an
INSERT into an unpartitioned table, or one with no identity or bucket
partition field, runs one writer task per input partition. Into such a
table,
INSERT INTO t SELECT … FROM <iceberg table>writes from a singletask.
feat(datafusion): Add opt-in eager file scan planning with output partitioning iceberg-rust#2671 (@toutane, closes Enable parallel file-level scanning for IcebergTableScan Datafusion Integration iceberg-rust#2220)
implements this:
scan()plans the file tasks eagerly, groups them acrosstarget_partitions, andexecute(i)reads group i. It went through severalreview rounds there but was stranded by the move to this repo, and
[Epic] Sort-order reporting for iceberg-datafusion scans iceberg-rust#3126 (sort-order reporting) builds on it. Porting it
here also needs, for Ballista:
like feat: make the plan nodes inspectable and rebuildable #20's nodes. The files must be planned once, on the scheduler, and
shipped to executors: an executor that planned again would read a
different grouping, since file planning order isn't deterministic.
FileScanTaskis public and serializable.TableScan::arrow_reader_builder(), which #2671 adds toiceberg-rust but neither iceberg-rust
mainnor the revision chore(deps): align with current iceberg-rust revision #21 movesto has.
Build the
ArrowReaderBuilderhere instead (its setters are public, andRuntime::try_current()works at execute time), or upstream theaccessor first (see iceberg-rust below).
iceberg.enable_eager_scan_planning), the setting reaches the Ballistascheduler, or iceberg-ballista turns it on.
Nice to have
Fixes and improvements that don't block apache/datafusion-ballista#2217. Most
of the iceberg-rust items remove a workaround in iceberg-ballista or close a
known gap.
datafusion-iceberg
partitioned INSERT into optional columns from a NOT NULL source (such as a
VALUESlist with no NULLs, whose columns DataFusion infers asnon-nullable) fails at plan time.
partitioned tables reject nullable sources at plan time
repartition.rshashes on
_partitiononly when the spec has an identity or bucket field,and otherwise round-robins, which Ballista drops. Each writer task can then
write to every table partition, and the INSERT produces many small files.
_partitionalready holds the transformed value, so hashing on it worksfor any transform.
nearly for free, which helps join planning. An earlier attempt,
feat(datafusion): Expose DataFusion statistics on an IcebergTableScan iceberg-rust#880, went stale.
batch_sizeto the scan's reader, and bound per-partitionfile-read concurrency once a scan has several partitions per executor.
Correctness bugs
Opened by @andygrove. None of these blocks apache/datafusion-ballista#2217, but
each returns wrong results or fails in a way a distributed run hits too.
Writes:
the input plan ends in a projection.
PartitionExpr::children()is empty, soProjectionPushdownmoves it onto a different column layout and data filesget the wrong partition values.
IcebergWriteExecdeclares norequired input ordering or distribution, so the optimizer removes the
partition sort. Also covers unpartitioned tables (INSERT into an unpartitioned table fails when fanout is disabled #43, closed as a
duplicate). A required hash distribution on
_partitionwould also helpthe hash-partitioned INSERT item above, since Ballista keeps hash
repartitions but drops round-robin ones.
RecordBatchPartitionSplitteremits a batch's partitions inHashMaporder.Needed on top of INSERT into a partitioned table with fanout disabled always fails: the optimizer removes the partition sort #38 for fanout-disabled writes to work.
Scans and filter pushdown:
fails valid queries (for example a REAL column compared with a DOUBLE
literal)
LIKEprefix pushdown ignores backslash escapes and drops matching rowsconvert_filters_to_predicatenegates a weakened AND, producing a predicate that drops matching rows #44convert_filters_to_predicatenegates a weakenedAND, producing apredicate that drops matching rows
queries return wrong columns or fail. The fix adds a projection to
IcebergMetadataScan, which feat: make the plan nodes inspectable and rebuildable #20 makes public, so its accessors andconstructor need to carry it.
SELECT *fails after a column is added until the next write to the table #42SELECT *fails after a column is added until the next write to thetable. The datafusion-iceberg side of fix(scan): don't project a stale schema after a schema update with no write iceberg-rust#2904 below.
Catalog:
CREATE TABLE/DROP TABLEhang on a current-thread tokio runtime and panic outside a runtime #37CREATE TABLE/DROP TABLEhang on a current-thread tokio runtime andpanic outside a runtime
iceberg-rust
In review:
with no write since (fixes Scan projects a stale schema after a schema update with no write iceberg-rust#2905; surfaces here as
SELECT *fails after a column is added until the next write to the table #42)with_runtimeon the catalog loader'sBoxedCatalogBuilder, so a loaded catalog can be bound to a chosen runtimestrip_metadata_from_schemafails on list and mapcolumns. Blocks the list and map nullability tests staged in fix: accept NOT NULL input columns in partitioned INSERT #19.
Not opened yet, most important first:
fast_append, checked inside the commit retry loop.Closes the concurrent-commit gap above: today a conflicting commit is retried
and re-applied on the new head without re-checking the commit id.
fails to plan once its token expires
Debugoutput ofFileIO/StorageConfigprints storage credentialsTable::with_runtimethat keeps the table's manifest cacheTableScan::arrow_reader_builder(), from feat(datafusion): Add opt-in eager file scan planning with output partitioning iceberg-rust#2671, so apartitioned scan reads with exactly the reader settings its
TableScanwas built with
To investigate: whether
FileIOcan refresh catalog-vended storage credentialsduring a long job.
Merge order
new_with_schemauses feat: make the plan nodes inspectable and rebuildable #20's publictry_newandtable_ident()IcebergCommitExec, which feat: make the plan nodes inspectable and rebuildable #20 makes publicIcebergTableScancode feat: make the plan nodes inspectable and rebuildable #20 changesThe nice-to-haves can land in any order.