mbutrovich commented on code in PR #20:
URL: https://github.com/apache/datafusion-iceberg/pull/20#discussion_r4133959675
##########
crates/datafusion/tests/integration_datafusion_test.rs:
##########
@@ -977,3 +987,150 @@ async fn test_insert_into_partitioned() -> Result<(),
Box<dyn Error>> {
Ok(())
}
+
+/// Executes the single partition of `plan` and renders its rows as a table.
+async fn run(
+ plan: &dyn ExecutionPlan,
+ ctx: &SessionContext,
+) -> Result<String, Box<dyn Error>> {
+ let stream = plan.execute(0, ctx.task_ctx())?;
+ let batches = datafusion::physical_plan::common::collect(stream).await?;
+ Ok(pretty_format_batches(&batches)?.to_string())
+}
+
+/// Returns the first node of type `T` in `plan`, depth first.
+fn find_node<T: ExecutionPlan + 'static>(plan: &Arc<dyn ExecutionPlan>) ->
Option<&T> {
+ plan.downcast_ref::<T>()
+ .or_else(|| plan.children().into_iter().find_map(find_node::<T>))
+}
+
+/// The plan nodes and providers can be named and inspected from outside this
+/// crate, and rebuilt from their parts, as a codec that serializes them does.
+#[tokio::test]
+async fn test_plan_nodes_are_inspectable() -> Result<(), Box<dyn Error>> {
+ let iceberg_catalog = get_iceberg_catalog().await;
+ let namespace = NamespaceIdent::new("test_plan_nodes".to_string());
+ set_test_namespace(&iceberg_catalog, &namespace).await?;
+ let creation = get_table_creation(temp_path(), "my_table", None)?;
+ iceberg_catalog.create_table(&namespace, creation).await?;
+ let ident = TableIdent::new(namespace.clone(), "my_table".to_string());
+ let client: Arc<dyn Catalog> = Arc::new(iceberg_catalog);
+
+ let ctx = SessionContext::new();
+ let catalog = IcebergCatalogProvider::try_new(client.clone()).await?;
+ ctx.register_catalog("catalog", Arc::new(catalog));
+ let provider = ctx
+ .table_provider("catalog.test_plan_nodes.my_table")
+ .await?;
+ let provider = provider
+ .downcast_ref::<IcebergTableProvider>()
+ .expect("a catalog-backed provider");
+ assert_eq!(provider.table_ident(), &ident);
+ assert!(Arc::ptr_eq(provider.catalog(), &client));
+ let rebuilt = IcebergTableProvider::try_new(
+ provider.catalog().clone(),
+ provider.table_ident().namespace().clone(),
+ provider.table_ident().name(),
+ )
+ .await?;
+ assert_eq!(rebuilt.table_ident(), &ident);
+ assert_eq!(rebuilt.schema(), provider.schema());
+
+ // Write path: a commit above a write, both holding the table, and the
+ // commit going through the provider's catalog. The plan that runs is
+ // rebuilt from their accessors and children alone. The optimizer drops
+ // the coalesce above a single-partition write, so the rebuilt commit
+ // always gets one, as a codec would.
+ let insert = ctx
+ .sql("INSERT INTO catalog.test_plan_nodes.my_table VALUES (1, 'alan'),
(2, 'turing')")
+ .await?
+ .create_physical_plan()
+ .await?;
+ let commit = insert
+ .downcast_ref::<IcebergCommitExec>()
+ .expect("the insert plan is rooted at a commit");
+ assert_eq!(commit.table().identifier(), &ident);
+ assert!(Arc::ptr_eq(commit.catalog(), &client));
+ let write = find_node::<IcebergWriteExec>(&insert).expect("a write below
the commit");
+ assert_eq!(write.table().identifier(), &ident);
+ let rebuilt_write: Arc<dyn ExecutionPlan> = Arc::new(IcebergWriteExec::new(
+ write.table().clone(),
+ write.children()[0].clone(),
+ ));
+ let rebuilt_commit = IcebergCommitExec::new(
+ commit.table().clone(),
+ commit.catalog().clone(),
+ Arc::new(CoalescePartitionsExec::new(rebuilt_write)),
+ );
+ let inserted = run(&rebuilt_commit, &ctx).await?;
+ assert!(inserted.contains("| 2 |"), "{inserted}");
+
+ // Read path: a scan pinned to a snapshot, rebuilt from its accessors,
+ // returns the same rows.
+ let table = client.load_table(&ident).await?;
+ let snapshot_id = table.metadata().current_snapshot_id().unwrap();
+ let pinned =
+ IcebergStaticTableProvider::try_new_from_table_snapshot(table,
snapshot_id)
+ .await?;
+ assert_eq!(pinned.snapshot_id(), Some(snapshot_id));
+ ctx.register_table("pinned", Arc::new(pinned.clone()))?;
+ let plan = ctx
+ .sql("SELECT foo2 FROM pinned WHERE foo1 = 1")
+ .await?
+ .create_physical_plan()
+ .await?;
+ let scan = find_node::<IcebergTableScan>(&plan).expect("a scan");
+ assert!(scan.predicates().is_some(), "the filter is pushed down");
+ let rebuilt = IcebergTableScan::new_with_predicate(
+ scan.table().clone(),
+ scan.snapshot_id(),
+ scan.schema(),
+ scan.projection().map(<[String]>::to_vec),
+ scan.predicates().cloned(),
+ scan.limit(),
+ )?;
+ assert_eq!(rebuilt.schema(), scan.schema());
+ let expected = run(scan, &ctx).await?;
+ assert!(
+ expected.contains("alan") && !expected.contains("turing"),
+ "{expected}"
+ );
+ assert_eq!(run(&rebuilt, &ctx).await?, expected);
+
+ // Without a projection the scan reads every column, and its limit is kept.
+ let plan = pinned.scan(&ctx.state(), None, &[], Some(1)).await?;
+ let scan = plan.downcast_ref::<IcebergTableScan>().expect("a scan");
+ assert_eq!(scan.projection(), None);
+ let rebuilt = IcebergTableScan::new_with_predicate(
+ scan.table().clone(),
+ scan.snapshot_id(),
+ scan.schema(),
+ scan.projection().map(<[String]>::to_vec),
+ scan.predicates().cloned(),
+ scan.limit(),
+ )?;
+ assert_eq!(rebuilt.limit(), Some(1));
+ let expected = run(scan, &ctx).await?;
+ assert_eq!(expected.lines().count(), 5, "one row:\n{expected}");
Review Comment:
Comparing each rebuilt plan's output to the original's checks the round
trip. The checks on what the originals return are looser: `inserted.contains("|
2 |")`, `contains("alan")`, and here a line count of 5, which doesn't show
that both columns come back without a projection. The rest of this file asserts
results with `expect!`. Could these assert the full pretty-printed tables with
`expect![[...]].assert_eq(&expected)` instead? Then the test states the rows
each plan returns, and the equality with the rebuilt plan covers the rest.
##########
crates/datafusion/src/physical_plan/scan.rs:
##########
@@ -64,23 +65,141 @@ impl IcebergTableScan {
projection: Option<&Vec<usize>>,
filters: &[Expr],
limit: Option<usize>,
- ) -> Self {
- let output_schema = match projection {
- None => schema.clone(),
- Some(projection) => Arc::new(schema.project(projection).unwrap()),
+ ) -> Result<Self> {
+ let (output_schema, projection) = match projection {
+ None => (schema, None),
+ Some(projection) => {
+ let output_schema = Arc::new(schema.project(projection)?);
+ let names = output_schema
+ .fields()
+ .iter()
+ .map(|field| field.name().clone())
+ .collect();
+ (output_schema, Some(names))
+ }
};
- let plan_properties = Self::compute_properties(output_schema.clone());
- let projection = get_column_names(schema.clone(), projection);
- let predicates = convert_filters_to_predicate(filters);
+ Self::new_with_predicate(
+ table,
+ snapshot_id,
+ output_schema,
+ projection,
+ 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.
+ ///
+ /// Each argument takes what the matching accessor returns (`predicate`
+ /// what [`Self::predicates`] does), so a scan is rebuilt from
+ /// [`ExecutionPlan::schema`] and [`Self::projection`]:
+ ///
+ /// - `snapshot_id`: the snapshot to read, or `None` for the table's
current
+ /// snapshot.
+ /// - `schema`: the Arrow schema the scan outputs. With `projection` set,
+ /// its fields must be the projected columns, in order. Without one, it
+ /// must be the full schema of the table as read, which is not checked.
+ /// - `projection`: the names of the columns to read from the table, or
+ /// `None` for all of them.
+ /// - `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` does not name the fields of `schema`,
+ /// in order.
+ ///
+ /// # 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 accessors alone.
+ /// let rebuilt = IcebergTableScan::new_with_predicate(
+ /// scan.table().clone(),
+ /// scan.snapshot_id(),
+ /// scan.schema(),
+ /// scan.projection().map(<[String]>::to_vec),
+ /// scan.predicates().cloned(),
+ /// scan.limit(),
+ /// )?;
+ /// assert_eq!(rebuilt.schema(), scan.schema());
+ /// assert_eq!(rebuilt.projection(), scan.projection());
+ /// 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<Vec<String>>,
+ predicate: Option<Predicate>,
+ limit: Option<usize>,
+ ) -> Result<Self> {
+ // The scan reads the named columns and reports `schema`, so the two
+ // must agree, or its batches would not match its schema.
+ if let Some(projection) = &projection {
+ let fields: Vec<&String> =
+ schema.fields().iter().map(|field| field.name()).collect();
+ if !projection.iter().eq(fields.iter().copied()) {
+ return plan_err!(
+ "IcebergTableScan projection {projection:?} does not match
the \
+ fields of its schema {fields:?}"
+ );
+ }
Review Comment:
This follows up on the signature I suggested last round. When `projection`
is `Some`, it has to equal `schema`'s field names, so it carries nothing that
`schema` doesn't. When it's `None`, `schema` has to be the full table schema,
and the doc says that isn't checked. On the head commit, passing a one-field
projected schema with `None` still builds a scan whose `schema()` has one field
while its batches have two columns (`foo1`, `foo2`).
What do you think about dropping the `projection` parameter and always
selecting `schema`'s field names? For a full schema, that reads the same
columns as `None`. In the iceberg-rust rev this repo uses, `select_all` and a
`select` of every top-level name resolve to the same [field
ids](https://github.com/apache/iceberg/blob/48330b8dacab6662242d252b39c8444190979bb2/format/spec.md?plain=1#L401)
in
[`collect_scan_field_ids`](https://github.com/apache/iceberg-rust/blob/665c64e48e8d33797ecb1a421f327edd9b024879/crates/iceberg/src/scan/mod.rs#L63-L99),
and I got identical output from both on the head commit. That removes the
unchecked precondition, and the mismatch error and
`test_scan_rejects_projection_not_matching_schema` go away with it.
The trade-off I see is that `projection()` would return every column name
for a scan built without a projection, where it returns `None` today, and the
display would list every column. Is anything relying on that `None`?
--
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]