Skip to content

feat(datafusion): Add opt-in eager file scan planning with output partitioning - #2671

Open
toutane wants to merge 12 commits into
apache:mainfrom
toutane:datafusion-eager-file-scan-planning
Open

toutane wants to merge 12 commits into
apache:mainfrom
toutane:datafusion-eager-file-scan-planning

Conversation

@toutane

@toutane toutane commented Jun 18, 2026 •

Copy link
Copy Markdown
Contributor

Which issue does this PR close?

What changes are included in this PR?

This PR adds an opt-in eager scan planning path for the DataFusion integration.

When iceberg.enable_eager_scan_planning is enabled, IcebergTableProvider::scan() plans FileScanTasks during physical planning, groups them across DataFusion target_partitions, and exposes the resulting output partition count as UnknownPartitioning(N). Each execute(partition) call then reads only the task group assigned to that output partition through ArrowReaderBuilder.

The default behavior is unchanged: eager planning is disabled by default, so scans keep the existing lazy single-partition planning path unless explicitly enabled.

Eager planning builds a single TableScan per query: the scan that plans the FileScanTasks is retained and reused by every execute(partition) call to derive its ArrowReaderBuilder. Both the eager and the lazy path therefore source their reader settings from TableScan::arrow_reader_builder(), so the two cannot silently drift apart.

When a LIMIT is pushed into the scan, eager planning applies it as a per-partition bound. The final global bound is enforced above the scan by the optimized DataFusion plan, currently through CoalescePartitionsExec::fetch.

This is a narrower split from #2298, focused on file-level eager planning and output partition count reporting. Some of the broader design discussions and review feedback happened in #2298; this PR keeps only the first scoped step and intentionally does not include hash partitioning, size-aware bin-packing, or row-group/sub-file planning.

Public API changes:

Adds one public method to the core crate: iceberg::scan::TableScan::arrow_reader_builder() -> ArrowReaderBuilder. It extracts the reader construction that to_arrow() already performed, so the DataFusion integration can build a reader from a planned scan without duplicating the reader defaults.

Trade-offs:

  • Eager planning does catalog/metadata work during TableProvider::scan(), so it is kept opt-in
  • Task grouping is round-robin and count-based, not size-aware. This is simple and deterministic, but can be imbalanced when file sizes vary (tracked by [iceberg-datafusion] Balance eager scan task groups by estimated read size #2962)
  • Parallelism is at the FileScanTask level only. A table with one large file will not benefit from this change
  • The scan reports UnknownPartitioning(N), not hash partitioning. This exposes the number of output partitions without claiming stronger partitioning semantics
  • Eager planning retains the planning TableScan, and hence its PlanContext and warm evaluator caches, for the lifetime of the physical plan

Follow-up work:

Are these changes tested?

Yes. Integration tests cover:

  • eager scan planning disabled by default (test_multi_file_scan_defaults_to_single_lazy_partition)
  • enabling eager scan planning through iceberg.enable_eager_scan_planning, including at runtime via SET (test_set_enable_eager_scan_planning)
  • exposing multiple scan output partitions, clamped to the planned task count so that no empty partition is exposed (test_multi_file_scan_produces_multiple_partitions)
  • preserving query results between the lazy path, single-partition eager planning, and multi-partition eager planning (test_multi_partition_scan_matches_single_partition_results)
  • enforcing a global LIMIT across multiple output partitions (test_multi_partition_scan_enforces_global_limit)
  • enforcing a global ORDER BY ... LIMIT across multiple output partitions (test_multi_partition_ordered_scan_enforces_global_limit)
  • rejecting non-empty children on the leaf scan node (test_iceberg_table_scan_rejects_non_empty_children)

AI Disclosure

This PR was developed with AI assistance, using Claude Code (Opus 4.8)

@toutane

toutane commented Jun 30, 2026

Copy link
Copy Markdown
Contributor Author

👋 cc @mbutrovich @timsaucer

@mbutrovich
mbutrovich self-requested a review July 9, 2026 15:05

@mbutrovich mbutrovich left a comment

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

First pass, thanks @toutane!

Comment thread crates/integrations/datafusion/src/physical_plan/scan.rs Outdated
Comment thread crates/integrations/datafusion/src/table/mod.rs Outdated
Comment thread crates/integrations/datafusion/src/physical_plan/scan_planning.rs Outdated
Comment thread crates/integrations/datafusion/src/physical_plan/scan.rs Outdated
@toutane

toutane commented Jul 10, 2026 •

Copy link
Copy Markdown
Contributor Author

Thanks @mbutrovich, addressed all four comments in d99603f

@toutane
toutane force-pushed the datafusion-eager-file-scan-planning branch 2 times, most recently from 6a4f1a1 to be6d53e Compare July 10, 2026 13:44
@toutane
toutane requested a review from mbutrovich July 15, 2026 08:28
@toutane

toutane commented Jul 16, 2026

Copy link
Copy Markdown
Contributor Author

@mbutrovich Could you please give this another look?

@toutane
toutane force-pushed the datafusion-eager-file-scan-planning branch from be6d53e to 528019a Compare July 17, 2026 14:28

@mbutrovich mbutrovich left a comment

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

Latest pass. Thanks for continuing to chip away at this, @toutane!

Comment thread crates/integrations/datafusion/src/physical_plan/scan.rs
Comment thread crates/integrations/datafusion/src/physical_plan/scan.rs Outdated
Comment thread crates/integrations/datafusion/src/physical_plan/scan.rs
Comment thread crates/integrations/datafusion/tests/integration_datafusion_test.rs
@toutane
toutane force-pushed the datafusion-eager-file-scan-planning branch from 528019a to 2380cbf Compare July 22, 2026 15:26
@toutane
toutane force-pushed the datafusion-eager-file-scan-planning branch from 2380cbf to 1b0504d Compare July 27, 2026 15:37
@toutane

toutane commented Jul 27, 2026

Copy link
Copy Markdown
Contributor Author

Hello @mbutrovich, thanks for this second pass! I addressed your comments in 1b0504d.

@mbutrovich mbutrovich left a comment

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

Latest pass, thanks @toutane!

Comment thread crates/integrations/datafusion/src/config.rs
Comment thread crates/integrations/datafusion/src/physical_plan/scan.rs
toutane added a commit to DataDog/iceberg-rust that referenced this pull request Jul 29, 2026
* feat(io): route FileIO storage through runtime

(cherry picked from commit 331effe)

* feat(catalog): propagate runtime to table FileIO

(cherry picked from commit 96deb60)

* feat(datafusion): propagate Iceberg runtime through providers

(cherry picked from commit 481b283)

* feat(runtime): add Datadog runtime integration hooks

* feat(datafusion): add opt-in eager scan planning

(cherry picked from commit f0b4c27)

* test(datafusion): cover lazy and eager scan planning

(cherry picked from commit 8ae6eb6)

* fix(datafusion): address eager scan review feedback

(cherry picked from commit 2380cbf)

* fix(datafusion): address eager scan review feedback

Second review round for apache#2671.

---------

Co-authored-by: Geoffrey Claude <geoffrey.claude@datadoghq.com>
toutane added a commit to DataDog/iceberg-rust that referenced this pull request Jul 29, 2026
* feat(datafusion): add opt-in eager scan planning

(cherry picked from commit f0b4c27)

* test(datafusion): cover lazy and eager scan planning

(cherry picked from commit 8ae6eb6)

* fix(datafusion): address eager scan review feedback

(cherry picked from commit 2380cbf)

* fix(datafusion): address eager scan review feedback

Second review round for apache#2671.

---------

Co-authored-by: Geoffrey Claude <geoffrey.claude@datadoghq.com>
toutane added a commit to DataDog/iceberg-rust that referenced this pull request Jul 29, 2026
* feat(datafusion): add opt-in eager scan planning

(cherry picked from commit f0b4c27)

* test(datafusion): cover lazy and eager scan planning

(cherry picked from commit 8ae6eb6)

* fix(datafusion): address eager scan review feedback

(cherry picked from commit 2380cbf)

* fix(datafusion): address eager scan review feedback

Second review round for apache#2671.
@toutane
toutane requested a review from mbutrovich July 29, 2026 15:13
@toutane
toutane force-pushed the datafusion-eager-file-scan-planning branch from 5cd31e5 to 5ae5544 Compare July 30, 2026 08:18
@toutane

toutane commented Jul 30, 2026

Copy link
Copy Markdown
Contributor Author

Thank for this pass @mbutrovich!

@mbutrovich mbutrovich left a comment

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

Another pass, thanks for your patience on this PR, @toutane! I think it's getting close!

Comment thread crates/integrations/datafusion/src/physical_plan/scan.rs Outdated
toutane added 2 commits August 3, 2026 11:24
Eager scan planning built N+1 TableScans per query: one in
plan_file_task_groups to call plan_files(), then one per output partition
in execute() purely to read .arrow_reader_builder() off it.

arrow_reader_builder() reads nothing from the PlanContext, only file_io,
runtime, batch size and the row-group/row-selection flags. So the work
TableScanBuilder::build() does on those per-partition calls (snapshot
resolution, per-column schema validation, field id resolution, predicate
binding, name-mapping JSON parse, three fresh caches) was entirely
discarded.

Introduce EagerScanPlan, which keeps the Arc<TableScan> that planned the
tasks next to the tasks themselves, and have execute() derive its reader
from it. The eager path now builds one TableScan per query.

This keeps both scan paths routed through build_table_scan and
TableScan::arrow_reader_builder, so their reader settings cannot drift
apart. The eager reader is now derived from the very instance that
planned the tasks, making that equivalence structural rather than a call
convention. .build() is still called per partition, so each partition
keeps its own ArrowReader and DeleteFilter as before.
@toutane
toutane requested a review from mbutrovich August 3, 2026 13:38

@mbutrovich mbutrovich left a comment

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

Thanks for staying with this through several rounds, this is close. Two small suggestions inline while we wait for a committer to review, both about test coverage rather than correctness of the feature itself.

/// round-robin assignment. Non-empty groups are bounded by `tasks.len()`.
// TODO: Replace this naive round-robin grouping with size-based grouping once the
// first parallel scan path is stable. Keep this v1 simple and deterministic.
fn group_file_scan_tasks_round_robin(

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

Every test that exercises this function uses a task count equal to (or clamped down to) the partition count, so each group ends up with exactly one task. Would it be worth adding a direct #[test] here for an uneven split, e.g. 5 tasks over 3 partitions landing as [2, 2, 1]? It's a small pure function with no I/O, so a unit test seems cheap, and right now nothing would catch a change that preserves total task/row counts but breaks the round-robin balance itself.

@toutane toutane Aug 5, 2026 •

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Addressed in f66c9b9 viatest_group_file_scan_tasks_round_robin_uneven_split

// The physical optimizer absorbs the initial GlobalLimitExec into
// CoalescePartitionsExec; its fetch enforces the global bound in the final plan.
let global_limit_coalescer = plan
.downcast_ref::<CoalescePartitionsExec>()

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

This is a good end-to-end check that the global limit is actually enforced. One thought: it pins today's DataFusion optimizer output shape (GlobalLimitExec absorbed into CoalescePartitionsExec::fetch). If a future DataFusion bump changes how that composes, this assertion could fail with no regression in this crate's code. Might be worth a short comment here noting that a failure after a datafusion version bump likely means the plan shape changed, not that eager scanning broke, so whoever hits it later doesn't have to rediscover that.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Done in ab4e8ff

@toutane

toutane commented Aug 5, 2026 •

Copy link
Copy Markdown
Contributor Author

Thanks for staying with this through several rounds, this is close. Two small suggestions inline while we wait for a committer to review, both about test coverage rather than correctness of the feature itself.

Thank you for continuing to review this code, @mbutrovich.
Do you know of any committer who might be able to take this on?

@toutane
toutane requested a review from mbutrovich August 5, 2026 11:29

@timsaucer timsaucer left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

I think you've done a great job pushing through on this topic. I know there have been many iterations and I think it's paid off because this looks like a very clean PR. I don't have committer access, so I can only provide an extra "approve" that doesn't carry weight.

The DataFusion parts specifically look very good to me. I am unable to test at scale myself, so I will have to rely on others or on the included tests.

The comments I have are all documentation, not code issues.

// https://github.com/apache/datafusion/blob/ad8e7b7f2babe3fcddc3a4f9b5cd1ac0d1b16ad9/datafusion/datasource/src/file_stream/scan_state.rs#L42-L43
.with_data_file_concurrency_limit(1)
.build()
// TODO: Avoid cloning FileScanTasks here once ArrowReader can accept shared tasks.

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

Is there an open issue for this? If so we should link here so future developers can track this TODO status.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Done: #2964

// Do not cache planned FileScanTasks in the provider in v1. They are query-specific
// because projection, predicate binding, snapshot schema, and delete planning can differ
// between scans. Catalog-backed providers also need fresh metadata on each scan.
// TODO: Revisit provider-level caching for static tables with a precise cache key.

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

Recommend opening an issue and linking here.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Done: #2963

Comment on lines +164 to +165
// TODO: Replace this naive round-robin grouping with size-based grouping once the
// first parallel scan path is stable. Keep this v1 simple and deterministic.

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

Recommend opening an issue and linking here.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Done: #2962

@toutane

toutane commented Aug 5, 2026

Copy link
Copy Markdown
Contributor Author

I think you've done a great job pushing through on this topic. I know there have been many iterations and I think it's paid off because this looks like a very clean PR. I don't have committer access, so I can only provide an extra "approve" that doesn't carry weight.

The DataFusion parts specifically look very good to me. I am unable to test at scale myself, so I will have to rely on others or on the included tests.

The comments I have are all documentation, not code issues.

@timsaucer, thanks a lot for your message, really appreciate it!

For testing, we're going to shadow real traffic in our infra. It won't be a standard benchmark, but it should at least tell us whether things are improving on real-world queries.

That said, I think having a standard way to benchmark this kind of optimization would be really valuable for the iceberg-datafusion crate. Not sure if that's already been discussed somewhere?

Thanks again for the comments - I opened three issues to track the TODOs.

@anuragmantri

Copy link
Copy Markdown
Contributor

Thanks for this work @toutane. Eager planning is what I needed for my work on using the sort orders in #3126

I did a review and it looks good to my untrained eyes (I'm new to Rust and this repo).

@anuragmantri

Copy link
Copy Markdown
Contributor

@mbutrovich @toutane - Are there any pending items in this PR? This is a prerequisite for my work in #3126.

@anuragmantri

anuragmantri commented Sep 15, 2026 •

Copy link
Copy Markdown
Contributor

Hi @toutane, I've been building on top of it for a follow-up (sort-order reporting), so I rebased this branch onto current main to test it and found two things.

  1. main picked up feat(scan): make FileScanTask serializable #3091 (FileScanTask fields went private behind a serde adapter, build() now returns Result<FileScanTask>). Your round-robin test in scan_planning.rs builds a FileScanTask with .build() and no .unwrap(), which no longer compiles.
  2. main also picked up a DataFusion bump to 55.0, which deprecated ExecutionPlan::with_new_children in favor of replace_children/ReplaceChildrenOptions. test_iceberg_table_scan_rejects_non_empty_children calls .with_new_children() directly, which now trips clippy -D warnings.

Both are one-line fixes, but you'll hit them as soon as you rebase onto current main.

@NoahKusaba

NoahKusaba commented Sep 24, 2026 •

Copy link
Copy Markdown
Contributor

Hey @toutane can you re-open your PR in the new datafusion-iceberg repo.
https://github.com/apache/datafusion-iceberg/pulls

@toutane

toutane commented Sep 30, 2026

Copy link
Copy Markdown
Contributor Author

Hi @toutane, I've been building on top of it for a follow-up (sort-order reporting), so I rebased this branch onto current main to test it and found two things.

  1. main picked up feat(scan): make FileScanTask serializable #3091 (FileScanTask fields went private behind a serde adapter, build() now returns Result<FileScanTask>). Your round-robin test in scan_planning.rs builds a FileScanTask with .build() and no .unwrap(), which no longer compiles.
  2. main also picked up a DataFusion bump to 55.0, which deprecated ExecutionPlan::with_new_children in favor of replace_children/ReplaceChildrenOptions. test_iceberg_table_scan_rejects_non_empty_children calls .with_new_children() directly, which now trips clippy -D warnings.

Both are one-line fixes, but you'll hit them as soon as you rebase onto current main.

Hey @anuragmantri, thanks for your comment and for your suggestions! It's so nice to see that you built something on top of it.

I will rebase this branch on top of main and apply the fixes you mentioned - I will probably do it directly into the new repository.

@laskoviymishka laskoviymishka left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Thanks for grinding through this one, @toutane — the eager path has come together nicely over the rounds, and the parity and partition-count tests give me real confidence the feature is doing the right thing. Reading it as a first committer-grade pass, there's one decision I'd want to make deliberately before it merges, and one quick fix.

The decision: TableScan::arrow_reader_builder() is now public on crates/iceberg, which is publish = true. The moment this ships to crates.io we've committed ArrowReaderBuilder — an impl type — to the crate's stable API, and we'd be doing it solely to feed the DataFusion integration. After that, reshaping or removing it is a breaking change, and every future reader setting on TableScan has to be mirrored into it or external callers who captured the builder silently miss it. I'd rather expose a small value type — a TableScanReaderConfig handed back by something like TableScan::reader_config() — that the DataFusion crate turns into its own builder. That keeps the stable surface a plain config that can't be misused as a decoupled reader, and it happens to also fix the memory-retention point I left inline, since the eager plan could store that config instead of the whole TableScan. If we do keep the builder public, the docs need to spell out the plan_files() + build().read() pairing.

The quick fix: the schema.project(projection).unwrap() in IcebergScanConfig::new turns a recoverable error into a panic. Making new fallible and propagating is local, and both call sites already return DFResult.

Before merge I'd also want the eager path's delete-file behavior looked at — build() runs per partition today, so a delete file spanning partitions gets re-fetched N times — and at least one test that scans a table with delete files in eager mode, since none of the new tests do. Neither is a correctness bug, but the first is a real cost behind a performance flag and the second is the gap I'd least want to ship blind. Everything else is inline and minor.

This is close — sort the API shape and I think it's genuinely there.

/// Returns an [`ArrowRecordBatchStream`].
pub async fn to_arrow(&self) -> Result<ArrowRecordBatchStream> {
/// Returns an [`ArrowReaderBuilder`] configured for this table scan.
pub fn arrow_reader_builder(&self) -> ArrowReaderBuilder {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Since crates/iceberg is publish = true, this pins ArrowReaderBuilder — an impl type — into the crate's stable API the moment it ships to crates.io, and we'd be doing it purely to feed the DataFusion integration. Once it's out, reshaping or removing it is a breaking change, and every future reader knob on TableScan has to be mirrored here or external callers who captured the builder silently miss it.

I'd rather not expose the builder itself. Could we hand back a small value type — a TableScanReaderConfig via something like TableScan::reader_config() — that the DataFusion crate turns into its own ArrowReaderBuilder? That keeps the stable surface a plain config that can't be misused as a decoupled reader, and it also solves the Arc<TableScan> retention point I left on scan_planning.rs. If we do keep the builder public, the doc needs to spell out the plan_files() + build().read() pairing, since to_arrow() is the intended path for everyone else.

}

impl ConfigExtension for IcebergDataFusionConfig {
const PREFIX: &'static str = "iceberg";

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

PREFIX = "iceberg" claims the whole iceberg.* option namespace for this one extension. DataFusion wants a unique prefix per registered ExtensionOptions, so the next Iceberg DF extension — write config, catalog auth, whatever — that also picks iceberg collides at registration. I'd scope it now, iceberg.scan or iceberg.datafusion, while nothing depends on the bare prefix yet.

// flight per output partition here.
// https://github.com/apache/datafusion/blob/ad8e7b7f2babe3fcddc3a4f9b5cd1ac0d1b16ad9/datafusion/datasource/src/file_stream/scan_state.rs#L42-L43
.with_data_file_concurrency_limit(1)
.build()

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

build() constructs a fresh CachingDeleteFileLoader per execute(partition), so each output partition gets its own cache. A delete file referenced by tasks in several partitions then gets fetched from object storage once per partition — with target_partitions = 16 and equality deletes that's roughly 16x the delete I/O the lazy path does with its single loader. Results stay correct, but it's a surprising cost behind a flag people flip on for speed. I'd build the reader once in plan_eager_scan() and share it across partitions (ArrowReader is Clone), or if that's follow-up material, call it out in the PR description and a tracking issue alongside the other TODOs here.

While we're on this chain: with_data_file_concurrency_limit(1) silently overrides whatever the user configured (the builder just picked that up inside arrow_reader_builder()). It's a deliberate choice per the comment, but it should be documented on the flag so nobody's surprised their concurrency setting does nothing in eager mode.


Box::pin(stream)
}
None => {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

The eager arm guards partition (the else return a few lines up), but this lazy arm ignores it and always returns the full-scan stream. IcebergTableScan is public, so a caller holding an Arc and calling execute(1) directly gets a valid stream for a partition that doesn't exist — duplicate rows instead of an error. Normal DataFusion flow keeps partition_count at 1 here so it never fires, but the asymmetry is a latent trap; I'd add the same partition != 0 guard returning Internal at the top of this arm.

) -> Self {
let output_schema = match projection {
None => schema.clone(),
Some(projection) => Arc::new(schema.project(projection).unwrap()),

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

I'd make IcebergScanConfig::new fallible and drop this unwrap(). Schema::project returns a Result, and an out-of-bounds projection index turns into a process panic here with nowhere to recover — the kind of thing we've been stamping out across the crate. Both call sites in table/mod.rs are already inside async fn -> DFResult<…>, so returning DFResult<Self> and propagating with .map_err(|e| DataFusionError::ArrowError(e, None))? is a clean, local change. DataFusion validating indices upstream makes it unreachable today, but I'd rather not leave a panic on the read path betting on that staying true.

/// The [`TableScan`] used to plan `task_groups`. Retained so that every output
/// partition builds its reader from this same scan, instead of rebuilding a
/// throwaway `TableScan` on each `execute()` call.
table_scan: Arc<TableScan>,

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Holding Arc<TableScan> here keeps its PlanContext / ObjectCache alive for the whole physical-plan lifetime — that's a moka cache sized to DEFAULT_CACHE_SIZE_BYTES (32MB), populated during plan_files() and never read again after planning. A query joining N eager-scanned tables pins up to N×32MB for its full duration; the lazy path avoids this since its TableScan is built and dropped inside execute(). All execute() actually needs off the scan is the reader settings, so I'd pull those into a small struct at planning time and store that instead, letting the TableScan drop once plan_eager_scan() returns. This is the same reader-config value type that would let us keep arrow_reader_builder off the public API — one change covers both.

use super::*;

#[test]
fn test_group_file_scan_tasks_round_robin_uneven_split() {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

The one unit test covers the 5-tasks / 3-partitions uneven case, but the load-bearing line is the .max(1).min(tasks.len()) clamp — it's what stops us handing DataFusion empty partitions. I'd add direct cases for it: empty input (one empty group), target_partitions = 0 (clamped to 1), and target_partitions > tasks.len() (e.g. 3 tasks / 10 partitions producing 3 groups, not 10). That last one is the case an integration test won't reliably catch, since it leans on the real file count.

let query =
format!("SELECT foo1, foo2 FROM catalog.{namespace_name}.{table_name} LIMIT {limit}");
let dataframe = ctx.sql(&query).await.unwrap();
let plan = dataframe.create_physical_plan().await.unwrap();

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

This calls create_physical_plan() to inspect the plan and then collect() on the dataframe, which re-plans — so plan_eager_scan() runs twice and the plan you assert on isn't the object that executes; the shape assertions don't actually constrain what ran. I'd plan once and execute that same plan. Separately, the downcast_ref::<CoalescePartitionsExec>() pins the DF55 plan shape — a future bump that stacks a GlobalLimitExec on top would panic the downcast rather than fail readably — so I'd lean on the summed-row-count == limit assertion as the real invariant and treat the shape check as best-effort.

let dataframe = ctx.sql(&query).await.unwrap();
let plan = dataframe.create_physical_plan().await.unwrap();

let scan = find_iceberg_scan(plan.as_ref()).expect("Expected IcebergTableScan in ordered plan");

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

This test can't actually catch the risk it's guarding. Each file holds one row, so per-partition truncation to LIMIT 2 can never drop a globally-smaller row — a regression where the limit got pushed into the scan would still pass. The per-partition bound is only safe because DataFusion currently doesn't push a limit through the sort, which leaves self.limit as None for ORDER BY … LIMIT. I'd assert that assumption right here with assert_eq!(scan.limit(), None), and switch to a multi-row-per-file table with ORDER BY foo1 LIMIT K (K below a single file's row count) so per-partition truncation would produce wrong rows if it ever kicked in.

}

#[tokio::test]
async fn test_multi_partition_scan_matches_single_partition_results() -> Result<()> {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

None of the new eager tests exercise delete files — every table here is built without them. The association looks correct by construction (each task carries its own deletes, and round-robin grouping keeps that intact per task), but for a read-path change this is exactly where I'd want a test rather than reasoning. I'd add a positional-delete table scanned in eager mode across multiple partitions and assert the deleted rows are absent, mirroring the parity check this test does — equality deletes too if it's cheap.

This branch has not been deployed

No deployments
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.

Enable parallel file-level scanning for IcebergTableScan Datafusion Integration

6 participants