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",