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]