NoahKusaba opened a new issue, #23: URL: https://github.com/apache/datafusion-iceberg/issues/23
## Goal Run Iceberg reads and writes in parallel on [Ballista](https://github.com/apache/datafusion-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. ## 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. - [x] #14 Make `PartitionExpr` reconstructible (closes #13) - [ ] #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 #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 #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. - [ ] Partitioned scans. `IcebergTableScan` reports `UnknownPartitioning(1)`, so every 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 single task. apache/iceberg-rust#2671 (@toutane, closes apache/iceberg-rust#2220) implements this: `scan()` plans the file tasks eagerly, groups them across `target_partitions`, and `execute(i)` reads group i. It went through several review rounds there but was stranded by the move to this repo, and apache/iceberg-rust#3126 (sort-order reporting) builds on it. Porting it here also needs, for Ballista: - [ ] The node exposes its task groups and has a constructor that takes them, like #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 #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 - [ ] #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. - [ ] #22 follow-up: unpartitioned tables write NULLs into required columns; partitioned tables reject nullable sources 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. - [ ] Scan statistics. With eager planning, file record counts give row counts nearly for free, which helps join planning. An earlier attempt, apache/iceberg-rust#880, went stale. - [ ] Pass the session's `batch_size` to the scan's reader, and bound per-partition file-read concurrency once a scan has several partitions per executor. ### iceberg-rust In review: - [ ] apache/iceberg-rust#2904 A scan projects a stale schema after a schema change with no write since - [ ] apache/iceberg-rust#3286 `with_runtime` on the catalog loader's `BoxedCatalogBuilder`, so a loaded catalog can be bound to a chosen runtime Fix ready, PR not opened yet: - [ ] apache/iceberg-rust#3297 `strip_metadata_from_schema` fails on list and map columns. Blocks the list and map nullability tests staged in #19. Not opened yet, most important first: - [ ] 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 - [ ] `TableScan::arrow_reader_builder()`, from apache/iceberg-rust#2671, so a partitioned scan reads with exactly the reader settings its `TableScan` was built with To investigate: whether `FileIO` can refresh catalog-vended storage credentials during a long job. ## Merge order 1. #20 2. The rest of the required items, which build on #20, in any order: - `new_with_schema` uses #20's public `try_new` and `table_ident()` - idempotent commits add to `IcebergCommitExec`, which #20 makes public - partitioned scans rework the `IcebergTableScan` code #20 changes The nice-to-haves can land in any order. -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected] --------------------------------------------------------------------- To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
