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",
+ );
+ }
+}