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


##########
crates/datafusion/src/physical_plan/commit.rs:
##########
@@ -259,35 +262,39 @@ impl ExecutionPlan for IcebergCommitExec {
                         )
                     })?
                     .as_any()
-                    .downcast_ref::<StringArray>()
+                    .downcast_ref::<LargeBinaryArray>()
                     .ok_or_else(|| {
                         internal_datafusion_err!(
-                            "Expected 'data_files' column to be StringArray"
+                            "Expected 'data_files' column to be 
LargeBinaryArray"
                         )
                     })?;
 
-                // Deserialize all data files from the StringArray
-                let batch_files: Vec<DataFile> = files_array
-                    .into_iter()
-                    .flatten()
-                    .map(|f| -> Result<DataFile> {
-                        // Parse JSON to DataFileSerde and convert to DataFile
-                        deserialize_data_file_from_json(
-                            f,
-                            spec_id,
-                            &partition_type,
-                            &current_schema,
+                // Each value is a separate Avro container with one or more 
data files.
+                files_array.iter().try_for_each(|payload| {
+                    let mut payload = payload.ok_or_else(|| {
+                        DataFusionError::Execution(
+                            "Null Avro data file payload".to_string(),
                         )
-                        .map_err(to_datafusion_error)
-                    })
-                    .collect::<Result<_>>()?;
-
-                // add record_counts from the current batch to total record 
count
-                total_record_count +=
-                    batch_files.iter().map(|f| f.record_count()).sum::<u64>();
-
-                // Add all deserialized files to our collection
-                data_files.extend(batch_files);
+                    })?;
+                    if payload.is_empty() {
+                        return Err(DataFusionError::Execution(
+                            "Empty Avro data file payload".to_string(),
+                        ));
+                    }

Review Comment:
   Should these two errors use `internal_datafusion_err!` like the column 
checks above? `IcebergWriteExec` declares the column non-nullable and never 
emits an empty payload, so reaching either branch means the input broke that 
contract rather than the user doing something wrong.
   
   Is the empty check needed? For an empty payload, `read_data_files_from_avro` 
already returns `Failed to read header: failed to fill whole buffer`. If you'd 
like to keep it for the clearer message, a commit test with 
`MockWriteExec::new(vec![vec![]])` would cover it.



##########
crates/datafusion/tests/integration_datafusion_test.rs:
##########
@@ -960,14 +973,41 @@ async fn test_insert_into_partitioned() -> Result<(), 
Box<dyn Error>> {
               "novel",
               "textbook",
               "shirt",
+              "coat",
+            ],
+            float_partition: PrimitiveArray<Float32>
+            [
+              123.456,
+              123.456,
+              NaN,
+              inf,
+              -inf,
+              null,
             ]"#]],
         &[],
         Some("id"),
     );
 
-    // Verify that data files exist under correct partition paths
+    // Verify all five partition files and six rows were committed.
     let table_ident = TableIdent::new(namespace.clone(), 
"partitioned_table".to_string());
     let table = client.load_table(&table_ident).await?;
+
+    let snapshot = table.metadata().current_snapshot().unwrap();
+    let manifest_list = table.manifest_list_reader(snapshot).load().await?;
+    let mut file_count = 0;
+    let mut row_count = 0;
+    for entry in manifest_list.entries() {
+        let manifest = table.manifest_reader().read(entry).await?;
+        file_count += manifest.entries().len();
+        row_count += manifest
+            .entries()
+            .iter()
+            .map(|entry| entry.data_file().record_count())
+            .sum::<u64>();
+    }
+    assert_eq!(file_count, 5);
+    assert_eq!(row_count, 6);

Review Comment:
   Should we assert the committed partition values here? These checks also pass 
on the base commit once the insert gets past the JSON error. I replaced 
`123.456` with `1.5` on the base commit, and the test passed with the NaN and 
both infinity partitions committed as null. The NaN and `Infinity` rows still 
land in two separate files, so the file and row counts don't change.
   
   Here is one way to check the values. It needs `HashSet`, `Literal`, and 
`Struct` added to the imports. With it, the base commit (using `1.5`) fails 
with three distinct partitions instead of five.
   
   ```suggestion
       // Verify the committed files, rows, and partition values.
       let table_ident = TableIdent::new(namespace.clone(), 
"partitioned_table".to_string());
       let table = client.load_table(&table_ident).await?;
   
       let snapshot = table.metadata().current_snapshot().unwrap();
       let manifest_list = table.manifest_list_reader(snapshot).load().await?;
       let mut file_count = 0;
       let mut row_count = 0;
       let mut partitions = HashSet::new();
       for entry in manifest_list.entries() {
           let manifest = table.manifest_reader().read(entry).await?;
           for entry in manifest.entries() {
               file_count += 1;
               row_count += entry.data_file().record_count();
               partitions.insert(entry.data_file().partition().clone());
           }
       }
       assert_eq!(file_count, 5);
       assert_eq!(row_count, 6);
       let partition = |category: &str, value: Option<f32>| {
           Struct::from_iter([Some(Literal::string(category)), 
value.map(Literal::float)])
       };
       assert_eq!(
           partitions,
           HashSet::from([
               partition("electronics", Some(123.456)),
               partition("books", Some(f32::NAN)),
               partition("books", Some(f32::INFINITY)),
               partition("clothing", Some(f32::NEG_INFINITY)),
               partition("clothing", None),
           ])
       );
   ```
   
   What do you think about covering a `double` partition column as well? On the 
base commit, a `double` column with NaN and the infinities commits the same 
null partition values.



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