This is an automated email from the ASF dual-hosted git repository.

JingsongLi pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/paimon-rust.git


The following commit(s) were added to refs/heads/main by this push:
     new c8e18ccb fix: write value stats for primary-key parquet files (#938)
c8e18ccb is described below

commit c8e18ccb91c5d1e31a3ee5531925762c9b0cdc1d
Author: Jingsong Lee <[email protected]>
AuthorDate: Thu Sep 24 20:04:56 2026 +0800

    fix: write value stats for primary-key parquet files (#938)
---
 crates/paimon/src/arrow/format/parquet.rs  |  58 +++++-
 crates/paimon/src/spec/core_options.rs     |  84 +++++++++
 crates/paimon/src/table/kv_file_writer.rs  | 284 +++++++++++++++++++++++++++--
 crates/paimon/src/table/pk_vector_scan.rs  |  32 ++--
 crates/paimon/src/table/table_write.rs     |  17 +-
 crates/paimon/tests/pk_value_stats_test.rs | 123 +++++++++++++
 6 files changed, 549 insertions(+), 49 deletions(-)

diff --git a/crates/paimon/src/arrow/format/parquet.rs 
b/crates/paimon/src/arrow/format/parquet.rs
index 2d203bba..b3b54358 100644
--- a/crates/paimon/src/arrow/format/parquet.rs
+++ b/crates/paimon/src/arrow/format/parquet.rs
@@ -1600,6 +1600,12 @@ fn build_row_group_column_indices(
 ) -> Vec<Option<usize>> {
     let mut by_root_name: HashMap<&str, Option<usize>> = HashMap::new();
     for (column_index, column) in columns.iter().enumerate() {
+        // Only a top-level, non-repeated leaf has row-level statistics for its
+        // logical field. Nested leaf null counts (e.g. payload.child) and
+        // repeated element counts cannot describe their parent column.
+        if column.column_path().parts().len() != 1 || 
column.column_descr().max_rep_level() != 0 {
+            continue;
+        }
         let Some(root_name) = column.column_path().parts().first() else {
             continue;
         };
@@ -1788,6 +1794,8 @@ fn supports_manifest_min_max(data_type: &DataType) -> 
bool {
             | DataType::BigInt(_)
             | DataType::Char(_)
             | DataType::VarChar(_)
+            | DataType::Binary(_)
+            | DataType::VarBinary(_)
             | DataType::Decimal(_)
             | DataType::Double(_)
             | DataType::Float(_)
@@ -1808,10 +1816,10 @@ fn apply_stats_mode(
         return (min_datum, max_datum);
     };
     match data_type {
-        DataType::Char(_) | DataType::VarChar(_) => {
-            let min = min_datum.map(|datum| truncate_string_min_datum(datum, 
length));
+        DataType::Char(_) | DataType::VarChar(_) | DataType::Binary(_) | 
DataType::VarBinary(_) => {
+            let min = min_datum.map(|datum| truncate_min_datum(datum, length));
             let max = match max_datum {
-                Some(datum) => match truncate_string_max_datum(datum, length) {
+                Some(datum) => match truncate_max_datum(datum, length) {
                     Some(max) => Some(max),
                     None => return (None, None),
                 },
@@ -1823,20 +1831,37 @@ fn apply_stats_mode(
     }
 }
 
-fn truncate_string_min_datum(datum: Datum, length: usize) -> Datum {
+fn truncate_min_datum(datum: Datum, length: usize) -> Datum {
     match datum {
         Datum::String(value) => Datum::String(truncate_string_min(&value, 
length)),
+        Datum::Bytes(value) => 
Datum::Bytes(value.into_iter().take(length).collect()),
         other => other,
     }
 }
 
-fn truncate_string_max_datum(datum: Datum, length: usize) -> Option<Datum> {
+fn truncate_max_datum(datum: Datum, length: usize) -> Option<Datum> {
     match datum {
         Datum::String(value) => truncate_string_max(&value, 
length).map(Datum::String),
+        Datum::Bytes(value) => truncate_binary_max(&value, 
length).map(Datum::Bytes),
         other => Some(other),
     }
 }
 
+fn truncate_binary_max(value: &[u8], length: usize) -> Option<Vec<u8>> {
+    if value.len() <= length {
+        return Some(value.to_vec());
+    }
+    let mut prefix = value[..length].to_vec();
+    for idx in (0..prefix.len()).rev() {
+        if prefix[idx] != u8::MAX {
+            prefix[idx] += 1;
+            prefix.truncate(idx + 1);
+            return Some(prefix);
+        }
+    }
+    None
+}
+
 fn truncate_string_min(value: &str, length: usize) -> String {
     value.chars().take(length).collect()
 }
@@ -2789,6 +2814,29 @@ mod tests {
         ]
     }
 
+    #[test]
+    fn test_truncate_binary_stats_matches_java_unsigned_upper_bound() {
+        use crate::spec::MetadataStatsMode;
+
+        let binary = 
DataType::VarBinary(crate::spec::VarBinaryType::new(8).unwrap());
+        let (min, max) = super::apply_stats_mode(
+            &binary,
+            MetadataStatsMode::Truncate(2),
+            Some(Datum::Bytes(vec![0x12, 0xff, 0x01])),
+            Some(Datum::Bytes(vec![0x12, 0xff, 0xfe])),
+        );
+        assert_eq!(min, Some(Datum::Bytes(vec![0x12, 0xff])));
+        assert_eq!(max, Some(Datum::Bytes(vec![0x13])));
+
+        let (min, max) = super::apply_stats_mode(
+            &binary,
+            MetadataStatsMode::Truncate(1),
+            Some(Datum::Bytes(vec![0xfe, 0x01])),
+            Some(Datum::Bytes(vec![0xff, 0x01])),
+        );
+        assert_eq!((min, max), (None, None));
+    }
+
     fn test_parquet_schema() -> SchemaDescriptor {
         SchemaDescriptor::new(Arc::new(
             parse_message_type(
diff --git a/crates/paimon/src/spec/core_options.rs 
b/crates/paimon/src/spec/core_options.rs
index 86c28c1e..f100c6c3 100644
--- a/crates/paimon/src/spec/core_options.rs
+++ b/crates/paimon/src/spec/core_options.rs
@@ -73,6 +73,7 @@ const CHANGELOG_FILE_FORMAT_OPTION: &str = 
"changelog-file.format";
 const CHANGELOG_FILE_COMPRESSION_OPTION: &str = "changelog-file.compression";
 const CHANGELOG_FILE_STATS_MODE_OPTION: &str = "changelog-file.stats-mode";
 const METADATA_STATS_MODE_OPTION: &str = "metadata.stats-mode";
+const METADATA_STATS_MODE_PER_LEVEL_OPTION: &str = 
"metadata.stats-mode.per.level";
 const METADATA_STATS_DENSE_STORE_OPTION: &str = "metadata.stats-dense-store";
 const METADATA_STATS_KEEP_FIRST_N_COLUMNS_OPTION: &str = 
"metadata.stats-keep-first-n-columns";
 const DEFAULT_METADATA_STATS_MODE: &str = "truncate(16)";
@@ -1412,6 +1413,55 @@ impl<'a> CoreOptions<'a> {
         MetadataStatsMode::parse(METADATA_STATS_MODE_OPTION, value)
     }
 
+    /// Match Java's PK file stats precedence: changelog override, level, 
table.
+    pub(crate) fn pk_file_metadata_stats_mode(
+        &self,
+        level: i32,
+        is_changelog: bool,
+    ) -> crate::Result<&str> {
+        let mut level_mode = None;
+        if let Some(raw) = 
self.options.get(METADATA_STATS_MODE_PER_LEVEL_OPTION) {
+            for entry in raw
+                .split(',')
+                .map(str::trim)
+                .filter(|entry| !entry.is_empty())
+            {
+                let (key, value) =
+                    entry
+                        .split_once(':')
+                        .ok_or_else(|| crate::Error::DataInvalid {
+                            message: format!(
+                                "Invalid 
{METADATA_STATS_MODE_PER_LEVEL_OPTION} entry: '{entry}'"
+                            ),
+                            source: None,
+                        })?;
+                let parsed_level =
+                    key.trim()
+                        .parse::<i32>()
+                        .map_err(|error| crate::Error::DataInvalid {
+                            message: format!(
+                                "Invalid level in 
{METADATA_STATS_MODE_PER_LEVEL_OPTION}: '{key}'"
+                            ),
+                            source: Some(Box::new(error)),
+                        })?;
+                if parsed_level == level {
+                    level_mode = Some(value.trim());
+                }
+            }
+        }
+        Ok(if is_changelog {
+            self.changelog_file_stats_mode().or(level_mode)
+        } else {
+            level_mode
+        }
+        .unwrap_or_else(|| {
+            self.options
+                .get(METADATA_STATS_MODE_OPTION)
+                .map(String::as_str)
+                .unwrap_or(DEFAULT_METADATA_STATS_MODE)
+        }))
+    }
+
     /// Number of leading columns whose stats should be kept.
     ///
     /// A negative value means the option is ignored, matching Java Paimon.
@@ -2902,6 +2952,40 @@ mod tests {
         );
     }
 
+    #[test]
+    fn test_pk_file_metadata_stats_mode_follows_java_precedence() {
+        let options = HashMap::from([
+            (METADATA_STATS_MODE_OPTION.to_string(), "none".to_string()),
+            (
+                METADATA_STATS_MODE_PER_LEVEL_OPTION.to_string(),
+                "0:counts,1:truncate(8)".to_string(),
+            ),
+            (
+                CHANGELOG_FILE_STATS_MODE_OPTION.to_string(),
+                "full".to_string(),
+            ),
+        ]);
+        let core = CoreOptions::new(&options);
+        assert_eq!(
+            core.pk_file_metadata_stats_mode(0, false).unwrap(),
+            "counts"
+        );
+        assert_eq!(
+            core.pk_file_metadata_stats_mode(1, false).unwrap(),
+            "truncate(8)"
+        );
+        assert_eq!(core.pk_file_metadata_stats_mode(2, false).unwrap(), 
"none");
+        assert_eq!(core.pk_file_metadata_stats_mode(0, true).unwrap(), "full");
+
+        let invalid = HashMap::from([(
+            METADATA_STATS_MODE_PER_LEVEL_OPTION.to_string(),
+            "x:counts".to_string(),
+        )]);
+        assert!(CoreOptions::new(&invalid)
+            .pk_file_metadata_stats_mode(0, false)
+            .is_err());
+    }
+
     #[test]
     fn test_metadata_stats_mode_rejects_invalid_values() {
         let options = HashMap::from([(
diff --git a/crates/paimon/src/table/kv_file_writer.rs 
b/crates/paimon/src/table/kv_file_writer.rs
index 4f2987bb..3dc42928 100644
--- a/crates/paimon/src/table/kv_file_writer.rs
+++ b/crates/paimon/src/table/kv_file_writer.rs
@@ -27,13 +27,15 @@
 //! Reference: 
[org.apache.paimon.io.KeyValueDataFileWriterImpl](https://github.com/apache/paimon/blob/release-1.3/paimon-core/src/main/java/org/apache/paimon/io/KeyValueDataFileWriterImpl.java)
 
 use crate::arrow::arrow_fields_to_paimon;
-use crate::arrow::format::{create_format_writer, with_write_resources};
+use crate::arrow::format::{
+    create_format_writer, parquet::ParquetFormatWriter, with_write_resources, 
FormatFileWriter,
+};
 use crate::io::FileIO;
 use crate::resource::{MemoryReservation, ResourceContext};
 use crate::spec::stats::{compute_column_stats, BinaryTableStats};
 use crate::spec::{
     bucket_path_under, extract_datum_from_arrow, AggregationConfig, 
BinaryRowBuilder, CoreOptions,
-    DataFileMeta, DataType, MergeEngine, PartialUpdateConfig, RowKind, 
EMPTY_SERIALIZED_ROW,
+    DataField, DataFileMeta, DataType, MergeEngine, PartialUpdateConfig, 
RowKind,
     SEQUENCE_NUMBER_FIELD_NAME, VALUE_KIND_FIELD_NAME,
 };
 use crate::table::prepared_files::PreparedFiles;
@@ -91,6 +93,8 @@ pub(crate) struct KeyValueWriteConfig {
     pub primary_key_indices: Vec<usize>,
     /// Paimon DataTypes for each primary key column (same order as 
primary_key_indices).
     pub primary_key_types: Vec<DataType>,
+    /// Logical value fields, in file order, for Parquet footer statistics.
+    pub value_fields: Vec<DataField>,
     /// Sequence field column indices in the user schema (empty if not 
configured).
     pub sequence_field_indices: Vec<usize>,
     /// Merge engine for deduplication.
@@ -99,6 +103,7 @@ pub(crate) struct KeyValueWriteConfig {
 }
 
 struct IndexedFileWrite<'a> {
+    is_changelog: bool,
     file_prefix: &'a str,
     file_ordinal: usize,
     file_format: &'a str,
@@ -325,6 +330,7 @@ impl KeyValueFileWriter {
                 data_seq.as_ref(),
                 &data_indices,
                 IndexedFileWrite {
+                    is_changelog: false,
                     file_prefix: &self.config.data_file_prefix,
                     file_ordinal: self.written_files.len(),
                     file_format: &self.config.file_format,
@@ -344,6 +350,7 @@ impl KeyValueFileWriter {
                     seq_array.as_ref(),
                     &sorted_indices,
                     IndexedFileWrite {
+                        is_changelog: true,
                         file_prefix: &self.config.changelog_file_prefix,
                         file_ordinal: self.written_changelog_files.len(),
                         file_format: &self.config.changelog_file_format,
@@ -433,16 +440,39 @@ impl KeyValueFileWriter {
         self.file_io.mkdirs(&format!("{bucket_dir}/")).await?;
         let file_path = format!("{bucket_dir}/{file_name}");
         let output = self.file_io.new_output(&file_path)?;
-        let writer = create_format_writer(
-            &output,
-            physical_schema.clone(),
-            write.file_compression,
-            self.config.file_compression_zstd_level,
-            None,
-            None,
-            None,
-        )
-        .await?;
+        // The physical KV file also contains sequence and row-kind columns. 
Give
+        // Parquet only the logical value fields so metadata stats and their 
dense
+        // column mapping follow Java's value schema (and its stats options).
+        // Keep the existing unshredded KV layout for this writer.
+        let writer: Box<dyn FormatFileWriter> = if 
write.file_format.eq_ignore_ascii_case("parquet")
+        {
+            let mut stats_options = self.config.table_options.clone();
+            let core_options = CoreOptions::new(&self.config.table_options);
+            let stats_mode = core_options.pk_file_metadata_stats_mode(0, 
write.is_changelog)?;
+            stats_options.insert("metadata.stats-mode".to_string(), 
stats_mode.to_string());
+            Box::new(
+                ParquetFormatWriter::new(
+                    &output,
+                    physical_schema.clone(),
+                    write.file_compression,
+                    self.config.file_compression_zstd_level,
+                    Some(&self.config.value_fields),
+                    &stats_options,
+                )
+                .await?,
+            )
+        } else {
+            create_format_writer(
+                &output,
+                physical_schema.clone(),
+                write.file_compression,
+                self.config.file_compression_zstd_level,
+                None,
+                None,
+                None,
+            )
+            .await?
+        };
         let mut writer = with_write_resources(writer, self.resources.as_ref());
 
         let vk_idx = batch
@@ -509,7 +539,12 @@ impl KeyValueFileWriter {
             }
         }
 
-        let file_size = writer.close().await?.file_size as i64;
+        let write_result = writer.close().await?;
+        let file_size = write_result.file_size as i64;
+        let (value_stats, value_stats_cols) = match write_result.value_stats {
+            Some(stats) => (stats.stats, stats.columns),
+            None => (BinaryTableStats::empty(), Some(Vec::new())),
+        };
 
         let key_columns: Vec<Arc<dyn Array>> = self
             .config
@@ -552,11 +587,7 @@ impl KeyValueFileWriter {
             min_key,
             max_key,
             key_stats,
-            value_stats: BinaryTableStats::new(
-                EMPTY_SERIALIZED_ROW.clone(),
-                EMPTY_SERIALIZED_ROW.clone(),
-                vec![],
-            ),
+            value_stats,
             min_sequence_number: write.min_sequence_number,
             max_sequence_number: write.max_sequence_number,
             schema_id: self.config.schema_id,
@@ -566,7 +597,7 @@ impl KeyValueFileWriter {
             delete_row_count: Some(write.delete_row_count),
             embedded_index: None,
             file_source: Some(0), // FileSource.APPEND
-            value_stats_cols: Some(vec![]),
+            value_stats_cols,
             external_path: None,
             first_row_id: None,
             write_cols: None,
@@ -1023,6 +1054,15 @@ mod tests {
             primary_keys: vec!["id".into()],
             primary_key_indices: vec![0],
             primary_key_types: vec![DataType::Int(IntType::new())],
+            value_fields: vec![
+                DataField::new(0, "id".into(), DataType::Int(IntType::new())),
+                DataField::new(
+                    1,
+                    "seq".into(),
+                    DataType::BigInt(crate::spec::BigIntType::new()),
+                ),
+                DataField::new(2, "value".into(), 
DataType::Int(IntType::new())),
+            ],
             sequence_field_indices: vec![1],
             merge_engine,
             deletion_vectors_enabled: false,
@@ -1038,6 +1078,212 @@ mod tests {
         .unwrap()
     }
 
+    #[tokio::test]
+    async fn test_pk_value_stats_use_emitted_rows_and_logical_columns() {
+        let schema = Arc::new(ArrowSchema::new(vec![
+            ArrowField::new("id", ArrowDataType::Int32, false),
+            ArrowField::new("seq", ArrowDataType::Int64, false),
+            ArrowField::new("value", ArrowDataType::Int32, true),
+        ]));
+        let batch = RecordBatch::try_new(
+            schema,
+            vec![
+                Arc::new(Int32Array::from(vec![1, 2, 1])),
+                Arc::new(Int64Array::from(vec![10, 20, 30])),
+                Arc::new(Int32Array::from(vec![Some(100), None, Some(300)])),
+            ],
+        )
+        .unwrap();
+        let mut config = test_write_config(MergeEngine::Deduplicate);
+        config.input_changelog = true;
+        config.write_buffer_size = i64::MAX;
+        config
+            .table_options
+            .insert("metadata.stats-dense-store".to_string(), 
"true".to_string());
+        let mut writer =
+            
KeyValueFileWriter::new(FileIOBuilder::new("memory").build().unwrap(), config, 
0)
+                .unwrap();
+        writer.write(&batch).await.unwrap();
+        let prepared = writer.prepare_commit().await.unwrap();
+
+        let data = &prepared.data_files[0];
+        assert_eq!(data.row_count, 2);
+        assert_eq!(data.value_stats_cols, None);
+        assert_eq!(
+            data.value_stats.null_counts(),
+            &vec![Some(0), Some(0), Some(1)]
+        );
+        let min =
+            
crate::spec::BinaryRow::from_serialized_bytes(data.value_stats.min_values()).unwrap();
+        let max =
+            
crate::spec::BinaryRow::from_serialized_bytes(data.value_stats.max_values()).unwrap();
+        assert_eq!(min.arity(), 3);
+        assert_eq!(min.get_int(0).unwrap(), 1);
+        assert_eq!(max.get_int(0).unwrap(), 2);
+        assert_eq!(min.get_long(1).unwrap(), 20);
+        assert_eq!(max.get_long(1).unwrap(), 30);
+        assert_eq!(min.get_int(2).unwrap(), 300);
+        assert_eq!(max.get_int(2).unwrap(), 300);
+
+        let changelog = &prepared.changelog_files[0];
+        assert_eq!(changelog.row_count, 3);
+        assert_eq!(
+            changelog.value_stats.null_counts(),
+            &vec![Some(0), Some(0), Some(1)]
+        );
+        let min = 
crate::spec::BinaryRow::from_serialized_bytes(changelog.value_stats.min_values())
+            .unwrap();
+        assert_eq!(min.get_int(2).unwrap(), 100);
+    }
+
+    #[tokio::test]
+    async fn test_pk_value_stats_respect_dense_column_modes() {
+        let schema = Arc::new(ArrowSchema::new(vec![
+            ArrowField::new("id", ArrowDataType::Int32, false),
+            ArrowField::new("seq", ArrowDataType::Int64, false),
+            ArrowField::new("value", ArrowDataType::Int32, true),
+        ]));
+        let batch = RecordBatch::try_new(
+            schema,
+            vec![
+                Arc::new(Int32Array::from(vec![1, 2])),
+                Arc::new(Int64Array::from(vec![10, 20])),
+                Arc::new(Int32Array::from(vec![Some(100), None])),
+            ],
+        )
+        .unwrap();
+        let mut config = test_write_config(MergeEngine::Deduplicate);
+        config.write_buffer_size = i64::MAX;
+        config.table_options.extend([
+            ("metadata.stats-mode".to_string(), "none".to_string()),
+            ("metadata.stats-dense-store".to_string(), "true".to_string()),
+            ("fields.id.stats-mode".to_string(), "full".to_string()),
+            ("fields.value.stats-mode".to_string(), "counts".to_string()),
+        ]);
+        let mut writer =
+            
KeyValueFileWriter::new(FileIOBuilder::new("memory").build().unwrap(), config, 
0)
+                .unwrap();
+        writer.write(&batch).await.unwrap();
+        let prepared = writer.prepare_commit().await.unwrap();
+        let file = &prepared.data_files[0];
+
+        assert_eq!(
+            file.value_stats_cols,
+            Some(vec!["id".into(), "value".into()])
+        );
+        assert_eq!(file.value_stats.null_counts(), &vec![Some(0), Some(1)]);
+        let min =
+            
crate::spec::BinaryRow::from_serialized_bytes(file.value_stats.min_values()).unwrap();
+        assert_eq!(min.arity(), 2);
+        assert_eq!(min.get_int(0).unwrap(), 1);
+        assert!(min.is_null_at(1));
+    }
+
+    #[tokio::test]
+    async fn test_pk_binary_value_stats_match_full_and_truncate_modes() {
+        use crate::spec::{BinaryType, VarBinaryType};
+        use arrow_array::BinaryArray;
+
+        let schema = Arc::new(ArrowSchema::new(vec![
+            ArrowField::new("id", ArrowDataType::Int32, false),
+            ArrowField::new("payload", ArrowDataType::Binary, false),
+            ArrowField::new("raw", ArrowDataType::Binary, false),
+        ]));
+        let batch = RecordBatch::try_new(
+            schema,
+            vec![
+                Arc::new(Int32Array::from(vec![1, 2])),
+                Arc::new(BinaryArray::from_iter_values([b"ab1", b"ac0"])),
+                Arc::new(BinaryArray::from_iter_values([
+                    b"\xfe\x01".as_slice(),
+                    b"\xff\x01".as_slice(),
+                ])),
+            ],
+        )
+        .unwrap();
+        let mut config = test_write_config(MergeEngine::Deduplicate);
+        config.value_fields = vec![
+            DataField::new(0, "id".into(), DataType::Int(IntType::new())),
+            DataField::new(
+                1,
+                "payload".into(),
+                DataType::Binary(BinaryType::new(4).unwrap()),
+            ),
+            DataField::new(
+                2,
+                "raw".into(),
+                DataType::VarBinary(VarBinaryType::new(8).unwrap()),
+            ),
+        ];
+        config.write_buffer_size = i64::MAX;
+        config.table_options.extend([
+            ("metadata.stats-mode".to_string(), "full".to_string()),
+            (
+                "fields.payload.stats-mode".to_string(),
+                "truncate(2)".to_string(),
+            ),
+        ]);
+        let mut writer =
+            
KeyValueFileWriter::new(FileIOBuilder::new("memory").build().unwrap(), config, 
0)
+                .unwrap();
+        writer.write(&batch).await.unwrap();
+        let prepared = writer.prepare_commit().await.unwrap();
+        let stats = &prepared.data_files[0].value_stats;
+        assert_eq!(stats.null_counts(), &vec![Some(0); 3]);
+        let min = 
crate::spec::BinaryRow::from_serialized_bytes(stats.min_values()).unwrap();
+        let max = 
crate::spec::BinaryRow::from_serialized_bytes(stats.max_values()).unwrap();
+        assert_eq!(min.get_binary(1).unwrap(), b"ab");
+        assert_eq!(max.get_binary(1).unwrap(), b"ad");
+        assert_eq!(min.get_binary(2).unwrap(), b"\xfe\x01");
+        assert_eq!(max.get_binary(2).unwrap(), b"\xff\x01");
+    }
+
+    #[tokio::test]
+    async fn test_pk_value_stats_choose_level_and_changelog_modes() {
+        let schema = Arc::new(ArrowSchema::new(vec![
+            ArrowField::new("id", ArrowDataType::Int32, false),
+            ArrowField::new("seq", ArrowDataType::Int64, false),
+            ArrowField::new("value", ArrowDataType::Int32, false),
+        ]));
+        let batch = RecordBatch::try_new(
+            schema,
+            vec![
+                Arc::new(Int32Array::from(vec![1, 2])),
+                Arc::new(Int64Array::from(vec![10, 20])),
+                Arc::new(Int32Array::from(vec![100, 200])),
+            ],
+        )
+        .unwrap();
+        let mut config = test_write_config(MergeEngine::Deduplicate);
+        config.input_changelog = true;
+        config.write_buffer_size = i64::MAX;
+        config.table_options.extend([
+            ("metadata.stats-mode".to_string(), "none".to_string()),
+            (
+                "metadata.stats-mode.per.level".to_string(),
+                "0:counts".to_string(),
+            ),
+            ("changelog-file.stats-mode".to_string(), "full".to_string()),
+        ]);
+        let mut writer =
+            
KeyValueFileWriter::new(FileIOBuilder::new("memory").build().unwrap(), config, 
0)
+                .unwrap();
+        writer.write(&batch).await.unwrap();
+        let prepared = writer.prepare_commit().await.unwrap();
+        let data = &prepared.data_files[0].value_stats;
+        assert_eq!(data.null_counts(), &vec![Some(0); 3]);
+        let data_min = 
crate::spec::BinaryRow::from_serialized_bytes(data.min_values()).unwrap();
+        assert!(data_min.is_null_at(0));
+        assert!(data_min.is_null_at(2));
+
+        let changelog = &prepared.changelog_files[0].value_stats;
+        assert_eq!(changelog.null_counts(), &vec![Some(0); 3]);
+        let changelog_min =
+            
crate::spec::BinaryRow::from_serialized_bytes(changelog.min_values()).unwrap();
+        assert_eq!(changelog_min.get_int(0).unwrap(), 1);
+        assert_eq!(changelog_min.get_int(2).unwrap(), 100);
+    }
+
     #[test]
     fn test_dedup_sorted_indices_keeps_first_row_for_first_row_engine() {
         let schema = Arc::new(ArrowSchema::new(vec![
diff --git a/crates/paimon/src/table/pk_vector_scan.rs 
b/crates/paimon/src/table/pk_vector_scan.rs
index 9b845af7..8b998bd8 100644
--- a/crates/paimon/src/table/pk_vector_scan.rs
+++ b/crates/paimon/src/table/pk_vector_scan.rs
@@ -983,7 +983,6 @@ mod tests {
         use super::*;
         use crate::catalog::Identifier;
         use crate::io::{FileIO, FileIOBuilder};
-        use crate::spec::stats::compute_column_stats;
         use crate::spec::{
             DataType, Datum, FloatType, IntType, PredicateBuilder, Schema, 
TableSchema, VectorType,
         };
@@ -1088,9 +1087,8 @@ mod tests {
 
         /// Build a real single-file primary-key table via the public write 
path.
         ///
-        /// The Rust key-value writer records primary-key stats in 
`key_stats`. For the
-        /// deletion-vector test, also populate `value_stats` so its non-key 
predicate
-        /// has the same metadata a Java primary-key writer produces.
+        /// The Rust key-value writer records both key and value stats. Under
+        /// deletion vectors, the non-key predicate uses the real value stats.
         async fn build_pruning_test_table(
             with_deletion_vectors: bool,
         ) -> (tempfile::TempDir, Table) {
@@ -1122,23 +1120,17 @@ mod tests {
             let bucket = written.bucket;
             let partition = written.partition.clone();
             assert_eq!(base_meta.key_stats.null_counts(), &vec![Some(0)]);
-            assert!(base_meta.value_stats.null_counts().is_empty());
-            assert_eq!(base_meta.value_stats_cols, Some(vec![]));
-
-            let file_meta = if with_deletion_vectors {
-                let int = DataType::Int(IntType::new());
-                let value_stats: BinaryTableStats =
-                    compute_column_stats(&batch, &[0, 1], &[int.clone(), 
int]).unwrap();
-                DataFileMeta {
-                    value_stats,
-                    value_stats_cols: Some(vec!["id".to_string(), 
"score".to_string()]),
-                    ..base_meta
-                }
-            } else {
-                base_meta
-            };
+            assert_eq!(base_meta.value_stats.null_counts(), &vec![Some(0); 2]);
+            assert_eq!(
+                base_meta.value_stats_cols,
+                Some(vec!["id".to_string(), "score".to_string()])
+            );
+            let min = 
BinaryRow::from_serialized_bytes(base_meta.value_stats.min_values()).unwrap();
+            let max = 
BinaryRow::from_serialized_bytes(base_meta.value_stats.max_values()).unwrap();
+            assert_eq!(min.get_int(1).unwrap(), 0);
+            assert_eq!(max.get_int(1).unwrap(), PRUNE_ROWS - 1);
 
-            let message = CommitMessage::new(partition, bucket, 
vec![file_meta]);
+            let message = CommitMessage::new(partition, bucket, 
vec![base_meta]);
             TableCommit::new(table.clone(), "pkvector-prune".to_string())
                 .commit(vec![message])
                 .await
diff --git a/crates/paimon/src/table/table_write.rs 
b/crates/paimon/src/table/table_write.rs
index 46eea6e2..1c3da1b9 100644
--- a/crates/paimon/src/table/table_write.rs
+++ b/crates/paimon/src/table/table_write.rs
@@ -1196,6 +1196,7 @@ impl TableWrite {
                     primary_keys: self.table.schema().primary_keys().to_vec(),
                     primary_key_indices: self.primary_key_indices.clone(),
                     primary_key_types: self.primary_key_types.clone(),
+                    value_fields: self.table.schema().fields().to_vec(),
                     sequence_field_indices: 
self.sequence_field_indices.clone(),
                     merge_engine: self.merge_engine,
                     deletion_vectors_enabled: 
CoreOptions::new(self.table.schema().options())
@@ -2008,7 +2009,7 @@ pub(in crate::table) mod tests {
     }
 
     #[tokio::test]
-    async fn 
test_append_write_truncates_string_value_stats_and_keeps_binary_counts() {
+    async fn test_append_write_truncates_string_and_binary_value_stats() {
         let file_io = test_file_io();
         let table_path = "memory:/test_table_write_skip_variable_length_stats";
         setup_dirs(&file_io, table_path).await;
@@ -2067,8 +2068,8 @@ pub(in crate::table) mod tests {
         assert_eq!(max_values.get_int(0).unwrap(), 2);
         assert_eq!(min_values.get_string(1).unwrap(), "a long string va");
         assert_eq!(max_values.get_string(1).unwrap(), "another long sts");
-        assert!(min_values.is_null_at(2));
-        assert!(max_values.is_null_at(2));
+        assert_eq!(min_values.get_binary(2).unwrap(), b"another-large-bi");
+        assert_eq!(max_values.get_binary(2).unwrap(), b"large-binary-vam");
     }
 
     #[tokio::test]
@@ -2142,8 +2143,8 @@ pub(in crate::table) mod tests {
         assert!(max_values.is_null_at(0));
         assert_eq!(min_values.get_string(1).unwrap(), 
"alpha-long-value-12345");
         assert_eq!(max_values.get_string(1).unwrap(), "zeta-long-value-99999");
-        assert!(min_values.is_null_at(2));
-        assert!(max_values.is_null_at(2));
+        assert_eq!(min_values.get_binary(2).unwrap(), b"first-binary-value");
+        assert_eq!(max_values.get_binary(2).unwrap(), b"first-binary-value");
     }
 
     #[tokio::test]
@@ -3913,6 +3914,12 @@ pub(in crate::table) mod tests {
         assert_eq!(file.level, 0);
         assert_eq!(file.min_sequence_number, 0);
         assert_eq!(file.max_sequence_number, 2);
+        assert_eq!(file.value_stats_cols, None);
+        assert_eq!(file.value_stats.null_counts(), &vec![Some(0), Some(0)]);
+        let min_values = 
BinaryRow::from_serialized_bytes(file.value_stats.min_values()).unwrap();
+        let max_values = 
BinaryRow::from_serialized_bytes(file.value_stats.max_values()).unwrap();
+        assert_eq!(min_values.get_int(1).unwrap(), 10);
+        assert_eq!(max_values.get_int(1).unwrap(), 30);
         // min_key and max_key should be non-empty (serialized BinaryRow)
         assert!(!file.min_key.is_empty());
         assert!(!file.max_key.is_empty());
diff --git a/crates/paimon/tests/pk_value_stats_test.rs 
b/crates/paimon/tests/pk_value_stats_test.rs
new file mode 100644
index 00000000..e30e8b81
--- /dev/null
+++ b/crates/paimon/tests/pk_value_stats_test.rs
@@ -0,0 +1,123 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+//   http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied.  See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+mod common;
+
+use arrow_array::{ArrayRef, Int32Array, ListArray, RecordBatch, StructArray};
+use arrow_buffer::{OffsetBuffer, ScalarBuffer};
+use arrow_schema::{DataType as ArrowDataType, Field as ArrowField, Schema as 
ArrowSchema};
+use common::incremental_helpers::{memory_table, persist_table_schema, 
setup_dirs};
+use paimon::spec::{
+    ArrayType, DataField, DataType, IntType, PredicateBuilder, RowType, 
Schema, TableSchema,
+};
+use std::sync::Arc;
+
+#[tokio::test]
+async fn nested_leaf_nulls_do_not_prune_non_null_parent_columns() {
+    let path = "memory:/pk_value_stats/nested_leaf_nulls";
+    let schema = Schema::builder()
+        .column("id", DataType::Int(IntType::new()))
+        .column(
+            "payload",
+            DataType::Row(RowType::new(vec![DataField::new(
+                2,
+                "child".to_string(),
+                DataType::Int(IntType::new()),
+            )])),
+        )
+        .column(
+            "items",
+            DataType::Array(ArrayType::new(DataType::Int(IntType::new()))),
+        )
+        .primary_key(["id"])
+        .option("bucket", "1")
+        .option("deletion-vectors.enabled", "true")
+        .option("metadata.stats-mode", "full")
+        .build()
+        .unwrap();
+    let (file_io, table) = memory_table(path, TableSchema::new(0, &schema));
+    setup_dirs(&file_io, path).await;
+    persist_table_schema(&file_io, path, table.schema()).await;
+
+    let child_field = Arc::new(ArrowField::new("child", ArrowDataType::Int32, 
true));
+    let payload: ArrayRef = Arc::new(StructArray::from(vec![(
+        child_field.clone(),
+        Arc::new(Int32Array::from(vec![None, None])) as ArrayRef,
+    )]));
+    let items: ArrayRef = Arc::new(ListArray::new(
+        Arc::new(ArrowField::new("element", ArrowDataType::Int32, true)),
+        OffsetBuffer::new(ScalarBuffer::from(vec![0, 1, 2])),
+        Arc::new(Int32Array::from(vec![None, None])),
+        None,
+    ));
+    let batch = RecordBatch::try_new(
+        Arc::new(ArrowSchema::new(vec![
+            ArrowField::new("id", ArrowDataType::Int32, false),
+            ArrowField::new(
+                "payload",
+                ArrowDataType::Struct(vec![child_field].into()),
+                true,
+            ),
+            ArrowField::new("items", items.data_type().clone(), true),
+        ])),
+        vec![Arc::new(Int32Array::from(vec![1, 2])), payload, items],
+    )
+    .unwrap();
+
+    let builder = table.new_write_builder();
+    let mut writer = builder.new_write().unwrap();
+    writer.write_arrow_batch(&batch).await.unwrap();
+    let messages = writer.prepare_commit().await.unwrap();
+    assert_eq!(messages.len(), 1);
+    let file = &messages[0].new_files[0];
+    assert_eq!(file.value_stats_cols, Some(vec!["id".to_string()]));
+    assert_eq!(file.value_stats.null_counts(), &vec![Some(0)]);
+    builder.new_commit().commit(messages).await.unwrap();
+
+    let plain = table.new_read_builder();
+    assert_eq!(
+        plain
+            .new_scan()
+            .with_scan_all_files()
+            .plan()
+            .await
+            .unwrap()
+            .splits()
+            .len(),
+        1
+    );
+    for column in ["payload", "items"] {
+        let mut filtered = table.new_read_builder();
+        filtered.with_filter(
+            PredicateBuilder::new(table.schema().fields())
+                .is_not_null(column)
+                .unwrap(),
+        );
+        assert_eq!(
+            filtered
+                .new_scan()
+                .with_scan_all_files()
+                .plan()
+                .await
+                .unwrap()
+                .splits()
+                .len(),
+            1,
+            "{column} IS NOT NULL must retain the data file",
+        );
+    }
+}

Reply via email to