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


##########
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:
   The doc says the arguments mean what the matching accessors return, but I 
don't think a decoder can get `schema` or `projection` from the scan. No 
accessor returns the full table schema, `schema()` returns the projected one, 
and `projection()` returns column names while this takes indices. The doctest 
and `test_plan_nodes_are_inspectable` both get the full schema from 
`provider.schema()`, and on the executor side there's no provider to ask.
   
   Rebuilding the full schema from `table()` means repeating each provider's 
logic: `IcebergTableProvider` caches the current schema at `try_new` time 
([`table/mod.rs`](https://github.com/apache/datafusion-iceberg/blob/27d6c7cca7f4c2c14125f4da02f9799d213ec791/crates/datafusion/src/table/mod.rs#L94-L97)),
 while a pinned `IcebergStaticTableProvider` uses the snapshot's schema 
([`table/mod.rs`](https://github.com/apache/datafusion-iceberg/blob/27d6c7cca7f4c2c14125f4da02f9799d213ec791/crates/datafusion/src/table/mod.rs#L297-L301)).
   
   The obvious alternative builds a scan that returns the wrong columns. I 
planned `SELECT foo2 FROM t` against a catalog table, then called 
`new_with_predicate(scan.table().clone(), scan.snapshot_id(), scan.schema(), 
None, scan.predicates().cloned(), scan.limit())`. The rebuilt scan reports one 
field (`foo2`), but its batches have two columns (`foo1`, `foo2`), because 
`None` means `select_all`.
   
   Should this take what the accessors return instead: the projected schema 
plus `Option<Vec<String>>` column names, returning an error when the names 
don't match the schema's fields? Then a codec could round-trip the node with 
`scan.schema()` and `scan.projection()`, and the test could rebuild the scan 
without going through the provider.



##########
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:
   Since this makes `new` public, how should a codec supply `schema` when it 
rebuilds the node? There's an accessor for `table` and `catalog`, but not for 
`schema`. As far as I can tell, `schema` is only read by the verbose 
`DisplayAs` arm, and `insert_into` passes the provider's cached schema 
([`table/mod.rs`](https://github.com/apache/datafusion-iceberg/blob/27d6c7cca7f4c2c14125f4da02f9799d213ec791/crates/datafusion/src/table/mod.rs#L237-L242)).
 Should we drop the parameter and display the schema of `table` instead, or add 
a `schema()` accessor? Changing the signature is free now, and it won't be once 
the node is public.
   
   `test_plan_nodes_are_inspectable` checks the commit and write accessors 
([lines 
1045-1051](https://github.com/apache/datafusion-iceberg/blob/27d6c7cca7f4c2c14125f4da02f9799d213ec791/crates/datafusion/tests/integration_datafusion_test.rs#L1045-L1051))
 but doesn't rebuild either node. What do you think about rebuilding both from 
their accessors and `children()` and executing the rebuilt plan, like the scan 
and metadata scan cases do? That would have caught the missing `schema`.



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