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

Reply via email to