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


##########
crates/datafusion/tests/integration_datafusion_test.rs:
##########
@@ -977,3 +991,415 @@ async fn test_insert_into_partitioned() -> Result<(), 
Box<dyn Error>> {
 
     Ok(())
 }
+
+/// Executes `plan`, which must have a single partition, and returns its rows.
+async fn run_batches(
+    plan: &dyn ExecutionPlan,
+    ctx: &SessionContext,
+) -> Result<Vec<RecordBatch>, Box<dyn Error>> {
+    assert_eq!(plan.properties().partitioning.partition_count(), 1);
+    let stream = plan.execute(0, ctx.task_ctx())?;
+    Ok(datafusion::physical_plan::common::collect(stream).await?)

Review Comment:
   Could `collect` be imported at the top of the file instead of spelled out 
here? The commit tests import it the same way 
([`commit.rs`](https://github.com/apache/datafusion-iceberg/blob/342145911a849d4d62982b32b369d125506459e7/crates/datafusion/src/physical_plan/commit.rs#L334)).



##########
crates/datafusion/tests/integration_datafusion_test.rs:
##########
@@ -977,3 +991,415 @@ async fn test_insert_into_partitioned() -> Result<(), 
Box<dyn Error>> {
 
     Ok(())
 }
+
+/// Executes `plan`, which must have a single partition, and returns its rows.
+async fn run_batches(
+    plan: &dyn ExecutionPlan,
+    ctx: &SessionContext,
+) -> Result<Vec<RecordBatch>, Box<dyn Error>> {
+    assert_eq!(plan.properties().partitioning.partition_count(), 1);
+    let stream = plan.execute(0, ctx.task_ctx())?;
+    Ok(datafusion::physical_plan::common::collect(stream).await?)
+}
+
+/// Executes `plan`, which must have a single partition, and renders its rows
+/// as a table.
+async fn run(
+    plan: &dyn ExecutionPlan,
+    ctx: &SessionContext,
+) -> Result<String, Box<dyn Error>> {
+    Ok(pretty_format_batches(&run_batches(plan, ctx).await?)?.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)),
+    );
+    expect![[r#"
+        +-------+
+        | count |
+        +-------+
+        | 2     |
+        +-------+"#]]
+    .assert_eq(&run(&rebuilt_commit, &ctx).await?);
+
+    // 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()))?;
+    // A later write, so a scan reading the current snapshot rather than the
+    // pinned one would return its row too.
+    ctx.sql("INSERT INTO catalog.test_plan_nodes.my_table VALUES (3, 
'hopper')")
+        .await?
+        .collect()
+        .await?;
+    let latest_snapshot_id = client
+        .load_table(&ident)
+        .await?
+        .metadata()
+        .current_snapshot_id()
+        .unwrap();
+    assert_ne!(latest_snapshot_id, snapshot_id);
+    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_eq!(
+        scan.predicates().map(ToString::to_string).as_deref(),
+        Some("foo1 = 1")
+    );
+    let rebuilt = IcebergTableScan::new_with_predicate(
+        scan.table().clone(),
+        scan.snapshot_id(),
+        scan.schema(),
+        scan.predicates().cloned(),
+        scan.limit(),
+    );
+    assert_eq!(rebuilt.schema(), scan.schema());
+    assert_eq!(rebuilt.projection(), scan.projection());
+    let expected = run(scan, &ctx).await?;
+    expect![[r#"
+        +------+------+
+        | foo1 | foo2 |
+        +------+------+
+        | 1    | alan |
+        +------+------+"#]]
+    .assert_eq(&expected);
+    assert_eq!(run(&rebuilt, &ctx).await?, expected);
+
+    // Without a projection the scan reads every column by name, 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(),
+        Some(&["foo1".to_string(), "foo2".to_string()][..])
+    );
+    let rebuilt = IcebergTableScan::new_with_predicate(
+        scan.table().clone(),
+        scan.snapshot_id(),
+        scan.schema(),
+        scan.predicates().cloned(),
+        scan.limit(),
+    );
+    assert_eq!(rebuilt.limit(), Some(1));
+    let expected = run(scan, &ctx).await?;
+    expect![[r#"
+        +------+------+
+        | foo1 | foo2 |
+        +------+------+
+        | 1    | alan |
+        +------+------+"#]]
+    .assert_eq(&expected);
+    assert_eq!(run(&rebuilt, &ctx).await?, expected);
+
+    // A scan reads the columns of its schema and no others, so one built over
+    // part of the table returns only those columns, from the pinned snapshot.
+    let foo2_only = Arc::new(pinned.schema().project(&[1])?);
+    let partial = IcebergTableScan::new_with_predicate(
+        pinned.table().clone(),
+        Some(snapshot_id),
+        foo2_only,
+        None,
+        None,
+    );
+    expect![[r#"
+        +--------+
+        | foo2   |
+        +--------+
+        | alan   |
+        | turing |
+        +--------+"#]]
+    .assert_eq(&run(&partial, &ctx).await?);
+
+    // Metadata tables: a scan rebuilt from a metadata scan's parts reads the
+    // same rows.
+    let plan = ctx
+        .sql("SELECT * FROM catalog.test_plan_nodes.\"my_table$snapshots\"")
+        .await?
+        .create_physical_plan()
+        .await?;
+    let metadata_scan = find_node::<IcebergMetadataScan>(&plan).expect("a 
metadata scan");
+    let provider = metadata_scan.provider();
+    assert_eq!(provider.table().identifier(), &ident);
+    let rebuilt = IcebergMetadataScan::new(IcebergMetadataTableProvider::new(
+        provider.table().clone(),
+        provider.metadata_type().clone(),
+    ));
+    let batches = run_batches(metadata_scan, &ctx).await?;
+    // Snapshots come back in no set order, so put the first, which has no
+    // parent, first. Their ids, times and paths differ on every run, so each
+    // row is checked against its snapshot's metadata.
+    let batch = concat_batches(&batches[0].schema(), &batches)?;
+    let order = sort_to_indices(batch.column_by_name("parent_id").unwrap(), 
None, None)?;
+    let batch = take_record_batch(&batch, &order)?;
+    let column = |name: &str| batch.column_by_name(name).unwrap().clone();
+    let longs = |name: &str| -> Result<Vec<Option<i64>>, Box<dyn Error>> {
+        Ok(cast(&column(name), &DataType::Int64)?
+            .as_primitive::<Int64Type>()
+            .iter()
+            .collect())
+    };
+    let strings = |name: &str| -> Vec<Option<String>> {
+        column(name)
+            .as_string::<i32>()
+            .iter()
+            .map(|value| value.map(str::to_string))
+            .collect()
+    };
+    let metadata = provider.table().metadata();
+    let snapshots = [snapshot_id, latest_snapshot_id].map(|id| {
+        metadata
+            .snapshot_by_id(id)
+            .expect("a snapshot of the table")
+    });
+    assert_eq!(
+        batch
+            .schema()
+            .fields()
+            .iter()
+            .map(|field| field.name().as_str())
+            .collect::<Vec<_>>(),
+        [
+            "committed_at",
+            "snapshot_id",
+            "parent_id",
+            "operation",
+            "manifest_list",
+            "summary"
+        ]
+    );
+    assert_eq!(
+        longs("committed_at")?,
+        snapshots.map(|snapshot| Some(snapshot.timestamp_ms() * 1000))
+    );
+    assert_eq!(
+        longs("snapshot_id")?,
+        [Some(snapshot_id), Some(latest_snapshot_id)]
+    );
+    assert_eq!(longs("parent_id")?, [None, Some(snapshot_id)]);
+    assert_eq!(
+        strings("operation"),
+        [Some("append".to_string()), Some("append".to_string())]
+    );
+    assert_eq!(
+        strings("manifest_list"),
+        snapshots.map(|snapshot| Some(snapshot.manifest_list().to_string()))
+    );
+    // The summaries print their keys in no set order, so sort them.
+    let summaries = column("summary");
+    let summaries = summaries.as_map();
+    let summaries = (0..summaries.len())
+        .map(|row| {
+            let entries = summaries.value(row);
+            let keys = entries.column(0).as_string::<i32>();
+            let values = entries.column(1).as_string::<i32>();
+            (0..entries.len())
+                .map(|entry| format!("{}: {}", keys.value(entry), 
values.value(entry)))
+                .collect::<BTreeSet<_>>()
+        })
+        .collect::<Vec<_>>();
+    expect![[r#"
+        [
+            {
+                "added-data-files: 1",
+                "added-files-size: 913",
+                "added-records: 2",
+                "total-data-files: 1",
+                "total-delete-files: 0",
+                "total-equality-deletes: 0",
+                "total-files-size: 913",
+                "total-position-deletes: 0",
+                "total-records: 2",
+            },
+            {
+                "added-data-files: 1",
+                "added-files-size: 901",
+                "added-records: 1",
+                "total-data-files: 2",
+                "total-delete-files: 0",
+                "total-equality-deletes: 0",
+                "total-files-size: 1814",
+                "total-position-deletes: 0",
+                "total-records: 3",
+            },
+        ]
+    "#]]
+    .assert_debug_eq(&summaries);

Review Comment:
   Could the summaries be checked against each snapshot's metadata, like the 
other columns above? `added-files-size` and `total-files-size` are the sizes of 
the Parquet files the inserts wrote. As I read it, a dependency bump that 
changes the bytes the Parquet writer produces would fail this test with no 
change in behavior. This version passes on the head commit merged with current 
`main`:
   
   ```suggestion
       assert_eq!(
           summaries,
           snapshots.map(|snapshot| {
               snapshot
                   .summary()
                   .additional_properties
                   .iter()
                   .map(|(key, value)| format!("{key}: {value}"))
                   .collect::<BTreeSet<_>>()
           })
       );
   ```



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