This is an automated email from the ASF dual-hosted git repository.

JingsongLi pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/paimon.git


The following commit(s) were added to refs/heads/master by this push:
     new 7f0814d937 [python] Add the $consumers system table (#10174)
7f0814d937 is described below

commit 7f0814d9376a0b9f445a89e984c715a8bceb1f3f
Author: jackylee <[email protected]>
AuthorDate: Fri Sep 25 21:09:22 2026 +0800

    [python] Add the $consumers system table (#10174)
---
 .../pypaimon/table/system/consumers_table.py       | 63 ++++++++++++++
 .../pypaimon/table/system/system_table_loader.py   |  5 +-
 .../pypaimon/tests/system/consumers_table_test.py  | 99 ++++++++++++++++++++++
 .../tests/system/system_table_loader_test.py       |  2 +-
 4 files changed, 167 insertions(+), 2 deletions(-)

diff --git a/paimon-python/pypaimon/table/system/consumers_table.py 
b/paimon-python/pypaimon/table/system/consumers_table.py
new file mode 100644
index 0000000000..c5d517dfed
--- /dev/null
+++ b/paimon-python/pypaimon/table/system/consumers_table.py
@@ -0,0 +1,63 @@
+# 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.
+
+"""The ``$consumers`` system table — streaming consumer progress."""
+
+from typing import List
+
+import pyarrow
+
+from pypaimon.schema.data_types import AtomicType, DataField, RowType
+from pypaimon.table.system.system_table import SystemTable
+
+
+TABLE_TYPE = RowType(False, [
+    DataField(0, "consumer_id", AtomicType("STRING", nullable=False)),
+    DataField(1, "next_snapshot_id", AtomicType("BIGINT", nullable=False)),
+])
+
+
+class ConsumersTable(SystemTable):
+    """The ``$consumers`` system table: one ``(consumer_id,
+    next_snapshot_id)`` row per streaming consumer, so the consumption
+    progress persisted under ``{table}/consumer/`` is queryable the same
+    way Java's ``ConsumersTable`` exposes it (needed under a REST catalog,
+    where the internal consumer API is otherwise the only view).
+    """
+
+    def system_table_name(self) -> str:
+        return "consumers"
+
+    def row_type(self) -> RowType:
+        return TABLE_TYPE
+
+    def primary_keys(self) -> List[str]:
+        return ["consumer_id"]
+
+    def _build_arrow_table(self) -> pyarrow.Table:
+        # consumers() maps consumer_id -> next_snapshot; sort by id so the
+        # output is deterministic. Build with explicit column types so an
+        # empty table still carries the declared (string, int64) schema
+        # rather than pyarrow's null-type inference.
+        consumers = self.base_table.consumer_manager().consumers()
+        consumer_ids = sorted(consumers.keys())
+        next_snapshot_ids = [consumers[cid] for cid in consumer_ids]
+        return pyarrow.table({
+            "consumer_id": pyarrow.array(consumer_ids, type=pyarrow.string()),
+            "next_snapshot_id": pyarrow.array(
+                next_snapshot_ids, type=pyarrow.int64()),
+        })
diff --git a/paimon-python/pypaimon/table/system/system_table_loader.py 
b/paimon-python/pypaimon/table/system/system_table_loader.py
index c2f50649a5..b924fab248 100644
--- a/paimon-python/pypaimon/table/system/system_table_loader.py
+++ b/paimon-python/pypaimon/table/system/system_table_loader.py
@@ -23,7 +23,7 @@ new module.
 
 The following short names are intentionally not registered here yet:
 
-  audit_log, binlog, read_optimized, consumers, statistics,
+  audit_log, binlog, read_optimized, statistics,
   aggregation_fields, row_tracking, all_tables, all_partitions, 
all_table_options,
   catalog_options
 """
@@ -48,6 +48,7 @@ SYSTEM_TABLES: Tuple[str, ...] = (
     "branches",
     "file_key_ranges",
     "table_indexes",
+    "consumers",
 )
 
 
@@ -75,6 +76,8 @@ SYSTEM_TABLE_LOADERS: Dict[str, Callable[..., "SystemTable"]] 
= {
         "pypaimon.table.system.file_key_ranges_table", "FileKeyRangesTable"),
     "table_indexes": _lazy(
         "pypaimon.table.system.table_indexes_table", "TableIndexesTable"),
+    "consumers": _lazy(
+        "pypaimon.table.system.consumers_table", "ConsumersTable"),
 }
 
 
diff --git a/paimon-python/pypaimon/tests/system/consumers_table_test.py 
b/paimon-python/pypaimon/tests/system/consumers_table_test.py
new file mode 100644
index 0000000000..c73ef96af9
--- /dev/null
+++ b/paimon-python/pypaimon/tests/system/consumers_table_test.py
@@ -0,0 +1,99 @@
+# 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.
+
+"""End-to-end tests for the ``$consumers`` system table."""
+
+import os
+import shutil
+import tempfile
+import unittest
+
+from pypaimon import CatalogFactory, Schema
+from pypaimon.consumer.consumer import Consumer
+from pypaimon.schema.data_types import DataField
+from pypaimon.table.system.consumers_table import ConsumersTable
+
+
+def _read(table):
+    rb = table.new_read_builder()
+    return rb.new_read().to_arrow(rb.new_scan().plan().splits())
+
+
+class ConsumersTableTest(unittest.TestCase):
+
+    def setUp(self):
+        self.tmp = tempfile.mkdtemp(prefix="consumers_sys_")
+        warehouse = os.path.join(self.tmp, "warehouse")
+        self.catalog = CatalogFactory.create({"warehouse": warehouse})
+        self.catalog.create_database("db", False)
+        fields = [DataField.from_dict({"id": 0, "name": "v", "type": "INT"})]
+        self.catalog.create_table("db.t", Schema(fields=fields), False)
+
+    def tearDown(self):
+        shutil.rmtree(self.tmp, ignore_errors=True)
+
+    def _record_consumers(self, mapping):
+        manager = self.catalog.get_table("db.t").consumer_manager()
+        for consumer_id, next_snapshot in mapping.items():
+            manager.reset_consumer(consumer_id, Consumer(next_snapshot))
+
+    def test_consumers_table_is_loaded_via_catalog(self):
+        table = self.catalog.get_table("db.t$consumers")
+        self.assertIsInstance(table, ConsumersTable)
+
+    def test_schema_column_layout(self):
+        table = self.catalog.get_table("db.t$consumers")
+        row_type = table.row_type()
+        self.assertEqual(["consumer_id", "next_snapshot_id"],
+                         [f.name for f in row_type.fields])
+        for field in row_type.fields:
+            self.assertFalse(field.type.nullable,
+                             "{} should be NOT NULL".format(field.name))
+        self.assertEqual(["consumer_id"], table.primary_keys())
+
+    def test_read_returns_every_consumer(self):
+        self._record_consumers({"c1": 5, "c2": 8})
+        arrow_table = _read(self.catalog.get_table("db.t$consumers"))
+        self.assertEqual(["consumer_id", "next_snapshot_id"],
+                         arrow_table.schema.names)
+        rendered = dict(zip(
+            arrow_table.column("consumer_id").to_pylist(),
+            arrow_table.column("next_snapshot_id").to_pylist(),
+        ))
+        self.assertEqual({"c1": 5, "c2": 8}, rendered)
+
+    def test_empty_when_no_consumers(self):
+        arrow_table = _read(self.catalog.get_table("db.t$consumers"))
+        self.assertEqual(["consumer_id", "next_snapshot_id"],
+                         arrow_table.schema.names)
+        self.assertEqual(0, arrow_table.num_rows)
+        # The declared (string, int64) schema survives an empty table.
+        self.assertEqual("string", 
str(arrow_table.schema.field("consumer_id").type))
+        self.assertEqual("int64",
+                         
str(arrow_table.schema.field("next_snapshot_id").type))
+
+    def test_read_with_projection_keeps_requested_column_only(self):
+        self._record_consumers({"c1": 5})
+        table = self.catalog.get_table("db.t$consumers")
+        rb = table.new_read_builder().with_projection(["consumer_id"])
+        arrow_table = rb.new_read().to_arrow(rb.new_scan().plan().splits())
+        self.assertEqual(["consumer_id"], arrow_table.schema.names)
+        self.assertEqual(["c1"], arrow_table.column("consumer_id").to_pylist())
+
+
+if __name__ == "__main__":
+    unittest.main()
diff --git a/paimon-python/pypaimon/tests/system/system_table_loader_test.py 
b/paimon-python/pypaimon/tests/system/system_table_loader_test.py
index 3553d647a2..23585cf7b6 100644
--- a/paimon-python/pypaimon/tests/system/system_table_loader_test.py
+++ b/paimon-python/pypaimon/tests/system/system_table_loader_test.py
@@ -33,6 +33,7 @@ _EXPECTED_SYSTEM_TABLES = (
     "branches",
     "file_key_ranges",
     "table_indexes",
+    "consumers",
 )
 
 # Short names recognised by the Paimon catalog that this loader does
@@ -41,7 +42,6 @@ _UNREGISTERED_NAMES = {
     "audit_log",
     "binlog",
     "read_optimized",
-    "consumers",
     "statistics",
     "aggregation_fields",
     "row_tracking",

Reply via email to