Skip to content

feat(datafusion): parallel file scanning with eager task bucketing - #42

Closed
phillipleblanc wants to merge 1 commit into
spiceai-0.9.0from
phillip/parallel-file-scanning-2298
Closed

phillipleblanc wants to merge 1 commit into
spiceai-0.9.0from
phillip/parallel-file-scanning-2298

Conversation

@phillipleblanc

Copy link
Copy Markdown

Ports apache/iceberg-rust#2298 ("enable parallel file scanning with eager task bucketing") onto the spiceai-0.9.0 fork branch, ahead of the upstream merge.

What it does

IcebergTableProvider::scan() now plans files eagerly (at planning time) and distributes the resulting FileScanTasks into min(target_partitions, n_files) buckets — one bucket per DataFusion partition — so file reads are scheduled concurrently instead of streaming through a single UnknownPartitioning(1) partition.

When the table is identity-partitioned (single partition spec, supported column types, and the partition columns are present in the projection), the scan declares Partitioning::Hash keyed on those columns, so downstream joins/aggregates can skip a redundant RepartitionExec. Anything that doesn't fully satisfy those conditions falls back to UnknownPartitioning(N) while still bucketing for parallelism.

Key changes

  • TableScan::to_arrow_from_tasks — new public method that replays a pre-collected FileScanTask stream through the Arrow reader (instead of calling plan_files() internally). Reader-side config (concurrency, batch size, row-group filtering, row selection) is taken from self.
  • IcebergTableScan — gains new_with_tasks (eager, multi-partition) alongside new (lazy, single-partition, still used by IcebergStaticTableProvider). execute(i) streams buckets[i]. Constructors are now pub; with_new_children errors on non-empty children.
  • table/bucketing.rs (new) — identity-hash bucketing using REPARTITION_RANDOM_STATE + create_hashes so bucket assignment matches DataFusion's hash-repartition convention; falls back to hashing data_file_path.
  • Drops the unused convert_filters_to_predicate re-export from physical_plan/mod.rs.

Spice fork adaptations (vs. the upstream PR)

  • to_arrow_from_tasks uses this fork's ArrowReaderBuilder::new(file_io, runtime) signature and keeps with_row_selection_enabled.
  • The fork's limit pushdown is preserved: with_limit is threaded into both the planning builder (so plan_files() stops early) and build_table_scan, in addition to the per-partition post-stream truncation.

Validation

  • cargo test -p iceberg-datafusion --lib — 86 passed (incl. the 6 new bucketing tests and the existing limit-pushdown tests).
  • cargo test -p iceberg-sqllogictest — all 9 df_* schedules pass (EXPLAIN snapshots updated for the new buckets:[N] file_count:[M] display and the new input_partitions counts).
  • integration_datafusion_test::test_provider_plan_stream_schema passes (updated execute(1) → execute(0)).
  • clippy + rustfmt clean.
  • Builds cleanly inside the Spice runtime against the Spice DataFusion fork (cargo check -p data_components).

Port of apache#2298 onto the spiceai-0.9.0 fork.

IcebergTableProvider::scan() now plans files eagerly and distributes
FileScanTasks into min(target_partitions, n_files) buckets, one bucket
per DataFusion partition, so file reads are scheduled concurrently
instead of streaming through a single UnknownPartitioning(1) partition.
When the table is identity-partitioned (single spec, supported column
types, partition columns projected) the scan declares Partitioning::Hash
so downstream joins/aggregates can skip a RepartitionExec.

- TableScan::to_arrow_from_tasks: replay pre-collected FileScanTasks
  through the Arrow reader; preserves the spice fork's ArrowReaderBuilder
  (file_io, runtime) signature and row-selection config.
- IcebergTableScan gains new_with_tasks (eager) alongside new (lazy,
  used by IcebergStaticTableProvider); execute(i) streams buckets[i].
  Constructors made pub; with_new_children now errors on children.
- New table/bucketing.rs: identity-hash bucketing via
  REPARTITION_RANDOM_STATE + create_hashes, fallback to data_file_path.
- Spice limit pushdown preserved: with_limit threaded into the planning
  builder and build_table_scan.
- Drop the unused convert_filters_to_predicate re-export.

Copilot AI left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Pull request overview

This PR ports upstream changes to enable parallel Iceberg file scanning in DataFusion by eagerly planning FileScanTasks during TableProvider::scan(), bucketing them into multiple DataFusion partitions, and (when safe) declaring Partitioning::Hash for identity-partitioned tables to avoid redundant repartitioning downstream.

Changes:

  • Eagerly plans file scan tasks and buckets them into min(target_partitions, n_files) partitions for parallel execution.
  • Adds identity-hash bucketing to align bucket assignment with DataFusion’s hash repartitioning when projection/spec conditions are met.
  • Adds TableScan::to_arrow_from_tasks to replay a pre-collected task stream through the Arrow reader, plus associated plan/explain/test updates.

Reviewed changes

Copilot reviewed 11 out of 11 changed files in this pull request and generated 1 comment.

Show a summary per file
File Description
crates/integrations/datafusion/src/table/mod.rs Eager task planning + bucketing/partitioning selection; adds bucketing-focused unit tests.
crates/integrations/datafusion/src/table/bucketing.rs New helper for identity-hash bucketing aligned with DataFusion hashing; fallback hashing for non-identity cases.
crates/integrations/datafusion/src/physical_plan/scan.rs Adds eager multi-partition scan mode (new_with_tasks), partition-aware execute, and updated display output.
crates/iceberg/src/scan/mod.rs Introduces to_arrow_from_tasks and refactors to_arrow to use it.
crates/integrations/datafusion/src/physical_plan/mod.rs Removes unused convert_filters_to_predicate re-export.
crates/integrations/datafusion/tests/integration_datafusion_test.rs Fixes test to execute partition 0 (now guaranteed to exist).
crates/sqllogictest/testdata/slts/df_test/basic_queries.slt Updates EXPLAIN expectations for new scan display (buckets/file_count).
crates/sqllogictest/testdata/slts/df_test/boolean_predicate_pushdown.slt Updates EXPLAIN expectations for new scan display (buckets/file_count).
crates/sqllogictest/testdata/slts/df_test/binary_predicate_pushdown.slt Updates EXPLAIN expectations for new scan display (buckets/file_count).
crates/sqllogictest/testdata/slts/df_test/like_predicate_pushdown.slt Updates EXPLAIN expectations incl. new input_partitions and scan display.
crates/sqllogictest/testdata/slts/df_test/timestamp_predicate_pushdown.slt Updates EXPLAIN expectations for new scan display (buckets/file_count).

💡 Add Copilot custom instructions for smarter, more guided reviews. Learn how to get started.

Comment thread crates/integrations/datafusion/src/table/mod.rs
@phillipleblanc

Copy link
Copy Markdown
Author

Superseded by #43, which targets the renamed spiceai-0.9.1-df-53 base branch (following the new spiceai-<iceberg>-df-<datafusion> convention). Same port, re-based onto the version-corrected (0.9.1) base.

@phillipleblanc
phillipleblanc deleted the phillip/parallel-file-scanning-2298 branch June 15, 2026 02:34
pull Bot pushed a commit to TheRakeshPurohit/spiceai that referenced this pull request Jun 15, 2026
… fork branch naming (spiceai#11328)

* chore(deps): bump iceberg-rust to parallel file scanning fork

Bumps the spiceai/iceberg-rust pin from e519b221 to b652de2e, picking up
the eager task-bucketing parallel-scan change (spiceai/iceberg-rust#42,
port of apache/iceberg-rust#2298).

IcebergTableProvider::scan() now plans files eagerly and distributes
FileScanTasks across min(target_partitions, n_files) DataFusion
partitions, so Iceberg file reads are scheduled concurrently instead of
through a single partition. Identity-partitioned tables additionally
declare Partitioning::Hash so downstream joins/aggregates can skip a
RepartitionExec.

No Spice code changes are required: Spice consumes the provider only via
IcebergTableProvider::try_new, and the change is internal to the scan
path. The fork preserves Spice's existing iceberg limit-pushdown.

* chore(deps): adopt spiceai-<iceberg>-df-<df> fork naming; re-pin to spiceai-0.9.1-df-53

Re-pins iceberg-rust to the renamed `spiceai-0.9.1-df-53` branch (rev
9e6e1a00) and bumps the version requirement to `=0.9.1` to match the
corrected crate version. Adds docs/dev/fork-branch-naming.md documenting
the `spiceai-<iceberg>-df-<datafusion>` convention.

The DF53 fork line was previously named `spiceai-0.9.0` but actually
tracks a post-0.9.1 `main` snapshot; the rename + version correction make
the branch name reflect its real contents.

* chore(deps): re-pin iceberg-rust to merged spiceai-0.9.1-df-53 commit

Fork PR spiceai/iceberg-rust#43 merged; re-pin from the branch head
(9e6e1a00) to the merge commit b8fc79d9 on spiceai-0.9.1-df-53.

* chore(deps): annotate iceberg pins with fork branch; drop naming doc

Address Copilot review on spiceai#11328: add `# branch: spiceai-0.9.1-df-53`
to the iceberg-rust pins (matching the convention used by
delta_kernel/duckdb/arrow). Remove docs/dev/fork-branch-naming.md per
review feedback (no separate doc).
github-actions Bot pushed a commit to spiceai/spiceai that referenced this pull request Jun 16, 2026
… fork branch naming (#11328)

* chore(deps): bump iceberg-rust to parallel file scanning fork

Bumps the spiceai/iceberg-rust pin from e519b221 to b652de2e, picking up
the eager task-bucketing parallel-scan change (spiceai/iceberg-rust#42,
port of apache/iceberg-rust#2298).

IcebergTableProvider::scan() now plans files eagerly and distributes
FileScanTasks across min(target_partitions, n_files) DataFusion
partitions, so Iceberg file reads are scheduled concurrently instead of
through a single partition. Identity-partitioned tables additionally
declare Partitioning::Hash so downstream joins/aggregates can skip a
RepartitionExec.

No Spice code changes are required: Spice consumes the provider only via
IcebergTableProvider::try_new, and the change is internal to the scan
path. The fork preserves Spice's existing iceberg limit-pushdown.

* chore(deps): adopt spiceai-<iceberg>-df-<df> fork naming; re-pin to spiceai-0.9.1-df-53

Re-pins iceberg-rust to the renamed `spiceai-0.9.1-df-53` branch (rev
9e6e1a00) and bumps the version requirement to `=0.9.1` to match the
corrected crate version. Adds docs/dev/fork-branch-naming.md documenting
the `spiceai-<iceberg>-df-<datafusion>` convention.

The DF53 fork line was previously named `spiceai-0.9.0` but actually
tracks a post-0.9.1 `main` snapshot; the rename + version correction make
the branch name reflect its real contents.

* chore(deps): re-pin iceberg-rust to merged spiceai-0.9.1-df-53 commit

Fork PR spiceai/iceberg-rust#43 merged; re-pin from the branch head
(9e6e1a00) to the merge commit b8fc79d9 on spiceai-0.9.1-df-53.

* chore(deps): annotate iceberg pins with fork branch; drop naming doc

Address Copilot review on #11328: add `# branch: spiceai-0.9.1-df-53`
to the iceberg-rust pins (matching the convention used by
delta_kernel/duckdb/arrow). Remove docs/dev/fork-branch-naming.md per
review feedback (no separate doc).
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants