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]

Reply via email to