NoahKusaba commented on code in PR #20:
URL: https://github.com/apache/datafusion-iceberg/pull/20#discussion_r4132959353


##########
crates/datafusion/src/physical_plan/commit.rs:
##########
@@ -55,6 +59,13 @@ pub(crate) struct IcebergCommitExec {
 }
 
 impl IcebergCommitExec {
+    /// Commits the data files `input` produces to `table` through `catalog`.
+    ///
+    /// `input` must have a single partition, such as a
+    /// 
[`CoalescePartitionsExec`](datafusion::physical_plan::coalesce_partitions::CoalescePartitionsExec)
+    /// over an [`IcebergWriteExec`](super::IcebergWriteExec); executing the 
node
+    /// fails otherwise. `schema` is the table's Arrow schema, shown in the
+    /// verbose plan display.

Review Comment:
   Agreed, dropped it in 4640919. It was only read by the verbose display, so 
`new` is now `(table, catalog, input)`, and the verbose arm shows `table`'s 
current schema on one line. That's also the schema `execute` reads the data 
files against, so it's closer to what the node does than the provider's cached 
one.
   
   `test_plan_nodes_are_inspectable` now rebuilds the write and the commit from 
their accessors and children and executes the rebuilt plan instead of the 
original; the later reads check the rows landed. One thing that came up: the 
optimizer drops the coalesce when the write already has a single partition, so 
the planned insert is just commit over write. The rebuilt commit always gets 
its own `CoalescePartitionsExec`, as a codec would, and the partition guard 
covers the case where it doesn't.
   



##########
crates/datafusion/src/physical_plan/scan.rs:
##########
@@ -64,23 +64,134 @@ impl IcebergTableScan {
         projection: Option<&Vec<usize>>,
         filters: &[Expr],
         limit: Option<usize>,
-    ) -> Self {
+    ) -> Result<Self> {
+        Self::new_with_predicate(
+            table,
+            snapshot_id,
+            schema,
+            projection.map(Vec::as_slice),
+            convert_filters_to_predicate(filters),
+            limit,
+        )
+    }
+
+    /// Creates a scan of `table` from an already-converted Iceberg
+    /// [`Predicate`] rather than DataFusion filters, for rebuilding a scan 
from
+    /// its parts, such as after sending them to another process. A predicate
+    /// cannot be converted back to the filters it came from.
+    ///
+    /// The arguments mean what the matching accessors return:
+    ///
+    /// - `snapshot_id`: the snapshot to read, or `None` for the table's 
current
+    ///   snapshot.
+    /// - `schema`: the Arrow schema of the table the scan reads, as its
+    ///   provider reports it.
+    /// - `projection`: indices into `schema` of the columns to read, or `None`
+    ///   for all. The columns are read from the table by name.
+    /// - `predicate`: pushed down to Iceberg to skip data files and rows. The
+    ///   table providers report their filters as
+    ///   
[`Inexact`](datafusion::logical_expr::TableProviderFilterPushDown::Inexact),
+    ///   so DataFusion still applies them above the scan.
+    ///
+    /// # Errors
+    ///
+    /// Returns an error if `projection` holds an index outside `schema`.
+    ///
+    /// # Example
+    ///
+    /// ```
+    /// use std::collections::HashMap;
+    ///
+    /// use datafusion::catalog::TableProvider;
+    /// use datafusion::physical_plan::ExecutionPlan;
+    /// use datafusion::prelude::{SessionContext, col, lit};
+    /// use datafusion_iceberg::IcebergStaticTableProvider;
+    /// use datafusion_iceberg::physical_plan::IcebergTableScan;
+    /// use iceberg::memory::{MEMORY_CATALOG_WAREHOUSE, MemoryCatalogBuilder};
+    /// use iceberg::spec::{NestedField, PrimitiveType, Schema, Type};
+    /// use iceberg::{Catalog, CatalogBuilder, NamespaceIdent, TableCreation};
+    ///
+    /// # tokio::runtime::Runtime::new()?.block_on(async {
+    /// # let warehouse = tempfile::tempdir()?;
+    /// # let props = HashMap::from([(
+    /// #     MEMORY_CATALOG_WAREHOUSE.to_string(),
+    /// #     warehouse.path().display().to_string(),
+    /// # )]);
+    /// # let catalog = MemoryCatalogBuilder::default().load("memory", 
props).await?;
+    /// # let namespace = NamespaceIdent::new("ns".to_string());
+    /// # catalog.create_namespace(&namespace, HashMap::new()).await?;
+    /// # let schema = Schema::builder()
+    /// #     .with_fields(vec![
+    /// #         NestedField::required(1, "id", 
Type::Primitive(PrimitiveType::Int)).into(),
+    /// #         NestedField::optional(2, "name", 
Type::Primitive(PrimitiveType::String))
+    /// #             .into(),
+    /// #     ])
+    /// #     .build()?;
+    /// # let creation = 
TableCreation::builder().name("t".to_string()).schema(schema).build();
+    /// # let table = catalog.create_table(&namespace, creation).await?;
+    /// let provider = 
IcebergStaticTableProvider::try_new_from_table(table).await?;
+    /// let ctx = SessionContext::new();
+    /// let filters = [col("id").gt(lit(1))];
+    /// let plan = provider
+    ///     .scan(&ctx.state(), Some(&vec![1]), &filters, None)
+    ///     .await?;
+    /// let scan = plan.downcast_ref::<IcebergTableScan>().unwrap();
+    ///
+    /// // Rebuild an equivalent scan from the original's parts.
+    /// let schema = provider.schema();
+    /// let projection = scan
+    ///     .projection()
+    ///     .map(|names| {
+    ///         names
+    ///             .iter()
+    ///             .map(|name| schema.index_of(name))
+    ///             .collect::<Result<Vec<_>, _>>()
+    ///     })
+    ///     .transpose()?;
+    /// let rebuilt = IcebergTableScan::new_with_predicate(
+    ///     scan.table().clone(),
+    ///     scan.snapshot_id(),
+    ///     schema,
+    ///     projection.as_deref(),
+    ///     scan.predicates().cloned(),
+    ///     scan.limit(),
+    /// )?;
+    /// assert_eq!(rebuilt.schema(), scan.schema());
+    /// assert_eq!(rebuilt.predicates(), scan.predicates());
+    /// # Ok::<(), Box<dyn std::error::Error>>(())
+    /// # })?;
+    /// # Ok::<(), Box<dyn std::error::Error>>(())
+    /// ```
+    pub fn new_with_predicate(
+        table: Table,
+        snapshot_id: Option<i64>,
+        schema: ArrowSchemaRef,
+        projection: Option<&[usize]>,
+        predicate: Option<Predicate>,
+        limit: Option<usize>,
+    ) -> Result<Self> {

Review Comment:
   Good catch, and thanks for trying the rebuild without the provider. You're 
right that a decoder has no full schema to pass, and that `scan.schema()` with 
`None` silently reads every column. In 4640919 it takes the output schema and 
`Option<Vec<String>>`, as `schema()` and `projection()` return them, and 
returns an error when the names aren't the schema's fields in order. The 
crate-private `new` does the index-to-name projection itself.
   
   The doctest and `test_plan_nodes_are_inspectable` now rebuild from the 
scan's accessors alone, including a scan with no projection and a limit. With 
`None`, `schema` still has to be the full table schema; I've documented that 
rather than checked it, since checking would mean repeating each provider's 
schema logic, and a codec passing `scan.projection()` never hits it.
   



-- 
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