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 2148923f feat(datafusion): support the $aggregation_fields system
table (#937)
2148923f is described below
commit 2148923fffcb50b5a0d7a325940d81cd829c99ea
Author: jackylee <[email protected]>
AuthorDate: Fri Sep 25 20:50:46 2026 +0800
feat(datafusion): support the $aggregation_fields system table (#937)
---
.../src/system_tables/aggregation_fields.rs | 133 +++++++++++++++++++++
.../datafusion/src/system_tables/mod.rs | 6 +-
.../integrations/datafusion/tests/system_tables.rs | 96 +++++++++++++++
3 files changed, 234 insertions(+), 1 deletion(-)
diff --git
a/crates/integrations/datafusion/src/system_tables/aggregation_fields.rs
b/crates/integrations/datafusion/src/system_tables/aggregation_fields.rs
new file mode 100644
index 00000000..9aac6084
--- /dev/null
+++ b/crates/integrations/datafusion/src/system_tables/aggregation_fields.rs
@@ -0,0 +1,133 @@
+// 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
[AggregationFieldsTable](https://github.com/apache/paimon/blob/release-1.4/paimon-core/src/main/java/org/apache/paimon/table/system/AggregationFieldsTable.java).
+
+use std::collections::HashMap;
+use std::sync::{Arc, OnceLock};
+
+use async_trait::async_trait;
+use datafusion::arrow::array::{RecordBatch, StringArray};
+use datafusion::arrow::datatypes::{DataType as ArrowDataType, Field, Schema,
SchemaRef};
+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::table::Table;
+
+pub(super) fn build(table: Table) -> DFResult<Arc<dyn TableProvider>> {
+ Ok(Arc::new(AggregationFieldsTable { table }))
+}
+
+fn aggregation_fields_schema() -> SchemaRef {
+ static SCHEMA: OnceLock<SchemaRef> = OnceLock::new();
+ SCHEMA
+ .get_or_init(|| {
+ Arc::new(Schema::new(vec![
+ Field::new("field_name", ArrowDataType::Utf8, false),
+ Field::new("field_type", ArrowDataType::Utf8, false),
+ Field::new("function", ArrowDataType::Utf8, true),
+ Field::new("function_options", ArrowDataType::Utf8, true),
+ Field::new("comment", ArrowDataType::Utf8, true),
+ ]))
+ })
+ .clone()
+}
+
+#[derive(Debug)]
+pub(super) struct AggregationFieldsTable {
+ table: Table,
+}
+
+#[async_trait]
+impl TableProvider for AggregationFieldsTable {
+ fn schema(&self) -> SchemaRef {
+ aggregation_fields_schema()
+ }
+
+ fn table_type(&self) -> TableType {
+ TableType::View
+ }
+
+ async fn scan(
+ &self,
+ _state: &dyn Session,
+ projection: Option<&Vec<usize>>,
+ _filters: &[Expr],
+ _limit: Option<usize>,
+ ) -> DFResult<Arc<dyn ExecutionPlan>> {
+ let table_schema = self.table.schema();
+ let options = table_schema.options();
+
+ let n = table_schema.fields().len();
+ let mut field_names = Vec::with_capacity(n);
+ let mut field_types = Vec::with_capacity(n);
+ let mut functions = Vec::with_capacity(n);
+ let mut function_options = Vec::with_capacity(n);
+ let mut comments: Vec<Option<String>> = Vec::with_capacity(n);
+ for field in table_schema.fields() {
+ field_names.push(field.name().to_string());
+ field_types.push(field.data_type().to_string());
+ let (keys, values) = field_scoped_options(options, field.name());
+ functions.push(format!("[{}]", values.join(", ")));
+ function_options.push(format!("[{}]", keys.join(", ")));
+ comments.push(field.description().map(str::to_string));
+ }
+
+ let schema = aggregation_fields_schema();
+ let batch = RecordBatch::try_new(
+ schema.clone(),
+ vec![
+ Arc::new(StringArray::from(field_names)),
+ Arc::new(StringArray::from(field_types)),
+ Arc::new(StringArray::from(functions)),
+ Arc::new(StringArray::from(function_options)),
+ Arc::new(StringArray::from(comments)),
+ ],
+ )?;
+
+ Ok(MemorySourceConfig::try_new_exec(
+ &[vec![batch]],
+ schema,
+ projection.cloned(),
+ )?)
+ }
+}
+
+/// Keys and values of the options scoped to `fields.<field_name>.*`, mirroring
+/// Java `AggregationFieldsTable.extractFieldMultimap`. Java renders each as a
+/// collection `toString` (`[a, b]`); we sort by key so the output is stable
+/// (Java iterates its options map in unspecified order — single-option fields,
+/// the common case, render identically either way).
+fn field_scoped_options(
+ options: &HashMap<String, String>,
+ field_name: &str,
+) -> (Vec<String>, Vec<String>) {
+ let mut pairs: Vec<(&String, &String)> = options
+ .iter()
+ .filter(|(key, _)| {
+ let parts: Vec<&str> = key.split('.').collect();
+ parts.len() > 2 && parts[0] == "fields" && parts[1] == field_name
+ })
+ .collect();
+ pairs.sort_by(|a, b| a.0.cmp(b.0));
+ let keys = pairs.iter().map(|(k, _)| (*k).clone()).collect();
+ let values = pairs.iter().map(|(_, v)| (*v).clone()).collect();
+ (keys, values)
+}
diff --git a/crates/integrations/datafusion/src/system_tables/mod.rs
b/crates/integrations/datafusion/src/system_tables/mod.rs
index 22fbc508..6e338f78 100644
--- a/crates/integrations/datafusion/src/system_tables/mod.rs
+++ b/crates/integrations/datafusion/src/system_tables/mod.rs
@@ -30,6 +30,7 @@ use paimon::table::Table;
use crate::error::to_datafusion_error;
+mod aggregation_fields;
mod audit_log;
mod branches;
mod consumers;
@@ -51,6 +52,7 @@ type Builder = fn(Table) -> DFResult<Arc<dyn TableProvider>>;
// in `load` because it needs the catalog handle (for metastore-tracked audit
// metadata via `Catalog::list_partitions`).
const TABLES: &[(&str, Builder)] = &[
+ ("aggregation_fields", aggregation_fields::build),
("audit_log", audit_log::build),
("branches", branches::build),
("consumers", consumers::build),
@@ -66,6 +68,7 @@ const TABLES: &[(&str, Builder)] = &[
];
const SYSTEM_TABLE_NAMES: &[&str] = &[
+ "aggregation_fields",
"audit_log",
"branches",
"consumers",
@@ -107,7 +110,8 @@ pub(crate) fn is_registered(name: &str) -> bool {
/// with a `$`. Keep in step with `TABLES`: one missing here goes back to
/// having its time-travel clause silently dropped.
pub(crate) fn is_system_table_provider(provider: &dyn TableProvider) -> bool {
- provider.is::<branches::BranchesTable>()
+ provider.is::<aggregation_fields::AggregationFieldsTable>()
+ || provider.is::<branches::BranchesTable>()
|| provider.is::<consumers::ConsumersTable>()
|| provider.is::<files::FilesTable>()
|| provider.is::<manifests::ManifestsTable>()
diff --git a/crates/integrations/datafusion/tests/system_tables.rs
b/crates/integrations/datafusion/tests/system_tables.rs
index dae555f7..84e5aa20 100644
--- a/crates/integrations/datafusion/tests/system_tables.rs
+++ b/crates/integrations/datafusion/tests/system_tables.rs
@@ -1515,3 +1515,99 @@ async fn test_partitions_system_table() {
}
}
}
+
+#[tokio::test]
+async fn test_aggregation_fields_system_table() {
+ let (ctx, _catalog, _tmp) = create_context().await;
+ run_sql(
+ &ctx,
+ "CREATE TABLE paimon.default.agg_fields (
+ id INT NOT NULL,
+ total BIGINT,
+ PRIMARY KEY (id)
+ ) WITH (
+ 'bucket' = '1',
+ 'merge-engine' = 'aggregation',
+ 'fields.total.aggregate-function' = 'sum'
+ )",
+ )
+ .await;
+
+ let batches = run_sql(
+ &ctx,
+ "SELECT * FROM paimon.default.agg_fields$aggregation_fields",
+ )
+ .await;
+ assert!(
+ !batches.is_empty(),
+ "$aggregation_fields should return ≥1 batch"
+ );
+
+ let arrow_schema = batches[0].schema();
+ let expected_columns = [
+ ("field_name", DataType::Utf8),
+ ("field_type", DataType::Utf8),
+ ("function", DataType::Utf8),
+ ("function_options", DataType::Utf8),
+ ("comment", DataType::Utf8),
+ ];
+ 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");
+ }
+ // MARKER_AGG_FIELDS
+ let mut rows: std::collections::BTreeMap<String, (String, String, String)>
=
+ std::collections::BTreeMap::new();
+ for batch in &batches {
+ let names = batch
+ .column(0)
+ .as_any()
+ .downcast_ref::<StringArray>()
+ .unwrap();
+ let types = batch
+ .column(1)
+ .as_any()
+ .downcast_ref::<StringArray>()
+ .unwrap();
+ let functions = batch
+ .column(2)
+ .as_any()
+ .downcast_ref::<StringArray>()
+ .unwrap();
+ let function_options = batch
+ .column(3)
+ .as_any()
+ .downcast_ref::<StringArray>()
+ .unwrap();
+ for i in 0..batch.num_rows() {
+ rows.insert(
+ names.value(i).to_string(),
+ (
+ types.value(i).to_string(),
+ functions.value(i).to_string(),
+ function_options.value(i).to_string(),
+ ),
+ );
+ }
+ }
+
+ let total = rows.get("total").expect("total field must appear");
+ assert!(!total.0.is_empty(), "field_type must be rendered");
+ assert_eq!(total.1, "[sum]", "configured aggregate function");
+ assert_eq!(
+ total.2, "[fields.total.aggregate-function]",
+ "the option key that set the function"
+ );
+
+ let id = rows.get("id").expect("id field must appear");
+ assert!(!id.0.is_empty(), "field_type must be rendered");
+ assert_eq!(
+ id.1, "[]",
+ "a field with no fields.* option has no function"
+ );
+ assert_eq!(
+ id.2, "[]",
+ "a field with no fields.* option has no option key"
+ );
+}