laskoviymishka commented on code in PR #2671:
URL: https://github.com/apache/iceberg-rust/pull/2671#discussion_r4153186316


##########
crates/iceberg/src/scan/mod.rs:
##########
@@ -498,8 +498,8 @@ impl TableScan {
         Ok(file_scan_task_rx.boxed())
     }
 
-    /// 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 {

Review Comment:
   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.



##########
crates/integrations/datafusion/src/physical_plan/scan_planning.rs:
##########
@@ -0,0 +1,235 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+//   http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied.  See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+use std::sync::Arc;
+
+use datafusion::arrow::datatypes::SchemaRef as ArrowSchemaRef;
+use datafusion::error::Result as DFResult;
+use datafusion::prelude::Expr;
+use futures::TryStreamExt;
+use iceberg::arrow::ArrowReaderBuilder;
+use iceberg::expr::Predicate;
+use iceberg::scan::{FileScanTask, TableScan};
+use iceberg::table::Table;
+
+use super::expr_to_predicate::convert_filters_to_predicate;
+use crate::to_datafusion_error;
+
+#[derive(Debug, Clone)]
+pub(crate) struct IcebergScanConfig {
+    /// Snapshot of the table to scan.
+    snapshot_id: Option<i64>,
+    /// Output schema after projection.
+    output_schema: ArrowSchemaRef,
+    /// Projection column names, None means all columns.
+    column_names: Option<Vec<String>>,
+    /// Filters to apply to the table scan.
+    predicates: Option<Predicate>,
+}
+
+impl IcebergScanConfig {
+    pub(crate) fn new(
+        schema: ArrowSchemaRef,
+        snapshot_id: Option<i64>,
+        projection: Option<&Vec<usize>>,
+        filters: &[Expr],
+    ) -> Self {
+        let output_schema = match projection {
+            None => schema.clone(),
+            Some(projection) => Arc::new(schema.project(projection).unwrap()),

Review Comment:
   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.



##########
crates/integrations/datafusion/src/config.rs:
##########
@@ -0,0 +1,34 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+//   http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied.  See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+use datafusion::common::config::ConfigExtension;
+use datafusion::common::extensions_options;
+
+extensions_options! {
+    /// Configuration options for Iceberg's DataFusion integration.
+    ///
+    /// Register this extension with 
[`SessionConfig::with_option_extension`](datafusion::prelude::SessionConfig::with_option_extension)
+    /// before creating a session to make `SET iceberg.<option> = <value>` 
available.
+    pub struct IcebergDataFusionConfig {
+        /// Plan Iceberg file scan tasks during TableProvider::scan().
+        pub enable_eager_scan_planning: bool, default = false
+    }
+}
+
+impl ConfigExtension for IcebergDataFusionConfig {
+    const PREFIX: &'static str = "iceberg";

Review Comment:
   `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.



##########
crates/integrations/datafusion/src/physical_plan/scan.rs:
##########
@@ -138,18 +134,61 @@ impl ExecutionPlan for IcebergTableScan {
 
     fn execute(
         &self,
-        _partition: usize,
+        partition: usize,
         _context: Arc<TaskContext>,
     ) -> DFResult<SendableRecordBatchStream> {
-        let fut = get_batch_stream(
-            self.table.clone(),
-            self.snapshot_id,
-            self.projection.clone(),
-            self.predicates.clone(),
-        );
-        let stream = futures::stream::once(fut).try_flatten();
-
-        // Apply limit if specified
+        let stream: Pin<Box<dyn Stream<Item = DFResult<RecordBatch>> + Send>> 
= match &self
+            .eager_plan
+        {
+            Some(eager_plan) => {
+                let Some(file_task_group) = eager_plan.task_group(partition) 
else {
+                    return 
Err(datafusion::common::DataFusionError::Internal(format!(
+                        "IcebergTableScan partition {partition} does not 
exist; scan has {} partitions",
+                        eager_plan.partition_count()
+                    )));
+                };
+
+                let tasks: FileScanTaskStream = Box::pin(futures::stream::iter(
+                    (0..file_task_group.len()).map(move |idx| 
Ok(file_task_group[idx].clone())),
+                ));
+                let stream = eager_plan
+                    .arrow_reader_builder()
+                    // Eager planning lets DataFusion drive scan concurrency 
via output
+                    // partitions. Match DataFusion's FileStream model, where 
each
+                    // output partition owns one ScanState; keep one data file 
in
+                    // 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()

Review Comment:
   `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.



##########
crates/integrations/datafusion/src/physical_plan/scan.rs:
##########
@@ -138,18 +134,61 @@ impl ExecutionPlan for IcebergTableScan {
 
     fn execute(
         &self,
-        _partition: usize,
+        partition: usize,
         _context: Arc<TaskContext>,
     ) -> DFResult<SendableRecordBatchStream> {
-        let fut = get_batch_stream(
-            self.table.clone(),
-            self.snapshot_id,
-            self.projection.clone(),
-            self.predicates.clone(),
-        );
-        let stream = futures::stream::once(fut).try_flatten();
-
-        // Apply limit if specified
+        let stream: Pin<Box<dyn Stream<Item = DFResult<RecordBatch>> + Send>> 
= match &self
+            .eager_plan
+        {
+            Some(eager_plan) => {
+                let Some(file_task_group) = eager_plan.task_group(partition) 
else {
+                    return 
Err(datafusion::common::DataFusionError::Internal(format!(
+                        "IcebergTableScan partition {partition} does not 
exist; scan has {} partitions",
+                        eager_plan.partition_count()
+                    )));
+                };
+
+                let tasks: FileScanTaskStream = Box::pin(futures::stream::iter(
+                    (0..file_task_group.len()).map(move |idx| 
Ok(file_task_group[idx].clone())),
+                ));
+                let stream = eager_plan
+                    .arrow_reader_builder()
+                    // Eager planning lets DataFusion drive scan concurrency 
via output
+                    // partitions. Match DataFusion's FileStream model, where 
each
+                    // output partition owns one ScanState; keep one data file 
in
+                    // 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()
+                    // TODO: Avoid cloning FileScanTasks here once ArrowReader 
can accept shared tasks.
+                    // Tracked in 
https://github.com/apache/iceberg-rust/issues/2964.
+                    .read(tasks)
+                    .map_err(to_datafusion_error)?
+                    .stream()
+                    .map_err(to_datafusion_error);
+
+                Box::pin(stream)
+            }
+            None => {

Review Comment:
   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.



##########
crates/integrations/datafusion/src/physical_plan/scan_planning.rs:
##########
@@ -0,0 +1,235 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+//   http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied.  See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+use std::sync::Arc;
+
+use datafusion::arrow::datatypes::SchemaRef as ArrowSchemaRef;
+use datafusion::error::Result as DFResult;
+use datafusion::prelude::Expr;
+use futures::TryStreamExt;
+use iceberg::arrow::ArrowReaderBuilder;
+use iceberg::expr::Predicate;
+use iceberg::scan::{FileScanTask, TableScan};
+use iceberg::table::Table;
+
+use super::expr_to_predicate::convert_filters_to_predicate;
+use crate::to_datafusion_error;
+
+#[derive(Debug, Clone)]
+pub(crate) struct IcebergScanConfig {
+    /// Snapshot of the table to scan.
+    snapshot_id: Option<i64>,
+    /// Output schema after projection.
+    output_schema: ArrowSchemaRef,
+    /// Projection column names, None means all columns.
+    column_names: Option<Vec<String>>,
+    /// Filters to apply to the table scan.
+    predicates: Option<Predicate>,
+}
+
+impl IcebergScanConfig {
+    pub(crate) fn new(
+        schema: ArrowSchemaRef,
+        snapshot_id: Option<i64>,
+        projection: Option<&Vec<usize>>,
+        filters: &[Expr],
+    ) -> Self {
+        let output_schema = match projection {
+            None => schema.clone(),
+            Some(projection) => Arc::new(schema.project(projection).unwrap()),
+        };
+
+        Self {
+            snapshot_id,
+            output_schema,
+            column_names: get_column_names(schema, projection),
+            predicates: convert_filters_to_predicate(filters),
+        }
+    }
+
+    pub(crate) fn snapshot_id(&self) -> Option<i64> {
+        self.snapshot_id
+    }
+
+    pub(crate) fn output_schema(&self) -> ArrowSchemaRef {
+        self.output_schema.clone()
+    }
+
+    pub(crate) fn column_names(&self) -> Option<&[String]> {
+        self.column_names.as_deref()
+    }
+
+    pub(crate) fn predicates(&self) -> Option<&Predicate> {
+        self.predicates.as_ref()
+    }
+}
+
+/// Result of eager scan planning: the [`TableScan`] that planned the file scan
+/// tasks, alongside those tasks grouped per output partition.
+#[derive(Debug)]
+pub(crate) struct EagerScanPlan {
+    /// 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>,

Review Comment:
   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.



##########
crates/integrations/datafusion/src/physical_plan/scan_planning.rs:
##########
@@ -0,0 +1,235 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+//   http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied.  See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+use std::sync::Arc;
+
+use datafusion::arrow::datatypes::SchemaRef as ArrowSchemaRef;
+use datafusion::error::Result as DFResult;
+use datafusion::prelude::Expr;
+use futures::TryStreamExt;
+use iceberg::arrow::ArrowReaderBuilder;
+use iceberg::expr::Predicate;
+use iceberg::scan::{FileScanTask, TableScan};
+use iceberg::table::Table;
+
+use super::expr_to_predicate::convert_filters_to_predicate;
+use crate::to_datafusion_error;
+
+#[derive(Debug, Clone)]
+pub(crate) struct IcebergScanConfig {
+    /// Snapshot of the table to scan.
+    snapshot_id: Option<i64>,
+    /// Output schema after projection.
+    output_schema: ArrowSchemaRef,
+    /// Projection column names, None means all columns.
+    column_names: Option<Vec<String>>,
+    /// Filters to apply to the table scan.
+    predicates: Option<Predicate>,
+}
+
+impl IcebergScanConfig {
+    pub(crate) fn new(
+        schema: ArrowSchemaRef,
+        snapshot_id: Option<i64>,
+        projection: Option<&Vec<usize>>,
+        filters: &[Expr],
+    ) -> Self {
+        let output_schema = match projection {
+            None => schema.clone(),
+            Some(projection) => Arc::new(schema.project(projection).unwrap()),
+        };
+
+        Self {
+            snapshot_id,
+            output_schema,
+            column_names: get_column_names(schema, projection),
+            predicates: convert_filters_to_predicate(filters),
+        }
+    }
+
+    pub(crate) fn snapshot_id(&self) -> Option<i64> {
+        self.snapshot_id
+    }
+
+    pub(crate) fn output_schema(&self) -> ArrowSchemaRef {
+        self.output_schema.clone()
+    }
+
+    pub(crate) fn column_names(&self) -> Option<&[String]> {
+        self.column_names.as_deref()
+    }
+
+    pub(crate) fn predicates(&self) -> Option<&Predicate> {
+        self.predicates.as_ref()
+    }
+}
+
+/// Result of eager scan planning: the [`TableScan`] that planned the file scan
+/// tasks, alongside those tasks grouped per output partition.
+#[derive(Debug)]
+pub(crate) struct EagerScanPlan {
+    /// 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>,
+    /// Planned file scan tasks, one group per output partition.
+    task_groups: Vec<Arc<[FileScanTask]>>,
+}
+
+impl EagerScanPlan {
+    /// Number of output partitions, i.e. the number of task groups.
+    pub(crate) fn partition_count(&self) -> usize {
+        self.task_groups.len()
+    }
+
+    /// Total number of planned file scan tasks across all partitions.
+    pub(crate) fn task_count(&self) -> usize {
+        self.task_groups.iter().map(|group| group.len()).sum()
+    }
+
+    /// Returns the task group assigned to `partition`, or `None` if out of 
range.
+    pub(crate) fn task_group(&self, partition: usize) -> 
Option<Arc<[FileScanTask]>> {
+        self.task_groups.get(partition).cloned()
+    }
+
+    /// Returns an [`ArrowReaderBuilder`] configured for this scan.
+    ///
+    /// This deliberately routes through [`TableScan::arrow_reader_builder`] 
rather
+    /// than constructing an [`ArrowReaderBuilder`] directly: it keeps the 
reader
+    /// settings (batch size, row group filtering, row selection) sourced from 
the
+    /// same place as the lazy path's `TableScan::to_arrow`, so the two scan 
paths
+    /// cannot silently drift apart.
+    pub(crate) fn arrow_reader_builder(&self) -> ArrowReaderBuilder {
+        self.table_scan.arrow_reader_builder()
+    }
+}
+
+pub(crate) async fn plan_eager_scan(
+    table: &Table,
+    scan_config: &IcebergScanConfig,
+    target_partitions: usize,
+) -> DFResult<EagerScanPlan> {
+    // TODO: Cache eager scan planning results across equivalent provider scan 
calls with a
+    // precise cache key. Catalog-backed providers must still reload table 
metadata before
+    // cache lookup so that a new snapshot naturally causes a cache miss.
+    // Tracked in https://github.com/apache/iceberg-rust/issues/2963.
+    let table_scan = Arc::new(build_table_scan(table, scan_config)?);
+
+    let tasks: Vec<FileScanTask> = table_scan
+        .plan_files()
+        .await
+        .map_err(to_datafusion_error)?
+        .try_collect::<Vec<_>>()
+        .await
+        .map_err(to_datafusion_error)?;
+
+    let task_groups = group_file_scan_tasks_round_robin(tasks, 
target_partitions)
+        .into_iter()
+        .map(Arc::<[FileScanTask]>::from)
+        .collect();
+
+    Ok(EagerScanPlan {
+        table_scan,
+        task_groups,
+    })
+}
+
+fn get_column_names(
+    schema: ArrowSchemaRef,
+    projection: Option<&Vec<usize>>,
+) -> Option<Vec<String>> {
+    projection.map(|v| {
+        v.iter()
+            .map(|p| schema.field(*p).name().clone())
+            .collect::<Vec<String>>()
+    })
+}
+
+/// Groups file scan tasks into `target_partitions` groups using a naive
+/// 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.
+// Tracked in https://github.com/apache/iceberg-rust/issues/2962.
+fn group_file_scan_tasks_round_robin(
+    tasks: Vec<FileScanTask>,
+    target_partitions: usize,
+) -> Vec<Vec<FileScanTask>> {
+    if tasks.is_empty() {
+        return vec![vec![]];
+    }
+
+    let target_partitions = target_partitions.max(1).min(tasks.len());
+
+    let mut groups: Vec<Vec<FileScanTask>> = vec![Vec::new(); 
target_partitions];
+    for (i, task) in tasks.into_iter().enumerate() {
+        groups[i % target_partitions].push(task);
+    }
+
+    groups
+}
+
+pub(crate) fn build_table_scan(
+    table: &Table,
+    scan_config: &IcebergScanConfig,
+) -> DFResult<TableScan> {
+    let builder = match scan_config.snapshot_id {
+        Some(id) => table.scan().snapshot_id(id),
+        None => table.scan(),
+    };
+    let mut builder = match scan_config.column_names.clone() {
+        Some(names) => builder.select(names),
+        None => builder.select_all(),
+    };
+    if let Some(pred) = scan_config.predicates.clone() {
+        builder = builder.with_filter(pred);
+    }
+    builder.build().map_err(to_datafusion_error)
+}
+
+#[cfg(test)]
+mod tests {
+    use iceberg::spec::{DataFileFormat, Schema};
+
+    use super::*;
+
+    #[test]
+    fn test_group_file_scan_tasks_round_robin_uneven_split() {

Review Comment:
   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.



##########
crates/integrations/datafusion/tests/integration_datafusion_test.rs:
##########
@@ -95,6 +112,532 @@ fn get_table_creation(
     Ok(creation)
 }
 
+async fn get_multi_file_table_context(
+    namespace_name: &str,
+    table_name: &str,
+    data_file_count: usize,
+) -> Result<(Arc<MemoryCatalog>, NamespaceIdent, String)> {
+    let iceberg_catalog = Arc::new(get_iceberg_catalog().await);
+    let namespace = NamespaceIdent::new(namespace_name.to_string());
+    set_test_namespace(&iceberg_catalog, &namespace).await?;
+
+    let creation = get_table_creation(temp_path(), table_name, None)?;
+    iceberg_catalog.create_table(&namespace, creation).await?;
+
+    let write_ctx = SessionContext::new_with_config(
+        SessionConfig::new().with_target_partitions(data_file_count),
+    );
+    let arrow_schema = Arc::new(ArrowSchema::new(vec![
+        Field::new("foo1", DataType::Int32, false),
+        Field::new("foo2", DataType::Utf8, false),
+    ]));
+
+    let batches: Vec<RecordBatch> = (1..=data_file_count as i32)
+        .map(|idx| {
+            RecordBatch::try_new(arrow_schema.clone(), vec![
+                Arc::new(Int32Array::from(vec![idx])) as ArrayRef,
+                Arc::new(StringArray::from(vec![format!("row-{idx}")])) as 
ArrayRef,
+            ])
+        })
+        .collect::<std::result::Result<_, _>>()?;
+
+    let partitions = batches.into_iter().map(|batch| vec![batch]).collect();
+    let source_table = Arc::new(MemTable::try_new(arrow_schema, 
partitions).unwrap());
+    write_ctx
+        .register_table("source_table", source_table)
+        .unwrap();
+
+    let catalog = 
Arc::new(IcebergCatalogProvider::try_new(iceberg_catalog.clone()).await?);
+    write_ctx.register_catalog("catalog", catalog);
+
+    let insert_sql =
+        format!("INSERT INTO catalog.{namespace_name}.{table_name} SELECT * 
FROM source_table");
+    let batches = write_ctx
+        .sql(&insert_sql)
+        .await
+        .unwrap()
+        .collect()
+        .await
+        .unwrap();
+    assert_eq!(batches.len(), 1);
+
+    let rows_inserted = batches[0]
+        .column(0)
+        .as_any()
+        .downcast_ref::<UInt64Array>()
+        .unwrap();
+    assert_eq!(rows_inserted.value(0), data_file_count as u64);
+
+    Ok((iceberg_catalog, namespace, table_name.to_string()))
+}
+
+async fn get_multi_row_group_table_context(
+    namespace_name: &str,
+    table_name: &str,
+) -> Result<(Arc<MemoryCatalog>, NamespaceIdent, String)> {
+    let iceberg_catalog = Arc::new(get_iceberg_catalog().await);
+    let namespace = NamespaceIdent::new(namespace_name.to_string());
+    set_test_namespace(&iceberg_catalog, &namespace).await?;
+
+    let creation = get_table_creation(temp_path(), table_name, None)?;
+    let table = iceberg_catalog.create_table(&namespace, creation).await?;
+    let arrow_schema: Arc<ArrowSchema> = Arc::new(
+        table
+            .metadata()
+            .current_schema()
+            .as_ref()
+            .try_into()
+            .unwrap(),
+    );
+
+    let file_rows = [
+        vec![
+            (1, "row-1"),
+            (2, "row-2"),
+            (100, "row-100"),
+            (101, "row-101"),
+        ],
+        vec![
+            (50, "row-50"),
+            (150, "row-150"),
+            (151, "row-151"),
+            (152, "row-152"),
+        ],
+        vec![
+            (99, "row-99"),
+            (1000, "row-1000"),
+            (1001, "row-1001"),
+            (1002, "row-1002"),
+        ],
+    ];
+
+    let location_generator = 
DefaultLocationGenerator::new(table.metadata()).unwrap();
+    let writer_properties = WriterProperties::builder()
+        .set_max_row_group_row_count(Some(2))
+        .build();
+    let mut data_files = Vec::with_capacity(file_rows.len());
+
+    for (file_idx, rows) in file_rows.into_iter().enumerate() {
+        let parquet_writer_builder = ParquetWriterBuilder::new(
+            writer_properties.clone(),
+            table.metadata().current_schema().clone(),
+        );
+        let file_name_generator = DefaultFileNameGenerator::new(
+            format!("multi-row-group-{file_idx}"),
+            None,
+            iceberg::spec::DataFileFormat::Parquet,
+        );
+        let rolling_file_writer_builder = 
RollingFileWriterBuilder::new_with_default_file_size(
+            parquet_writer_builder,
+            table.file_io().clone(),
+            location_generator.clone(),
+            file_name_generator,
+        );
+        let data_file_writer_builder = 
DataFileWriterBuilder::new(rolling_file_writer_builder);
+        let mut data_file_writer = data_file_writer_builder.build(None).await?;
+        let batch = RecordBatch::try_new(arrow_schema.clone(), vec![
+            Arc::new(Int32Array::from(
+                rows.iter().map(|(foo1, _)| *foo1).collect::<Vec<_>>(),
+            )) as ArrayRef,
+            Arc::new(StringArray::from(
+                rows.iter().map(|(_, foo2)| *foo2).collect::<Vec<_>>(),
+            )) as ArrayRef,
+        ])?;
+
+        data_file_writer.write(batch).await?;
+        let file_data_files = data_file_writer.close().await?;
+        assert_eq!(file_data_files.len(), 1);
+        data_files.extend(file_data_files);
+    }
+
+    for data_file in &data_files {
+        let file_path = data_file
+            .file_path()
+            .strip_prefix("file://")
+            .unwrap_or(data_file.file_path());
+        let reader = SerializedFileReader::new(File::open(file_path)?)?;
+        assert!(
+            reader.metadata().num_row_groups() > 1,
+            "expected multiple row groups for {}",
+            data_file.file_path()
+        );
+    }
+
+    let tx = Transaction::new(&table);
+    let action = tx.fast_append().add_data_files(data_files);
+    let tx = action.apply(tx)?;
+    tx.commit(iceberg_catalog.as_ref()).await?;
+
+    Ok((iceberg_catalog, namespace, table_name.to_string()))
+}
+
+async fn get_read_context(
+    catalog: Arc<MemoryCatalog>,
+    target_partitions: usize,
+    enable_eager_scan_planning: Option<bool>,
+) -> Result<SessionContext> {
+    let ctx = SessionContext::new_with_config(read_session_config(
+        target_partitions,
+        enable_eager_scan_planning,
+    ));
+    let catalog = Arc::new(IcebergCatalogProvider::try_new(catalog).await?);
+    ctx.register_catalog("catalog", catalog);
+    Ok(ctx)
+}
+
+fn read_session_config(
+    target_partitions: usize,
+    enable_eager_scan_planning: Option<bool>,
+) -> SessionConfig {
+    let config = 
SessionConfig::new().with_target_partitions(target_partitions);
+
+    match enable_eager_scan_planning {
+        Some(enabled) => {
+            let mut iceberg_config = IcebergDataFusionConfig::default();
+            iceberg_config.enable_eager_scan_planning = enabled;
+            config.with_option_extension(iceberg_config)
+        }
+        None => config,
+    }
+}
+
+async fn scan_plan(
+    ctx: &SessionContext,
+    namespace: &NamespaceIdent,
+    table_name: &str,
+) -> Arc<dyn ExecutionPlan> {
+    let provider = ctx.catalog("catalog").unwrap();
+    let namespace_name = &namespace[0];
+    let schema = provider.schema(namespace_name).unwrap();
+    let table = schema.table(table_name).await.unwrap().unwrap();
+
+    let state = ctx.state();
+    table.scan(&state, None, &[], None).await.unwrap()
+}
+
+async fn scan_partition_count(
+    ctx: &SessionContext,
+    namespace: &NamespaceIdent,
+    table_name: &str,
+) -> usize {
+    let plan = scan_plan(ctx, namespace, table_name).await;
+    plan.downcast_ref::<IcebergTableScan>()
+        .expect("Expected IcebergTableScan");
+    plan.properties().output_partitioning().partition_count()
+}
+
+fn find_iceberg_scan(plan: &dyn ExecutionPlan) -> Option<&IcebergTableScan> {
+    if let Some(scan) = plan.downcast_ref::<IcebergTableScan>() {
+        return Some(scan);
+    }
+
+    plan.children()
+        .into_iter()
+        .find_map(|child| find_iceberg_scan(child.as_ref()))
+}
+
+#[tokio::test]
+async fn test_multi_file_scan_produces_multiple_partitions() -> Result<()> {
+    let data_file_count = 3;
+    // Ask for more partitions than files to verify scan planning does not 
expose empty partitions.
+    let target_partitions = data_file_count + 1;
+    let (iceberg_catalog, namespace, table_name) = 
get_multi_file_table_context(
+        "test_multi_file_scan_partitions",
+        "my_table",
+        data_file_count,
+    )
+    .await?;
+    let ctx = get_read_context(iceberg_catalog, target_partitions, 
Some(true)).await?;
+    let plan = scan_plan(&ctx, &namespace, &table_name).await;
+    plan.downcast_ref::<IcebergTableScan>()
+        .expect("Expected IcebergTableScan");
+
+    let actual_partition_count = 
plan.properties().output_partitioning().partition_count();
+    assert_eq!(actual_partition_count, data_file_count);
+
+    // Pins the eager plan's task-group accounting. The reader settings it 
derives from
+    // the retained TableScan are not introspectable, so their equivalence 
with the lazy
+    // path is covered behaviorally by 
test_multi_partition_scan_matches_single_partition_results.
+    let display = datafusion::physical_plan::displayable(plan.as_ref())
+        .one_line()
+        .to_string();
+    assert!(
+        display.contains("task_groups:[3] tasks:[3]"),
+        "unexpected eager scan display: {display}"
+    );
+
+    Ok(())
+}
+
+#[tokio::test]
+async fn test_multi_file_scan_defaults_to_single_lazy_partition() -> 
Result<()> {
+    let data_file_count = 3;
+    let target_partitions = data_file_count + 1;
+    let (iceberg_catalog, namespace, table_name) = 
get_multi_file_table_context(
+        "test_multi_file_scan_default_lazy",
+        "my_table",
+        data_file_count,
+    )
+    .await?;
+    let ctx = get_read_context(iceberg_catalog, target_partitions, 
None).await?;
+
+    let actual_partition_count = scan_partition_count(&ctx, &namespace, 
&table_name).await;
+
+    assert_eq!(actual_partition_count, 1);
+
+    Ok(())
+}
+
+#[tokio::test]
+async fn test_set_enable_eager_scan_planning() -> Result<()> {
+    let data_file_count = 3;
+    let target_partitions = data_file_count + 1;
+    let (iceberg_catalog, namespace, table_name) =
+        get_multi_file_table_context("test_set_eager_scan_planning", 
"my_table", data_file_count)
+            .await?;
+    let ctx = get_read_context(iceberg_catalog, target_partitions, 
Some(false)).await?;
+
+    ctx.sql("SET iceberg.enable_eager_scan_planning = true")
+        .await
+        .unwrap()
+        .collect()
+        .await
+        .unwrap();
+
+    let actual_partition_count = scan_partition_count(&ctx, &namespace, 
&table_name).await;
+
+    assert_eq!(actual_partition_count, data_file_count);
+
+    Ok(())
+}
+
+#[tokio::test]
+async fn test_multi_partition_scan_enforces_global_limit() -> Result<()> {
+    let data_file_count = 3;
+    let limit = 2;
+    let target_partitions = data_file_count + 1;
+    let (iceberg_catalog, namespace, table_name) = 
get_multi_file_table_context(
+        "test_multi_partition_scan_limit",
+        "my_table",
+        data_file_count,
+    )
+    .await?;
+    let ctx = get_read_context(iceberg_catalog, target_partitions, 
Some(true)).await?;
+    let namespace_name = &namespace[0];
+    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();

Review Comment:
   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.



##########
crates/integrations/datafusion/tests/integration_datafusion_test.rs:
##########
@@ -95,6 +112,532 @@ fn get_table_creation(
     Ok(creation)
 }
 
+async fn get_multi_file_table_context(
+    namespace_name: &str,
+    table_name: &str,
+    data_file_count: usize,
+) -> Result<(Arc<MemoryCatalog>, NamespaceIdent, String)> {
+    let iceberg_catalog = Arc::new(get_iceberg_catalog().await);
+    let namespace = NamespaceIdent::new(namespace_name.to_string());
+    set_test_namespace(&iceberg_catalog, &namespace).await?;
+
+    let creation = get_table_creation(temp_path(), table_name, None)?;
+    iceberg_catalog.create_table(&namespace, creation).await?;
+
+    let write_ctx = SessionContext::new_with_config(
+        SessionConfig::new().with_target_partitions(data_file_count),
+    );
+    let arrow_schema = Arc::new(ArrowSchema::new(vec![
+        Field::new("foo1", DataType::Int32, false),
+        Field::new("foo2", DataType::Utf8, false),
+    ]));
+
+    let batches: Vec<RecordBatch> = (1..=data_file_count as i32)
+        .map(|idx| {
+            RecordBatch::try_new(arrow_schema.clone(), vec![
+                Arc::new(Int32Array::from(vec![idx])) as ArrayRef,
+                Arc::new(StringArray::from(vec![format!("row-{idx}")])) as 
ArrayRef,
+            ])
+        })
+        .collect::<std::result::Result<_, _>>()?;
+
+    let partitions = batches.into_iter().map(|batch| vec![batch]).collect();
+    let source_table = Arc::new(MemTable::try_new(arrow_schema, 
partitions).unwrap());
+    write_ctx
+        .register_table("source_table", source_table)
+        .unwrap();
+
+    let catalog = 
Arc::new(IcebergCatalogProvider::try_new(iceberg_catalog.clone()).await?);
+    write_ctx.register_catalog("catalog", catalog);
+
+    let insert_sql =
+        format!("INSERT INTO catalog.{namespace_name}.{table_name} SELECT * 
FROM source_table");
+    let batches = write_ctx
+        .sql(&insert_sql)
+        .await
+        .unwrap()
+        .collect()
+        .await
+        .unwrap();
+    assert_eq!(batches.len(), 1);
+
+    let rows_inserted = batches[0]
+        .column(0)
+        .as_any()
+        .downcast_ref::<UInt64Array>()
+        .unwrap();
+    assert_eq!(rows_inserted.value(0), data_file_count as u64);
+
+    Ok((iceberg_catalog, namespace, table_name.to_string()))
+}
+
+async fn get_multi_row_group_table_context(
+    namespace_name: &str,
+    table_name: &str,
+) -> Result<(Arc<MemoryCatalog>, NamespaceIdent, String)> {
+    let iceberg_catalog = Arc::new(get_iceberg_catalog().await);
+    let namespace = NamespaceIdent::new(namespace_name.to_string());
+    set_test_namespace(&iceberg_catalog, &namespace).await?;
+
+    let creation = get_table_creation(temp_path(), table_name, None)?;
+    let table = iceberg_catalog.create_table(&namespace, creation).await?;
+    let arrow_schema: Arc<ArrowSchema> = Arc::new(
+        table
+            .metadata()
+            .current_schema()
+            .as_ref()
+            .try_into()
+            .unwrap(),
+    );
+
+    let file_rows = [
+        vec![
+            (1, "row-1"),
+            (2, "row-2"),
+            (100, "row-100"),
+            (101, "row-101"),
+        ],
+        vec![
+            (50, "row-50"),
+            (150, "row-150"),
+            (151, "row-151"),
+            (152, "row-152"),
+        ],
+        vec![
+            (99, "row-99"),
+            (1000, "row-1000"),
+            (1001, "row-1001"),
+            (1002, "row-1002"),
+        ],
+    ];
+
+    let location_generator = 
DefaultLocationGenerator::new(table.metadata()).unwrap();
+    let writer_properties = WriterProperties::builder()
+        .set_max_row_group_row_count(Some(2))
+        .build();
+    let mut data_files = Vec::with_capacity(file_rows.len());
+
+    for (file_idx, rows) in file_rows.into_iter().enumerate() {
+        let parquet_writer_builder = ParquetWriterBuilder::new(
+            writer_properties.clone(),
+            table.metadata().current_schema().clone(),
+        );
+        let file_name_generator = DefaultFileNameGenerator::new(
+            format!("multi-row-group-{file_idx}"),
+            None,
+            iceberg::spec::DataFileFormat::Parquet,
+        );
+        let rolling_file_writer_builder = 
RollingFileWriterBuilder::new_with_default_file_size(
+            parquet_writer_builder,
+            table.file_io().clone(),
+            location_generator.clone(),
+            file_name_generator,
+        );
+        let data_file_writer_builder = 
DataFileWriterBuilder::new(rolling_file_writer_builder);
+        let mut data_file_writer = data_file_writer_builder.build(None).await?;
+        let batch = RecordBatch::try_new(arrow_schema.clone(), vec![
+            Arc::new(Int32Array::from(
+                rows.iter().map(|(foo1, _)| *foo1).collect::<Vec<_>>(),
+            )) as ArrayRef,
+            Arc::new(StringArray::from(
+                rows.iter().map(|(_, foo2)| *foo2).collect::<Vec<_>>(),
+            )) as ArrayRef,
+        ])?;
+
+        data_file_writer.write(batch).await?;
+        let file_data_files = data_file_writer.close().await?;
+        assert_eq!(file_data_files.len(), 1);
+        data_files.extend(file_data_files);
+    }
+
+    for data_file in &data_files {
+        let file_path = data_file
+            .file_path()
+            .strip_prefix("file://")
+            .unwrap_or(data_file.file_path());
+        let reader = SerializedFileReader::new(File::open(file_path)?)?;
+        assert!(
+            reader.metadata().num_row_groups() > 1,
+            "expected multiple row groups for {}",
+            data_file.file_path()
+        );
+    }
+
+    let tx = Transaction::new(&table);
+    let action = tx.fast_append().add_data_files(data_files);
+    let tx = action.apply(tx)?;
+    tx.commit(iceberg_catalog.as_ref()).await?;
+
+    Ok((iceberg_catalog, namespace, table_name.to_string()))
+}
+
+async fn get_read_context(
+    catalog: Arc<MemoryCatalog>,
+    target_partitions: usize,
+    enable_eager_scan_planning: Option<bool>,
+) -> Result<SessionContext> {
+    let ctx = SessionContext::new_with_config(read_session_config(
+        target_partitions,
+        enable_eager_scan_planning,
+    ));
+    let catalog = Arc::new(IcebergCatalogProvider::try_new(catalog).await?);
+    ctx.register_catalog("catalog", catalog);
+    Ok(ctx)
+}
+
+fn read_session_config(
+    target_partitions: usize,
+    enable_eager_scan_planning: Option<bool>,
+) -> SessionConfig {
+    let config = 
SessionConfig::new().with_target_partitions(target_partitions);
+
+    match enable_eager_scan_planning {
+        Some(enabled) => {
+            let mut iceberg_config = IcebergDataFusionConfig::default();
+            iceberg_config.enable_eager_scan_planning = enabled;
+            config.with_option_extension(iceberg_config)
+        }
+        None => config,
+    }
+}
+
+async fn scan_plan(
+    ctx: &SessionContext,
+    namespace: &NamespaceIdent,
+    table_name: &str,
+) -> Arc<dyn ExecutionPlan> {
+    let provider = ctx.catalog("catalog").unwrap();
+    let namespace_name = &namespace[0];
+    let schema = provider.schema(namespace_name).unwrap();
+    let table = schema.table(table_name).await.unwrap().unwrap();
+
+    let state = ctx.state();
+    table.scan(&state, None, &[], None).await.unwrap()
+}
+
+async fn scan_partition_count(
+    ctx: &SessionContext,
+    namespace: &NamespaceIdent,
+    table_name: &str,
+) -> usize {
+    let plan = scan_plan(ctx, namespace, table_name).await;
+    plan.downcast_ref::<IcebergTableScan>()
+        .expect("Expected IcebergTableScan");
+    plan.properties().output_partitioning().partition_count()
+}
+
+fn find_iceberg_scan(plan: &dyn ExecutionPlan) -> Option<&IcebergTableScan> {
+    if let Some(scan) = plan.downcast_ref::<IcebergTableScan>() {
+        return Some(scan);
+    }
+
+    plan.children()
+        .into_iter()
+        .find_map(|child| find_iceberg_scan(child.as_ref()))
+}
+
+#[tokio::test]
+async fn test_multi_file_scan_produces_multiple_partitions() -> Result<()> {
+    let data_file_count = 3;
+    // Ask for more partitions than files to verify scan planning does not 
expose empty partitions.
+    let target_partitions = data_file_count + 1;
+    let (iceberg_catalog, namespace, table_name) = 
get_multi_file_table_context(
+        "test_multi_file_scan_partitions",
+        "my_table",
+        data_file_count,
+    )
+    .await?;
+    let ctx = get_read_context(iceberg_catalog, target_partitions, 
Some(true)).await?;
+    let plan = scan_plan(&ctx, &namespace, &table_name).await;
+    plan.downcast_ref::<IcebergTableScan>()
+        .expect("Expected IcebergTableScan");
+
+    let actual_partition_count = 
plan.properties().output_partitioning().partition_count();
+    assert_eq!(actual_partition_count, data_file_count);
+
+    // Pins the eager plan's task-group accounting. The reader settings it 
derives from
+    // the retained TableScan are not introspectable, so their equivalence 
with the lazy
+    // path is covered behaviorally by 
test_multi_partition_scan_matches_single_partition_results.
+    let display = datafusion::physical_plan::displayable(plan.as_ref())
+        .one_line()
+        .to_string();
+    assert!(
+        display.contains("task_groups:[3] tasks:[3]"),
+        "unexpected eager scan display: {display}"
+    );
+
+    Ok(())
+}
+
+#[tokio::test]
+async fn test_multi_file_scan_defaults_to_single_lazy_partition() -> 
Result<()> {
+    let data_file_count = 3;
+    let target_partitions = data_file_count + 1;
+    let (iceberg_catalog, namespace, table_name) = 
get_multi_file_table_context(
+        "test_multi_file_scan_default_lazy",
+        "my_table",
+        data_file_count,
+    )
+    .await?;
+    let ctx = get_read_context(iceberg_catalog, target_partitions, 
None).await?;
+
+    let actual_partition_count = scan_partition_count(&ctx, &namespace, 
&table_name).await;
+
+    assert_eq!(actual_partition_count, 1);
+
+    Ok(())
+}
+
+#[tokio::test]
+async fn test_set_enable_eager_scan_planning() -> Result<()> {
+    let data_file_count = 3;
+    let target_partitions = data_file_count + 1;
+    let (iceberg_catalog, namespace, table_name) =
+        get_multi_file_table_context("test_set_eager_scan_planning", 
"my_table", data_file_count)
+            .await?;
+    let ctx = get_read_context(iceberg_catalog, target_partitions, 
Some(false)).await?;
+
+    ctx.sql("SET iceberg.enable_eager_scan_planning = true")
+        .await
+        .unwrap()
+        .collect()
+        .await
+        .unwrap();
+
+    let actual_partition_count = scan_partition_count(&ctx, &namespace, 
&table_name).await;
+
+    assert_eq!(actual_partition_count, data_file_count);
+
+    Ok(())
+}
+
+#[tokio::test]
+async fn test_multi_partition_scan_enforces_global_limit() -> Result<()> {
+    let data_file_count = 3;
+    let limit = 2;
+    let target_partitions = data_file_count + 1;
+    let (iceberg_catalog, namespace, table_name) = 
get_multi_file_table_context(
+        "test_multi_partition_scan_limit",
+        "my_table",
+        data_file_count,
+    )
+    .await?;
+    let ctx = get_read_context(iceberg_catalog, target_partitions, 
Some(true)).await?;
+    let namespace_name = &namespace[0];
+    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();
+
+    // The physical optimizer absorbs the initial GlobalLimitExec into
+    // CoalescePartitionsExec; its fetch enforces the global bound in the 
final plan.
+    // After a DataFusion bump, failure here likely means the plan shape 
changed, not eager scanning.
+    let global_limit_coalescer = plan
+        .downcast_ref::<CoalescePartitionsExec>()
+        .expect("Expected a globally limited CoalescePartitionsExec");
+    assert_eq!(global_limit_coalescer.fetch(), Some(limit));
+
+    let scan = find_iceberg_scan(global_limit_coalescer.input().as_ref())
+        .expect("Expected IcebergTableScan below the global limit");
+    assert_eq!(scan.limit(), Some(limit));
+    assert_eq!(
+        scan.properties().output_partitioning().partition_count(),
+        data_file_count
+    );
+
+    let batches = dataframe.collect().await.unwrap();
+    let mut seen_ids = std::collections::HashSet::new();
+    for batch in &batches {
+        let foo1 = batch
+            .column(0)
+            .as_any()
+            .downcast_ref::<Int32Array>()
+            .unwrap();
+        let foo2 = batch
+            .column(1)
+            .as_any()
+            .downcast_ref::<StringArray>()
+            .unwrap();
+
+        for row in 0..batch.num_rows() {
+            let id = foo1.value(row);
+            assert!((1..=data_file_count as i32).contains(&id));
+            assert_eq!(foo2.value(row), format!("row-{id}"));
+            assert!(seen_ids.insert(id), "duplicate row {id}");
+        }
+    }
+
+    assert_eq!(
+        batches.iter().map(RecordBatch::num_rows).sum::<usize>(),
+        limit
+    );
+    assert_eq!(seen_ids.len(), limit);
+
+    Ok(())
+}
+
+#[tokio::test]
+async fn test_multi_partition_ordered_scan_enforces_global_limit() -> 
Result<()> {
+    let data_file_count = 3;
+    let limit = 2;
+    let target_partitions = data_file_count + 1;
+    let (iceberg_catalog, namespace, table_name) = 
get_multi_file_table_context(
+        "test_multi_partition_ordered_scan_limit",
+        "my_table",
+        data_file_count,
+    )
+    .await?;
+    let ctx = get_read_context(iceberg_catalog, target_partitions, 
Some(true)).await?;
+    let namespace_name = &namespace[0];
+    let query = format!(
+        "SELECT foo1, foo2 FROM catalog.{namespace_name}.{table_name} ORDER BY 
foo1 LIMIT {limit}"
+    );
+    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");

Review Comment:
   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.



##########
crates/integrations/datafusion/tests/integration_datafusion_test.rs:
##########
@@ -95,6 +112,532 @@ fn get_table_creation(
     Ok(creation)
 }
 
+async fn get_multi_file_table_context(
+    namespace_name: &str,
+    table_name: &str,
+    data_file_count: usize,
+) -> Result<(Arc<MemoryCatalog>, NamespaceIdent, String)> {
+    let iceberg_catalog = Arc::new(get_iceberg_catalog().await);
+    let namespace = NamespaceIdent::new(namespace_name.to_string());
+    set_test_namespace(&iceberg_catalog, &namespace).await?;
+
+    let creation = get_table_creation(temp_path(), table_name, None)?;
+    iceberg_catalog.create_table(&namespace, creation).await?;
+
+    let write_ctx = SessionContext::new_with_config(
+        SessionConfig::new().with_target_partitions(data_file_count),
+    );
+    let arrow_schema = Arc::new(ArrowSchema::new(vec![
+        Field::new("foo1", DataType::Int32, false),
+        Field::new("foo2", DataType::Utf8, false),
+    ]));
+
+    let batches: Vec<RecordBatch> = (1..=data_file_count as i32)
+        .map(|idx| {
+            RecordBatch::try_new(arrow_schema.clone(), vec![
+                Arc::new(Int32Array::from(vec![idx])) as ArrayRef,
+                Arc::new(StringArray::from(vec![format!("row-{idx}")])) as 
ArrayRef,
+            ])
+        })
+        .collect::<std::result::Result<_, _>>()?;
+
+    let partitions = batches.into_iter().map(|batch| vec![batch]).collect();
+    let source_table = Arc::new(MemTable::try_new(arrow_schema, 
partitions).unwrap());
+    write_ctx
+        .register_table("source_table", source_table)
+        .unwrap();
+
+    let catalog = 
Arc::new(IcebergCatalogProvider::try_new(iceberg_catalog.clone()).await?);
+    write_ctx.register_catalog("catalog", catalog);
+
+    let insert_sql =
+        format!("INSERT INTO catalog.{namespace_name}.{table_name} SELECT * 
FROM source_table");
+    let batches = write_ctx
+        .sql(&insert_sql)
+        .await
+        .unwrap()
+        .collect()
+        .await
+        .unwrap();
+    assert_eq!(batches.len(), 1);
+
+    let rows_inserted = batches[0]
+        .column(0)
+        .as_any()
+        .downcast_ref::<UInt64Array>()
+        .unwrap();
+    assert_eq!(rows_inserted.value(0), data_file_count as u64);
+
+    Ok((iceberg_catalog, namespace, table_name.to_string()))
+}
+
+async fn get_multi_row_group_table_context(
+    namespace_name: &str,
+    table_name: &str,
+) -> Result<(Arc<MemoryCatalog>, NamespaceIdent, String)> {
+    let iceberg_catalog = Arc::new(get_iceberg_catalog().await);
+    let namespace = NamespaceIdent::new(namespace_name.to_string());
+    set_test_namespace(&iceberg_catalog, &namespace).await?;
+
+    let creation = get_table_creation(temp_path(), table_name, None)?;
+    let table = iceberg_catalog.create_table(&namespace, creation).await?;
+    let arrow_schema: Arc<ArrowSchema> = Arc::new(
+        table
+            .metadata()
+            .current_schema()
+            .as_ref()
+            .try_into()
+            .unwrap(),
+    );
+
+    let file_rows = [
+        vec![
+            (1, "row-1"),
+            (2, "row-2"),
+            (100, "row-100"),
+            (101, "row-101"),
+        ],
+        vec![
+            (50, "row-50"),
+            (150, "row-150"),
+            (151, "row-151"),
+            (152, "row-152"),
+        ],
+        vec![
+            (99, "row-99"),
+            (1000, "row-1000"),
+            (1001, "row-1001"),
+            (1002, "row-1002"),
+        ],
+    ];
+
+    let location_generator = 
DefaultLocationGenerator::new(table.metadata()).unwrap();
+    let writer_properties = WriterProperties::builder()
+        .set_max_row_group_row_count(Some(2))
+        .build();
+    let mut data_files = Vec::with_capacity(file_rows.len());
+
+    for (file_idx, rows) in file_rows.into_iter().enumerate() {
+        let parquet_writer_builder = ParquetWriterBuilder::new(
+            writer_properties.clone(),
+            table.metadata().current_schema().clone(),
+        );
+        let file_name_generator = DefaultFileNameGenerator::new(
+            format!("multi-row-group-{file_idx}"),
+            None,
+            iceberg::spec::DataFileFormat::Parquet,
+        );
+        let rolling_file_writer_builder = 
RollingFileWriterBuilder::new_with_default_file_size(
+            parquet_writer_builder,
+            table.file_io().clone(),
+            location_generator.clone(),
+            file_name_generator,
+        );
+        let data_file_writer_builder = 
DataFileWriterBuilder::new(rolling_file_writer_builder);
+        let mut data_file_writer = data_file_writer_builder.build(None).await?;
+        let batch = RecordBatch::try_new(arrow_schema.clone(), vec![
+            Arc::new(Int32Array::from(
+                rows.iter().map(|(foo1, _)| *foo1).collect::<Vec<_>>(),
+            )) as ArrayRef,
+            Arc::new(StringArray::from(
+                rows.iter().map(|(_, foo2)| *foo2).collect::<Vec<_>>(),
+            )) as ArrayRef,
+        ])?;
+
+        data_file_writer.write(batch).await?;
+        let file_data_files = data_file_writer.close().await?;
+        assert_eq!(file_data_files.len(), 1);
+        data_files.extend(file_data_files);
+    }
+
+    for data_file in &data_files {
+        let file_path = data_file
+            .file_path()
+            .strip_prefix("file://")
+            .unwrap_or(data_file.file_path());
+        let reader = SerializedFileReader::new(File::open(file_path)?)?;
+        assert!(
+            reader.metadata().num_row_groups() > 1,
+            "expected multiple row groups for {}",
+            data_file.file_path()
+        );
+    }
+
+    let tx = Transaction::new(&table);
+    let action = tx.fast_append().add_data_files(data_files);
+    let tx = action.apply(tx)?;
+    tx.commit(iceberg_catalog.as_ref()).await?;
+
+    Ok((iceberg_catalog, namespace, table_name.to_string()))
+}
+
+async fn get_read_context(
+    catalog: Arc<MemoryCatalog>,
+    target_partitions: usize,
+    enable_eager_scan_planning: Option<bool>,
+) -> Result<SessionContext> {
+    let ctx = SessionContext::new_with_config(read_session_config(
+        target_partitions,
+        enable_eager_scan_planning,
+    ));
+    let catalog = Arc::new(IcebergCatalogProvider::try_new(catalog).await?);
+    ctx.register_catalog("catalog", catalog);
+    Ok(ctx)
+}
+
+fn read_session_config(
+    target_partitions: usize,
+    enable_eager_scan_planning: Option<bool>,
+) -> SessionConfig {
+    let config = 
SessionConfig::new().with_target_partitions(target_partitions);
+
+    match enable_eager_scan_planning {
+        Some(enabled) => {
+            let mut iceberg_config = IcebergDataFusionConfig::default();
+            iceberg_config.enable_eager_scan_planning = enabled;
+            config.with_option_extension(iceberg_config)
+        }
+        None => config,
+    }
+}
+
+async fn scan_plan(
+    ctx: &SessionContext,
+    namespace: &NamespaceIdent,
+    table_name: &str,
+) -> Arc<dyn ExecutionPlan> {
+    let provider = ctx.catalog("catalog").unwrap();
+    let namespace_name = &namespace[0];
+    let schema = provider.schema(namespace_name).unwrap();
+    let table = schema.table(table_name).await.unwrap().unwrap();
+
+    let state = ctx.state();
+    table.scan(&state, None, &[], None).await.unwrap()
+}
+
+async fn scan_partition_count(
+    ctx: &SessionContext,
+    namespace: &NamespaceIdent,
+    table_name: &str,
+) -> usize {
+    let plan = scan_plan(ctx, namespace, table_name).await;
+    plan.downcast_ref::<IcebergTableScan>()
+        .expect("Expected IcebergTableScan");
+    plan.properties().output_partitioning().partition_count()
+}
+
+fn find_iceberg_scan(plan: &dyn ExecutionPlan) -> Option<&IcebergTableScan> {
+    if let Some(scan) = plan.downcast_ref::<IcebergTableScan>() {
+        return Some(scan);
+    }
+
+    plan.children()
+        .into_iter()
+        .find_map(|child| find_iceberg_scan(child.as_ref()))
+}
+
+#[tokio::test]
+async fn test_multi_file_scan_produces_multiple_partitions() -> Result<()> {
+    let data_file_count = 3;
+    // Ask for more partitions than files to verify scan planning does not 
expose empty partitions.
+    let target_partitions = data_file_count + 1;
+    let (iceberg_catalog, namespace, table_name) = 
get_multi_file_table_context(
+        "test_multi_file_scan_partitions",
+        "my_table",
+        data_file_count,
+    )
+    .await?;
+    let ctx = get_read_context(iceberg_catalog, target_partitions, 
Some(true)).await?;
+    let plan = scan_plan(&ctx, &namespace, &table_name).await;
+    plan.downcast_ref::<IcebergTableScan>()
+        .expect("Expected IcebergTableScan");
+
+    let actual_partition_count = 
plan.properties().output_partitioning().partition_count();
+    assert_eq!(actual_partition_count, data_file_count);
+
+    // Pins the eager plan's task-group accounting. The reader settings it 
derives from
+    // the retained TableScan are not introspectable, so their equivalence 
with the lazy
+    // path is covered behaviorally by 
test_multi_partition_scan_matches_single_partition_results.
+    let display = datafusion::physical_plan::displayable(plan.as_ref())
+        .one_line()
+        .to_string();
+    assert!(
+        display.contains("task_groups:[3] tasks:[3]"),
+        "unexpected eager scan display: {display}"
+    );
+
+    Ok(())
+}
+
+#[tokio::test]
+async fn test_multi_file_scan_defaults_to_single_lazy_partition() -> 
Result<()> {
+    let data_file_count = 3;
+    let target_partitions = data_file_count + 1;
+    let (iceberg_catalog, namespace, table_name) = 
get_multi_file_table_context(
+        "test_multi_file_scan_default_lazy",
+        "my_table",
+        data_file_count,
+    )
+    .await?;
+    let ctx = get_read_context(iceberg_catalog, target_partitions, 
None).await?;
+
+    let actual_partition_count = scan_partition_count(&ctx, &namespace, 
&table_name).await;
+
+    assert_eq!(actual_partition_count, 1);
+
+    Ok(())
+}
+
+#[tokio::test]
+async fn test_set_enable_eager_scan_planning() -> Result<()> {
+    let data_file_count = 3;
+    let target_partitions = data_file_count + 1;
+    let (iceberg_catalog, namespace, table_name) =
+        get_multi_file_table_context("test_set_eager_scan_planning", 
"my_table", data_file_count)
+            .await?;
+    let ctx = get_read_context(iceberg_catalog, target_partitions, 
Some(false)).await?;
+
+    ctx.sql("SET iceberg.enable_eager_scan_planning = true")
+        .await
+        .unwrap()
+        .collect()
+        .await
+        .unwrap();
+
+    let actual_partition_count = scan_partition_count(&ctx, &namespace, 
&table_name).await;
+
+    assert_eq!(actual_partition_count, data_file_count);
+
+    Ok(())
+}
+
+#[tokio::test]
+async fn test_multi_partition_scan_enforces_global_limit() -> Result<()> {
+    let data_file_count = 3;
+    let limit = 2;
+    let target_partitions = data_file_count + 1;
+    let (iceberg_catalog, namespace, table_name) = 
get_multi_file_table_context(
+        "test_multi_partition_scan_limit",
+        "my_table",
+        data_file_count,
+    )
+    .await?;
+    let ctx = get_read_context(iceberg_catalog, target_partitions, 
Some(true)).await?;
+    let namespace_name = &namespace[0];
+    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();
+
+    // The physical optimizer absorbs the initial GlobalLimitExec into
+    // CoalescePartitionsExec; its fetch enforces the global bound in the 
final plan.
+    // After a DataFusion bump, failure here likely means the plan shape 
changed, not eager scanning.
+    let global_limit_coalescer = plan
+        .downcast_ref::<CoalescePartitionsExec>()
+        .expect("Expected a globally limited CoalescePartitionsExec");
+    assert_eq!(global_limit_coalescer.fetch(), Some(limit));
+
+    let scan = find_iceberg_scan(global_limit_coalescer.input().as_ref())
+        .expect("Expected IcebergTableScan below the global limit");
+    assert_eq!(scan.limit(), Some(limit));
+    assert_eq!(
+        scan.properties().output_partitioning().partition_count(),
+        data_file_count
+    );
+
+    let batches = dataframe.collect().await.unwrap();
+    let mut seen_ids = std::collections::HashSet::new();
+    for batch in &batches {
+        let foo1 = batch
+            .column(0)
+            .as_any()
+            .downcast_ref::<Int32Array>()
+            .unwrap();
+        let foo2 = batch
+            .column(1)
+            .as_any()
+            .downcast_ref::<StringArray>()
+            .unwrap();
+
+        for row in 0..batch.num_rows() {
+            let id = foo1.value(row);
+            assert!((1..=data_file_count as i32).contains(&id));
+            assert_eq!(foo2.value(row), format!("row-{id}"));
+            assert!(seen_ids.insert(id), "duplicate row {id}");
+        }
+    }
+
+    assert_eq!(
+        batches.iter().map(RecordBatch::num_rows).sum::<usize>(),
+        limit
+    );
+    assert_eq!(seen_ids.len(), limit);
+
+    Ok(())
+}
+
+#[tokio::test]
+async fn test_multi_partition_ordered_scan_enforces_global_limit() -> 
Result<()> {
+    let data_file_count = 3;
+    let limit = 2;
+    let target_partitions = data_file_count + 1;
+    let (iceberg_catalog, namespace, table_name) = 
get_multi_file_table_context(
+        "test_multi_partition_ordered_scan_limit",
+        "my_table",
+        data_file_count,
+    )
+    .await?;
+    let ctx = get_read_context(iceberg_catalog, target_partitions, 
Some(true)).await?;
+    let namespace_name = &namespace[0];
+    let query = format!(
+        "SELECT foo1, foo2 FROM catalog.{namespace_name}.{table_name} ORDER BY 
foo1 LIMIT {limit}"
+    );
+    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");
+    assert_eq!(
+        scan.properties().output_partitioning().partition_count(),
+        data_file_count
+    );
+
+    let batches = dataframe.collect().await.unwrap();
+    let rows = batches
+        .iter()
+        .flat_map(|batch| {
+            let foo1 = batch
+                .column(0)
+                .as_any()
+                .downcast_ref::<Int32Array>()
+                .unwrap();
+            let foo2 = batch
+                .column(1)
+                .as_any()
+                .downcast_ref::<StringArray>()
+                .unwrap();
+
+            (0..batch.num_rows()).map(|row| (foo1.value(row), 
foo2.value(row).to_string()))
+        })
+        .collect::<Vec<_>>();
+
+    assert_eq!(rows, vec![
+        (1, "row-1".to_string()),
+        (2, "row-2".to_string())
+    ]);
+
+    Ok(())
+}
+
+#[tokio::test]
+async fn test_multi_partition_scan_matches_single_partition_results() -> 
Result<()> {

Review Comment:
   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 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