From dc84bca16ba5c60020364c6ce04d91e834e0287f Mon Sep 17 00:00:00 2001 From: jackylee-ch Date: Thu, 24 Sep 2026 17:31:09 +0800 Subject: [PATCH 1/2] feat(datafusion): support the $buckets system table Java exposes BucketsTable (`$buckets`) for per-bucket file statistics -- the standard way to diagnose bucket skew and decide whether a bucket needs compaction. The Rust DataFusion integration had no equivalent, so that diagnosis was unreachable from DataFusion. Add the provider: aggregate the table's data files by (partition, bucket) into record_count, file_size_in_bytes, file_count and last_update_time, mirroring Java's schema and partition/bucket ordering. The rows come from the same scan the $files table already uses; this only groups them. Reads fail closed under query-auth like the other metadata tables. --- .../datafusion/src/system_tables/buckets.rs | 187 ++++++++++++++++++ .../datafusion/src/system_tables/mod.rs | 4 + .../datafusion/tests/system_tables.rs | 118 +++++++++++ 3 files changed, 309 insertions(+) create mode 100644 crates/integrations/datafusion/src/system_tables/buckets.rs diff --git a/crates/integrations/datafusion/src/system_tables/buckets.rs b/crates/integrations/datafusion/src/system_tables/buckets.rs new file mode 100644 index 000000000..9d4ff87b7 --- /dev/null +++ b/crates/integrations/datafusion/src/system_tables/buckets.rs @@ -0,0 +1,187 @@ +// 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. + +//! Mirrors Java [BucketsTable](https://github.com/apache/paimon/blob/release-1.4/paimon-core/src/main/java/org/apache/paimon/table/system/BucketsTable.java). + +use std::collections::BTreeMap; +use std::sync::{Arc, OnceLock}; + +use async_trait::async_trait; +use datafusion::arrow::array::{ + Int32Array, Int64Array, RecordBatch, StringArray, TimestampMillisecondArray, +}; +use datafusion::arrow::datatypes::{DataType as ArrowDataType, Field, Schema, SchemaRef, TimeUnit}; +use datafusion::catalog::Session; +use datafusion::datasource::memory::MemorySourceConfig; +use datafusion::datasource::{TableProvider, TableType}; +use datafusion::error::Result as DFResult; +use datafusion::logical_expr::Expr; +use datafusion::physical_plan::ExecutionPlan; +use paimon::spec::{BinaryRow, DataField}; +use paimon::table::Table; + +use super::row_string_cast::format_row_as_java_cast_string; +use crate::error::to_datafusion_error; + +pub(super) fn build(table: Table) -> DFResult> { + Ok(Arc::new(BucketsTable { table })) +} + +fn buckets_schema() -> SchemaRef { + static SCHEMA: OnceLock = OnceLock::new(); + SCHEMA + .get_or_init(|| { + Arc::new(Schema::new(vec![ + Field::new("partition", ArrowDataType::Utf8, true), + Field::new("bucket", ArrowDataType::Int32, false), + Field::new("record_count", ArrowDataType::Int64, false), + Field::new("file_size_in_bytes", ArrowDataType::Int64, false), + Field::new("file_count", ArrowDataType::Int64, false), + Field::new( + "last_update_time", + ArrowDataType::Timestamp(TimeUnit::Millisecond, None), + true, + ), + ])) + }) + .clone() +} + +#[derive(Debug)] +pub(super) struct BucketsTable { + table: Table, +} + +#[async_trait] +impl TableProvider for BucketsTable { + fn schema(&self) -> SchemaRef { + buckets_schema() + } + + fn table_type(&self) -> TableType { + TableType::View + } + + async fn scan( + &self, + _state: &dyn Session, + projection: Option<&Vec>, + _filters: &[Expr], + _limit: Option, + ) -> DFResult> { + let table = self.table.clone(); + let rows = + crate::runtime::await_with_runtime(async move { collect_bucket_rows(&table).await }) + .await + .map_err(to_datafusion_error)?; + let batch = bucket_rows_to_record_batch(&rows)?; + Ok(MemorySourceConfig::try_new_exec( + &[vec![batch]], + buckets_schema(), + projection.cloned(), + )?) + } +} + +/// A bucket's files aggregated: mirrors Java `BucketEntry`. +#[derive(Default)] +struct BucketAgg { + record_count: i64, + file_size_in_bytes: i64, + file_count: i64, + last_update_time: Option, +} + +struct BucketRow { + partition: Option, + bucket: i32, + agg: BucketAgg, +} + +async fn collect_bucket_rows(table: &Table) -> paimon::Result> { + let scan = table + .new_read_builder() + .new_scan() + .with_scan_all_files() + .plan() + .await?; + let partition_fields = table.schema().partition_fields(); + // BTreeMap keys sort by partition string then bucket, matching Java BucketsTable's + // `Comparator.comparing(partition).thenComparing(bucket)`. + let mut aggs: BTreeMap<(Option, i32), BucketAgg> = BTreeMap::new(); + for split in scan.splits() { + let partition = format_partition(split.partition(), &partition_fields)?; + let agg = aggs.entry((partition, split.bucket())).or_default(); + for file in split.data_files() { + agg.record_count = agg.record_count.saturating_add(file.row_count); + agg.file_size_in_bytes = agg.file_size_in_bytes.saturating_add(file.file_size); + agg.file_count += 1; + if let Some(t) = file.creation_time.map(|t| t.timestamp_millis()) { + agg.last_update_time = Some(agg.last_update_time.map_or(t, |cur| cur.max(t))); + } + } + } + Ok(aggs + .into_iter() + .map(|((partition, bucket), agg)| BucketRow { + partition, + bucket, + agg, + }) + .collect()) +} + +fn bucket_rows_to_record_batch(rows: &[BucketRow]) -> DFResult { + let n = rows.len(); + let mut partitions = Vec::with_capacity(n); + let mut buckets = 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 last_update_times = Vec::with_capacity(n); + for row in rows { + partitions.push(row.partition.clone()); + buckets.push(row.bucket); + record_counts.push(row.agg.record_count); + file_sizes.push(row.agg.file_size_in_bytes); + file_counts.push(row.agg.file_count); + last_update_times.push(row.agg.last_update_time); + } + Ok(RecordBatch::try_new( + buckets_schema(), + vec![ + Arc::new(StringArray::from(partitions)), + Arc::new(Int32Array::from(buckets)), + Arc::new(Int64Array::from(record_counts)), + Arc::new(Int64Array::from(file_sizes)), + Arc::new(Int64Array::from(file_counts)), + Arc::new(TimestampMillisecondArray::from(last_update_times)), + ], + )?) +} + +/// Format `partition` as Java's cast-to-string, matching `$files`; `{}` when the +/// table is not partitioned. +fn format_partition( + partition: &BinaryRow, + partition_fields: &[DataField], +) -> paimon::Result> { + if partition_fields.is_empty() { + return Ok(Some("{}".to_string())); + } + format_row_as_java_cast_string(partition, partition_fields).map(Some) +} diff --git a/crates/integrations/datafusion/src/system_tables/mod.rs b/crates/integrations/datafusion/src/system_tables/mod.rs index 6e338f78c..f3fd36f9e 100644 --- a/crates/integrations/datafusion/src/system_tables/mod.rs +++ b/crates/integrations/datafusion/src/system_tables/mod.rs @@ -33,6 +33,7 @@ use crate::error::to_datafusion_error; mod aggregation_fields; mod audit_log; mod branches; +mod buckets; mod consumers; mod files; mod manifests; @@ -55,6 +56,7 @@ const TABLES: &[(&str, Builder)] = &[ ("aggregation_fields", aggregation_fields::build), ("audit_log", audit_log::build), ("branches", branches::build), + ("buckets", buckets::build), ("consumers", consumers::build), ("files", files::build), ("manifests", manifests::build), @@ -71,6 +73,7 @@ const SYSTEM_TABLE_NAMES: &[&str] = &[ "aggregation_fields", "audit_log", "branches", + "buckets", "consumers", "files", "manifests", @@ -112,6 +115,7 @@ pub(crate) fn is_registered(name: &str) -> bool { pub(crate) fn is_system_table_provider(provider: &dyn TableProvider) -> bool { provider.is::() || provider.is::() + || provider.is::() || provider.is::() || provider.is::() || provider.is::() diff --git a/crates/integrations/datafusion/tests/system_tables.rs b/crates/integrations/datafusion/tests/system_tables.rs index 84e5aa204..73dc6ba36 100644 --- a/crates/integrations/datafusion/tests/system_tables.rs +++ b/crates/integrations/datafusion/tests/system_tables.rs @@ -110,6 +110,7 @@ async fn test_query_auth_system_tables_fail_closed() { "SELECT * FROM paimon.default.qa$partitions", "SELECT * FROM paimon.default.qa$manifests", "SELECT * FROM paimon.default.qa$table_indexes", + "SELECT * FROM paimon.default.qa$buckets", ] { let err = query_error(&ctx, sql).await; assert!( @@ -129,6 +130,7 @@ async fn test_query_auth_system_tables_fail_closed() { "SELECT * FROM paimon.default.qa_dynamic$partitions", "SELECT * FROM paimon.default.qa_dynamic$manifests", "SELECT * FROM paimon.default.qa_dynamic$table_indexes", + "SELECT * FROM paimon.default.qa_dynamic$buckets", ] { let err = query_error(&ctx, sql).await; assert!( @@ -1611,3 +1613,119 @@ async fn test_aggregation_fields_system_table() { "a field with no fields.* option has no option key" ); } + +#[tokio::test] +async fn test_buckets_system_table() { + let (ctx, _catalog, _tmp) = create_context().await; + + // Ground truth: aggregate the tested $files table by (partition, bucket). + // $buckets must reproduce it exactly (same partition rendering, same per-file values). + let files = run_sql( + &ctx, + &format!( + "SELECT partition, bucket, record_count, file_size_in_bytes \ + FROM paimon.default.{FIXTURE_TABLE}$files" + ), + ) + .await; + let mut expected: std::collections::BTreeMap<(String, i32), (i64, i64, i64)> = + std::collections::BTreeMap::new(); + for batch in &files { + let parts = batch + .column(0) + .as_any() + .downcast_ref::() + .unwrap(); + let buckets = batch + .column(1) + .as_any() + .downcast_ref::() + .unwrap(); + let rc = batch + .column(2) + .as_any() + .downcast_ref::() + .unwrap(); + let sz = batch + .column(3) + .as_any() + .downcast_ref::() + .unwrap(); + for i in 0..batch.num_rows() { + let e = expected + .entry((parts.value(i).to_string(), buckets.value(i))) + .or_default(); + e.0 += rc.value(i); + e.1 += sz.value(i); + e.2 += 1; + } + } + assert!(!expected.is_empty(), "fixture should contain data files"); + + let sql = format!("SELECT * FROM paimon.default.{FIXTURE_TABLE}$buckets"); + let batches = run_sql(&ctx, &sql).await; + assert!(!batches.is_empty(), "$buckets should return ≥1 batch"); + + let arrow_schema = batches[0].schema(); + let expected_columns = [ + ("partition", DataType::Utf8), + ("bucket", DataType::Int32), + ("record_count", DataType::Int64), + ("file_size_in_bytes", DataType::Int64), + ("file_count", DataType::Int64), + ( + "last_update_time", + DataType::Timestamp(TimeUnit::Millisecond, None), + ), + ]; + for (i, (name, dtype)) in expected_columns.iter().enumerate() { + let field = arrow_schema.field(i); + assert_eq!(field.name(), name, "column {i} name"); + assert_eq!(field.data_type(), dtype, "column {i} type"); + } + + let mut actual: std::collections::BTreeMap<(String, i32), (i64, i64, i64)> = + std::collections::BTreeMap::new(); + for batch in &batches { + let parts = batch + .column(0) + .as_any() + .downcast_ref::() + .unwrap(); + let buckets = batch + .column(1) + .as_any() + .downcast_ref::() + .unwrap(); + let rc = batch + .column(2) + .as_any() + .downcast_ref::() + .unwrap(); + let sz = batch + .column(3) + .as_any() + .downcast_ref::() + .unwrap(); + let fc = batch + .column(4) + .as_any() + .downcast_ref::() + .unwrap(); + for i in 0..batch.num_rows() { + let prev = actual.insert( + (parts.value(i).to_string(), buckets.value(i)), + (rc.value(i), sz.value(i), fc.value(i)), + ); + assert!( + prev.is_none(), + "$buckets must emit one row per (partition, bucket)" + ); + } + } + + assert_eq!( + actual, expected, + "$buckets must aggregate $files by (partition, bucket)" + ); +} From 7d5c3de1f63281f5027816667eb933dd930150b4 Mon Sep 17 00:00:00 2001 From: jackylee-ch Date: Sat, 26 Sep 2026 20:53:35 +0800 Subject: [PATCH 2/2] fix(datafusion): group $buckets rows by the raw partition MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Review follow-up (#934): `collect_bucket_rows` keyed its aggregation on the rendered partition string, but the Java cast-to-string formatter is not injective — `(p1='a, b', p2='c')` and `(p1='a', p2='b, c')` both render `{a, b, c}` — so distinct partitions merged into one `$buckets` row with combined counts. Key on the serialized `BinaryRow` bytes instead, keeping the string for output and ordering; extract `aggregate_bucket_rows` so the look-alike case has a focused regression test. --- .../datafusion/src/system_tables/buckets.rs | 87 +++++++++++++++++-- 1 file changed, 79 insertions(+), 8 deletions(-) diff --git a/crates/integrations/datafusion/src/system_tables/buckets.rs b/crates/integrations/datafusion/src/system_tables/buckets.rs index 9d4ff87b7..f34b6a803 100644 --- a/crates/integrations/datafusion/src/system_tables/buckets.rs +++ b/crates/integrations/datafusion/src/system_tables/buckets.rs @@ -32,7 +32,7 @@ use datafusion::error::Result as DFResult; use datafusion::logical_expr::Expr; use datafusion::physical_plan::ExecutionPlan; use paimon::spec::{BinaryRow, DataField}; -use paimon::table::Table; +use paimon::table::{DataSplit, Table}; use super::row_string_cast::format_row_as_java_cast_string; use crate::error::to_datafusion_error; @@ -120,12 +120,33 @@ async fn collect_bucket_rows(table: &Table) -> paimon::Result> { .plan() .await?; let partition_fields = table.schema().partition_fields(); - // BTreeMap keys sort by partition string then bucket, matching Java BucketsTable's - // `Comparator.comparing(partition).thenComparing(bucket)`. - let mut aggs: BTreeMap<(Option, i32), BucketAgg> = BTreeMap::new(); - for split in scan.splits() { - let partition = format_partition(split.partition(), &partition_fields)?; - let agg = aggs.entry((partition, split.bucket())).or_default(); + aggregate_bucket_rows(scan.splits(), &partition_fields) +} + +/// Aggregate splits into one row per (partition, bucket). +/// +/// Group by the raw partition row, not its rendered string. The Java +/// cast-to-string formatter is not injective — e.g. `(p1='a, b', p2='c')` and +/// `(p1='a', p2='b, c')` both render `{a, b, c}` — so keying on the string would +/// merge unrelated partitions and report combined counts. The serialized +/// `BinaryRow` bytes are injective and form the grouping key; the string is +/// carried only for output. Keys still sort by partition string then bucket to +/// match Java BucketsTable's `Comparator.comparing(partition).thenComparing(bucket)`, +/// with the raw bytes as a final tie-break so distinct look-alike partitions keep +/// a deterministic order. +fn aggregate_bucket_rows( + splits: &[DataSplit], + partition_fields: &[DataField], +) -> paimon::Result> { + let mut aggs: BTreeMap<(Option, i32, Vec), BucketAgg> = BTreeMap::new(); + for split in splits { + let partition = format_partition(split.partition(), partition_fields)?; + let key = ( + partition, + split.bucket(), + split.partition().to_serialized_bytes(), + ); + let agg = aggs.entry(key).or_default(); for file in split.data_files() { agg.record_count = agg.record_count.saturating_add(file.row_count); agg.file_size_in_bytes = agg.file_size_in_bytes.saturating_add(file.file_size); @@ -137,7 +158,7 @@ async fn collect_bucket_rows(table: &Table) -> paimon::Result> { } Ok(aggs .into_iter() - .map(|((partition, bucket), agg)| BucketRow { + .map(|((partition, bucket, _), agg)| BucketRow { partition, bucket, agg, @@ -185,3 +206,53 @@ fn format_partition( } format_row_as_java_cast_string(partition, partition_fields).map(Some) } + +#[cfg(test)] +mod tests { + use super::*; + use paimon::spec::{DataType, Datum, VarCharType}; + + fn string_type() -> DataType { + DataType::VarChar(VarCharType::new(64).unwrap()) + } + + fn partition_row(a: &str, b: &str) -> BinaryRow { + let t = string_type(); + let da = Datum::String(a.to_string()); + let db = Datum::String(b.to_string()); + BinaryRow::from_datums(&[(Some(&da), &t), (Some(&db), &t)]) + } + + fn split(partition: BinaryRow, bucket: i32) -> DataSplit { + DataSplit::builder() + .with_snapshot(1) + .with_partition(partition) + .with_bucket(bucket) + .with_bucket_path("/warehouse/bucket".to_string()) + .with_data_files(Vec::new()) + .build() + .unwrap() + } + + // Two distinct partitions whose Java cast-string renders identically must + // still aggregate into separate rows. Grouping on the rendered string (the + // pre-fix behavior) would merge them into one row with combined counts. + #[test] + fn look_alike_partitions_do_not_merge() { + let t = string_type(); + let fields = vec![ + DataField::new(0, "p1".to_string(), t.clone()), + DataField::new(1, "p2".to_string(), t), + ]; + let a = partition_row("a, b", "c"); + let b = partition_row("a", "b, c"); + // Precondition: the rendered strings collide, so string keying would merge. + assert_eq!( + format_partition(&a, &fields).unwrap(), + format_partition(&b, &fields).unwrap(), + ); + + let rows = aggregate_bucket_rows(&[split(a, 0), split(b, 0)], &fields).unwrap(); + assert_eq!(rows.len(), 2, "distinct partitions must not be merged"); + } +}