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 79567f62 fix(datafusion): show an unreported partition statistic as
NULL, not as -1 (#697)
79567f62 is described below
commit 79567f625aead44425756d2d312f89224f641880
Author: Dapeng Sun(孙大鹏) <[email protected]>
AuthorDate: Sun Aug 9 21:04:36 2026 +0800
fix(datafusion): show an unreported partition statistic as NULL, not as -1
(#697)
---
.../datafusion/src/system_tables/partitions.rs | 31 +++++---
.../datafusion/tests/sql_context_tests.rs | 92 +++++++++++++++++++++-
crates/paimon/src/spec/partition.rs | 34 ++++++++
3 files changed, 144 insertions(+), 13 deletions(-)
diff --git a/crates/integrations/datafusion/src/system_tables/partitions.rs
b/crates/integrations/datafusion/src/system_tables/partitions.rs
index 97cf51b8..a0b86c2b 100644
--- a/crates/integrations/datafusion/src/system_tables/partitions.rs
+++ b/crates/integrations/datafusion/src/system_tables/partitions.rs
@@ -62,9 +62,12 @@ fn partitions_schema() -> SchemaRef {
.get_or_init(|| {
Arc::new(Schema::new(vec![
Field::new("partition", DataType::Utf8, true),
- Field::new("record_count", DataType::Int64, false),
- Field::new("file_size_in_bytes", DataType::Int64, false),
- Field::new("file_count", DataType::Int64, false),
+ // Nullable because a statistic the catalog never had reported
to it has no value to
+ // show. Declaring these non-null forced Partition::UNKNOWN to
be rendered as -1,
+ // which reads as a measurement.
+ Field::new("record_count", DataType::Int64, true),
+ Field::new("file_size_in_bytes", DataType::Int64, true),
+ Field::new("file_count", DataType::Int64, true),
Field::new(
"last_update_time",
DataType::Timestamp(TimeUnit::Millisecond, None),
@@ -78,7 +81,7 @@ fn partitions_schema() -> SchemaRef {
Field::new("created_by", DataType::Utf8, true),
Field::new("updated_by", DataType::Utf8, true),
Field::new("options", DataType::Utf8, true),
- Field::new("total_buckets", DataType::Int32, false),
+ Field::new("total_buckets", DataType::Int32, true),
Field::new("done", DataType::Boolean, false),
]))
})
@@ -159,22 +162,25 @@ impl TableProvider for PartitionsTable {
let n = rows.len();
let mut partition_strings: Vec<Option<String>> = Vec::with_capacity(n);
- let mut record_counts = Vec::with_capacity(n);
- let mut file_sizes = Vec::with_capacity(n);
- let mut file_counts = Vec::with_capacity(n);
+ let mut record_counts: Vec<Option<i64>> = Vec::with_capacity(n);
+ let mut file_sizes: Vec<Option<i64>> = Vec::with_capacity(n);
+ let mut file_counts: Vec<Option<i64>> = Vec::with_capacity(n);
let mut last_update_times: Vec<Option<i64>> = Vec::with_capacity(n);
let mut created_ats: Vec<Option<i64>> = Vec::with_capacity(n);
let mut created_bys: Vec<Option<String>> = Vec::with_capacity(n);
let mut updated_bys: Vec<Option<String>> = Vec::with_capacity(n);
let mut options_jsons: Vec<Option<String>> = Vec::with_capacity(n);
- let mut total_buckets = Vec::with_capacity(n);
+ let mut total_buckets: Vec<Option<i32>> = Vec::with_capacity(n);
let mut dones = Vec::with_capacity(n);
for (s, p) in rows {
partition_strings.push(Some(s));
- record_counts.push(p.record_count);
- file_sizes.push(p.file_size_in_bytes);
- file_counts.push(p.file_count);
+ // A field nobody reported on shows as NULL. Passing the
placeholder straight through
+ // would claim the partition holds -1 rows, and 0 would claim it
is empty.
+
record_counts.push(Partition::is_known(p.record_count).then_some(p.record_count));
+ file_sizes
+
.push(Partition::is_known(p.file_size_in_bytes).then_some(p.file_size_in_bytes));
+
file_counts.push(Partition::is_known(p.file_count).then_some(p.file_count));
// 0 marks "no creation_time on any file"; real wall-clock is never
// <= 0 in practice, so this never nullifies a genuine timestamp.
last_update_times.push(if p.last_file_creation_time > 0 {
@@ -197,7 +203,8 @@ impl TableProvider for PartitionsTable {
})
.transpose()?,
);
- total_buckets.push(p.total_buckets);
+ total_buckets
+ .push(Some(p.total_buckets).filter(|b| *b !=
Partition::UNKNOWN_TOTAL_BUCKETS));
dones.push(p.done);
}
diff --git a/crates/integrations/datafusion/tests/sql_context_tests.rs
b/crates/integrations/datafusion/tests/sql_context_tests.rs
index b1f37d85..37b08c20 100644
--- a/crates/integrations/datafusion/tests/sql_context_tests.rs
+++ b/crates/integrations/datafusion/tests/sql_context_tests.rs
@@ -152,6 +152,7 @@ impl Catalog for MetadataListingCatalog {
struct PartitionCatalog {
inner: Arc<FileSystemCatalog>,
fail_list_partitions: AtomicBool,
+ unknown_statistics: AtomicBool,
partition_identifiers: Mutex<Vec<Identifier>>,
}
@@ -160,6 +161,7 @@ impl PartitionCatalog {
Self {
inner,
fail_list_partitions: AtomicBool::new(false),
+ unknown_statistics: AtomicBool::new(false),
partition_identifiers: Mutex::new(Vec::new()),
}
}
@@ -168,6 +170,12 @@ impl PartitionCatalog {
self.fail_list_partitions.store(fail, Ordering::SeqCst);
}
+ /// Report every statistic as never measured, the way a catalog does for a
partition that was
+ /// registered but never had statistics reported to it.
+ fn set_unknown_statistics(&self, unknown: bool) {
+ self.unknown_statistics.store(unknown, Ordering::SeqCst);
+ }
+
fn take_partition_identifiers(&self) -> Vec<Identifier> {
std::mem::take(&mut *self.partition_identifiers.lock().unwrap())
}
@@ -266,7 +274,17 @@ impl Catalog for PartitionCatalog {
.push(identifier.clone());
let Some(branch) = identifier.branch_name()? else {
- return self.inner.list_partitions(identifier).await;
+ let mut partitions = self.inner.list_partitions(identifier).await?;
+ if self.unknown_statistics.load(Ordering::SeqCst) {
+ for partition in &mut partitions {
+ partition.record_count = paimon::spec::Partition::UNKNOWN;
+ partition.file_size_in_bytes =
paimon::spec::Partition::UNKNOWN;
+ partition.file_count = paimon::spec::Partition::UNKNOWN;
+ partition.last_file_creation_time =
paimon::spec::Partition::UNKNOWN;
+ partition.total_buckets =
paimon::spec::Partition::UNKNOWN_TOTAL_BUCKETS;
+ }
+ }
+ return Ok(partitions);
};
if self.fail_list_partitions.load(Ordering::SeqCst) {
return Err(paimon::Error::Unsupported {
@@ -580,6 +598,78 @@ async fn test_select_branch_table_reads_branch_snapshot() {
.await;
}
+/// A statistic the catalog never had reported to it has to read as NULL.
+///
+/// The alternative is what this used to do: declare the columns non-nullable
and let
+/// `Partition::UNKNOWN` through, so `$partitions` claimed the partition holds
-1 rows and -1 files.
+/// Reporting it as `0` instead would be worse still — that is a real
measurement meaning empty.
+#[tokio::test]
+async fn test_partitions_system_table_shows_unreported_statistics_as_null() {
+ let (_tmp, file_catalog) = create_test_env();
+ let catalog = Arc::new(PartitionCatalog::new(file_catalog.clone()));
+ let mut sql_context = SQLContext::new();
+ sql_context
+ .register_catalog("paimon", catalog.clone())
+ .await
+ .unwrap();
+
+ sql_context
+ .sql(
+ "CREATE TABLE paimon.default.unknown_stats_orders \
+ (id INT, name STRING) PARTITIONED BY (id)",
+ )
+ .await
+ .unwrap()
+ .collect()
+ .await
+ .unwrap();
+ sql_context
+ .sql("INSERT INTO paimon.default.unknown_stats_orders VALUES (1, 'a')")
+ .await
+ .unwrap()
+ .collect()
+ .await
+ .unwrap();
+
+ let sql = "SELECT record_count, file_size_in_bytes, file_count,
total_buckets \
+ FROM paimon.default.unknown_stats_orders$partitions";
+
+ // Measured statistics still arrive as values, so a NULL below means
unknown and not that the
+ // column stopped being populated at all.
+ let measured =
sql_context.sql(sql).await.unwrap().collect().await.unwrap();
+ let measured = &measured[0];
+ assert_eq!(measured.num_rows(), 1);
+ for column in 0..4 {
+ assert!(
+ !measured.column(column).is_null(0),
+ "column {column} should carry a measurement before the switch"
+ );
+ }
+
+ catalog.set_unknown_statistics(true);
+
+ let batches = sql_context.sql(sql).await.unwrap().collect().await.unwrap();
+ let batch = &batches[0];
+ assert_eq!(batch.num_rows(), 1);
+ for column in 0..4 {
+ assert!(
+ batch.column(column).is_null(0),
+ "column {column} was never measured and must read as NULL"
+ );
+ }
+
+ // The partition itself is still registered; only its statistics are
unknown.
+ assert_eq!(
+ collect_string_column(
+ &sql_context,
+ "SELECT \"partition\" FROM
paimon.default.unknown_stats_orders$partitions",
+ "partition",
+ )
+ .await,
+ vec!["id=1".to_string()]
+ );
+}
+
#[tokio::test]
async fn test_branch_partitions_system_table_reads_branch_snapshot() {
let (_tmp, file_catalog) = create_test_env();
diff --git a/crates/paimon/src/spec/partition.rs
b/crates/paimon/src/spec/partition.rs
index 478fa0b7..a78fe5b2 100644
--- a/crates/paimon/src/spec/partition.rs
+++ b/crates/paimon/src/spec/partition.rs
@@ -46,6 +46,27 @@ pub struct Partition {
pub options: Option<HashMap<String, String>>,
}
+impl Partition {
+ /// A statistic the catalog never had reported to it.
+ ///
+ /// On this plane — what `list_partitions` observed — any negative value
means the field was
+ /// never measured, and `-1` is the canonical spelling. `0` is an exact
zero and must stay
+ /// distinguishable from it: a partition nobody reported on is not an
empty partition. The
+ /// counterpart for file-level counts is
[`crate::spec::DataFileMeta::ROW_COUNT_UNKNOWN`].
+ ///
+ /// This says nothing about the delta plane used by snapshot commits,
where a negative value is
+ /// a decrement to apply rather than a missing measurement.
+ pub const UNKNOWN: i64 = -1;
+
+ /// A partition with no buckets, or none the catalog knows of. Format
tables have no buckets.
+ pub const UNKNOWN_TOTAL_BUCKETS: i32 = -1;
+
+ /// Whether a statistic read off a partition is a measurement rather than
[`Self::UNKNOWN`].
+ pub fn is_known(value: i64) -> bool {
+ value >= 0
+ }
+}
+
#[cfg(test)]
mod tests {
use super::*;
@@ -108,4 +129,17 @@ mod tests {
let decoded: Partition = serde_json::from_str(json).unwrap();
assert_eq!(decoded.spec.get("dt"), Some(&"2024-01-01".to_string()));
}
+
+ #[test]
+ fn test_unknown_is_negative_and_zero_is_a_measurement() {
+ assert_eq!(Partition::UNKNOWN, -1);
+ assert_eq!(Partition::UNKNOWN_TOTAL_BUCKETS, -1);
+
+ assert!(!Partition::is_known(Partition::UNKNOWN));
+ assert!(!Partition::is_known(-42));
+
+ // A partition nobody reported on is not an empty partition.
+ assert!(Partition::is_known(0));
+ assert!(Partition::is_known(1));
+ }
}