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 81f7bbe71c [python] Propagate dependency read context in REST headers
(#8908)
81f7bbe71c is described below
commit 81f7bbe71c6146e190f48d05b8408c1e2fd63b45
Author: YeJunHao <[email protected]>
AuthorDate: Thu Jul 30 11:59:04 2026 +0800
[python] Propagate dependency read context in REST headers (#8908)
---
paimon-python/pypaimon/api/rest_api.py | 1 +
.../pypaimon/catalog/catalog_environment.py | 42 +++++
paimon-python/pypaimon/catalog/catalog_factory.py | 17 ++
.../pypaimon/tests/catalog_environment_test.py | 202 +++++++++++++++++++++
paimon-python/pypaimon/utils/blob_view_lookup.py | 10 +-
5 files changed, 271 insertions(+), 1 deletion(-)
diff --git a/paimon-python/pypaimon/api/rest_api.py
b/paimon-python/pypaimon/api/rest_api.py
index 78afe0d914..e1ee43a616 100755
--- a/paimon-python/pypaimon/api/rest_api.py
+++ b/paimon-python/pypaimon/api/rest_api.py
@@ -58,6 +58,7 @@ from pypaimon.snapshot.snapshot_commit import
PartitionStatistics
class RESTApi:
HEADER_PREFIX = "header."
+ READ_VIA_HEADER = "X-Paimon-Read-Via"
MAX_RESULTS = "maxResults"
PAGE_TOKEN = "pageToken"
DATABASE_NAME_PATTERN = "databaseNamePattern"
diff --git a/paimon-python/pypaimon/catalog/catalog_environment.py
b/paimon-python/pypaimon/catalog/catalog_environment.py
index 0257a75c43..84530fbc0a 100644
--- a/paimon-python/pypaimon/catalog/catalog_environment.py
+++ b/paimon-python/pypaimon/catalog/catalog_environment.py
@@ -17,8 +17,13 @@
from typing import Optional
+from pypaimon.api.rest_api import RESTApi
+from pypaimon.api.rest_util import RESTUtil
+from pypaimon.catalog.catalog_context import CatalogContext
from pypaimon.catalog.catalog_loader import CatalogLoader
from pypaimon.common.identifier import Identifier
+from pypaimon.common.json_util import JSON
+from pypaimon.common.options.config import CatalogOptions
from pypaimon.snapshot.catalog_snapshot_commit import CatalogSnapshotCommit
from pypaimon.snapshot.renaming_snapshot_commit import RenamingSnapshotCommit
from pypaimon.snapshot.snapshot_commit import SnapshotCommit
@@ -27,6 +32,8 @@ from pypaimon.snapshot.snapshot_loader import SnapshotLoader
class CatalogEnvironment:
+ _READ_VIA_OPTION = RESTApi.HEADER_PREFIX + RESTApi.READ_VIA_HEADER
+
def __init__(
self,
identifier: Optional[Identifier] = None,
@@ -86,6 +93,41 @@ class CatalogEnvironment:
return SnapshotLoader(self.catalog_loader, self.identifier)
return None
+ def catalog_context(self) -> Optional[CatalogContext]:
+ if self.catalog_loader is None:
+ return None
+ context = getattr(self.catalog_loader, "context", None)
+ return context() if callable(context) else None
+
+ def dependency_read_context(self) -> Optional[CatalogContext]:
+ context = self.catalog_context()
+ if self.identifier is None or context is None:
+ return context
+
+ from pypaimon.catalog.rest.rest_catalog_loader import RESTCatalogLoader
+
+ rest_catalog = (
+ isinstance(self.catalog_loader, RESTCatalogLoader)
+ or context.options.get(CatalogOptions.METASTORE) == "rest"
+ )
+ if not rest_catalog:
+ return context
+ if context.options.contains_key(self._READ_VIA_OPTION):
+ return context
+
+ dependency_options = context.options.copy()
+ if not dependency_options.contains(CatalogOptions.METASTORE):
+ dependency_options.set(CatalogOptions.METASTORE, "rest")
+ dependency_options.to_map()[self._READ_VIA_OPTION] =
RESTUtil.encode_string(
+ JSON.to_json(self.identifier, separators=(",", ":"))
+ )
+ return CatalogContext.create(
+ dependency_options,
+ context.hadoop_conf,
+ context.prefer_io_loader,
+ context.fallback_io_loader,
+ )
+
def copy(self, identifier: Identifier) -> 'CatalogEnvironment':
"""
Create a copy of this CatalogEnvironment with a different identifier.
diff --git a/paimon-python/pypaimon/catalog/catalog_factory.py
b/paimon-python/pypaimon/catalog/catalog_factory.py
index 117c10b2e9..b28de94e70 100644
--- a/paimon-python/pypaimon/catalog/catalog_factory.py
+++ b/paimon-python/pypaimon/catalog/catalog_factory.py
@@ -44,3 +44,20 @@ class CatalogFactory:
if identifier in ("jdbc", "rest"):
return
catalog_class(CatalogContext.create_from_options(Options(catalog_options)))
return catalog_class(Options(catalog_options))
+
+ @staticmethod
+ def create_from_context(
+ context: CatalogContext,
+ config_required: bool = True
+ ) -> Catalog:
+ identifier = context.options.get(CatalogOptions.METASTORE)
+ catalog_class = CatalogFactory.CATALOG_REGISTRY.get(identifier)
+ if catalog_class is None:
+ raise ValueError(
+ "Unknown catalog identifier: {}. Available types: {}".format(
+ identifier, list(CatalogFactory.CATALOG_REGISTRY.keys())))
+ if identifier == "filesystem":
+ return catalog_class(context.options)
+ if identifier == "rest":
+ return catalog_class(context, config_required=config_required)
+ return catalog_class(context)
diff --git a/paimon-python/pypaimon/tests/catalog_environment_test.py
b/paimon-python/pypaimon/tests/catalog_environment_test.py
new file mode 100644
index 0000000000..3908b6cb13
--- /dev/null
+++ b/paimon-python/pypaimon/tests/catalog_environment_test.py
@@ -0,0 +1,202 @@
+# 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.
+
+import unittest
+from types import SimpleNamespace
+from unittest import mock
+
+from pypaimon.api.rest_api import RESTApi
+from pypaimon.api.rest_util import RESTUtil
+from pypaimon.catalog.catalog_context import CatalogContext
+from pypaimon.catalog.catalog_factory import CatalogFactory
+from pypaimon.catalog.catalog_environment import CatalogEnvironment
+from pypaimon.catalog.filesystem_catalog_loader import FileSystemCatalogLoader
+from pypaimon.catalog.rest.rest_catalog_loader import RESTCatalogLoader
+from pypaimon.common.identifier import Identifier
+from pypaimon.common.json_util import JSON
+from pypaimon.common.options import Options
+from pypaimon.common.options.config import CatalogOptions
+from pypaimon.utils.blob_view_lookup import BlobViewLookup
+
+
+class CatalogEnvironmentTest(unittest.TestCase):
+
+ _READ_VIA_OPTION = RESTApi.HEADER_PREFIX + RESTApi.READ_VIA_HEADER
+
+ def test_dependency_read_context_for_rest_catalog(self):
+ root = Identifier.create("db", "root", branch="dev")
+ options = Options({"other-option": "value"})
+ context = CatalogContext.create_from_options(options)
+ environment = CatalogEnvironment(
+ identifier=root,
+ catalog_loader=RESTCatalogLoader(context),
+ )
+
+ dependency_context = environment.dependency_read_context()
+
+ self.assertIsNot(dependency_context, context)
+ self.assertFalse(context.options.contains_key(self._READ_VIA_OPTION))
+ self.assertEqual(
+ dependency_context.options.get(CatalogOptions.METASTORE), "rest")
+ self.assertEqual(
+ dependency_context.options.to_map()["other-option"], "value")
+ read_via = JSON.from_json(
+ RESTUtil.decode_string(
+ dependency_context.options.to_map()[self._READ_VIA_OPTION]),
+ Identifier,
+ )
+ self.assertEqual(read_via, root)
+
+ def test_dependency_read_context_preserves_outermost_table(self):
+ outermost = Identifier.create("db", "outermost")
+ read_via = RESTUtil.encode_string(
+ JSON.to_json(outermost, separators=(",", ":")))
+ context = CatalogContext.create_from_options(Options({
+ CatalogOptions.METASTORE.key(): "rest",
+ self._READ_VIA_OPTION: read_via,
+ }))
+ environment = CatalogEnvironment(
+ identifier=Identifier.create("db", "intermediate"),
+ catalog_loader=RESTCatalogLoader(context),
+ )
+
+ self.assertIs(environment.dependency_read_context(), context)
+ self.assertEqual(
+ context.options.to_map()[self._READ_VIA_OPTION], read_via)
+
+ def test_dependency_read_context_does_not_affect_other_catalogs(self):
+ context = CatalogContext.create_from_options(Options({}))
+ environment = CatalogEnvironment(
+ identifier=Identifier.create("db", "table"),
+ catalog_loader=FileSystemCatalogLoader(context),
+ )
+
+ self.assertIs(environment.dependency_read_context(), context)
+ self.assertFalse(context.options.contains_key(self._READ_VIA_OPTION))
+
+ def test_dependency_read_context_for_external_rest_table(self):
+ context = CatalogContext.create_from_options(Options({
+ CatalogOptions.METASTORE.key(): "rest",
+ }))
+ environment = CatalogEnvironment(
+ identifier=Identifier.create("db", "external"),
+ catalog_loader=FileSystemCatalogLoader(context),
+ )
+
+ self.assertIsNot(environment.dependency_read_context(), context)
+
+ def test_dependency_read_context_preserves_custom_rest_metastore(self):
+ context = CatalogContext.create_from_options(Options({
+ CatalogOptions.METASTORE.key(): "custom-rest",
+ }))
+ environment = CatalogEnvironment(
+ identifier=Identifier.create("db", "table"),
+ catalog_loader=RESTCatalogLoader(context),
+ )
+
+ dependency_context = environment.dependency_read_context()
+
+ self.assertIsNot(dependency_context, context)
+ self.assertEqual(
+ dependency_context.options.get(CatalogOptions.METASTORE),
+ "custom-rest",
+ )
+
+ def test_catalog_factory_creates_rest_from_context(self):
+ context = CatalogContext.create_from_options(Options({
+ CatalogOptions.METASTORE.key(): "rest",
+ }))
+ dependency_catalog = mock.sentinel.dependency_catalog
+ rest_catalog = mock.Mock(return_value=dependency_catalog)
+
+ with mock.patch.dict(
+ CatalogFactory.CATALOG_REGISTRY, {"rest": rest_catalog}):
+ result = CatalogFactory.create_from_context(
+ context, config_required=False)
+
+ rest_catalog.assert_called_once_with(context, config_required=False)
+ self.assertIs(result, dependency_catalog)
+
+ def test_blob_view_lookup_loads_dependency_catalog(self):
+ root = Identifier.create("db", "root")
+ target = Identifier.create("db", "target")
+ context = CatalogContext.create_from_options(Options({
+ CatalogOptions.METASTORE.key(): "rest",
+ }))
+ original_loader = RESTCatalogLoader(context)
+ environment = CatalogEnvironment(
+ identifier=root,
+ catalog_loader=original_loader,
+ )
+ table = SimpleNamespace(catalog_environment=environment)
+ dependency_catalog = mock.MagicMock()
+ dependency_table = mock.sentinel.dependency_table
+ dependency_catalog.get_table.return_value = dependency_table
+
+ with mock.patch.object(
+ CatalogFactory,
+ "create_from_context",
+ return_value=dependency_catalog) as create_catalog:
+ result = BlobViewLookup(table)._load_table(target)
+
+ dependency_context = create_catalog.call_args.args[0]
+ create_catalog.assert_called_once_with(
+ dependency_context, config_required=False)
+ self.assertTrue(
+ dependency_context.options.contains_key(self._READ_VIA_OPTION))
+ dependency_catalog.get_table.assert_called_once_with(target)
+ self.assertIs(result, dependency_table)
+
+ def test_blob_view_lookup_preserves_custom_rest_catalog(self):
+ target = Identifier.create("db", "target")
+ dependency_table = mock.sentinel.dependency_table
+
+ class CustomRESTCatalog:
+ context = None
+
+ def __init__(self, context):
+ CustomRESTCatalog.context = context
+
+ def get_table(self, identifier):
+ self.identifier = identifier
+ return dependency_table
+
+ context = CatalogContext.create_from_options(Options({
+ CatalogOptions.METASTORE.key(): "custom-rest",
+ }))
+ environment = CatalogEnvironment(
+ identifier=Identifier.create("db", "root"),
+ catalog_loader=RESTCatalogLoader(context),
+ )
+ table = SimpleNamespace(catalog_environment=environment)
+
+ with mock.patch.dict(
+ CatalogFactory.CATALOG_REGISTRY,
+ {"custom-rest": CustomRESTCatalog}):
+ result = BlobViewLookup(table)._load_table(target)
+
+ dependency_context = CustomRESTCatalog.context
+ self.assertEqual(
+ dependency_context.options.get(CatalogOptions.METASTORE),
+ "custom-rest",
+ )
+ self.assertTrue(
+ dependency_context.options.contains_key(self._READ_VIA_OPTION))
+ self.assertIs(result, dependency_table)
+
+
+if __name__ == "__main__":
+ unittest.main()
diff --git a/paimon-python/pypaimon/utils/blob_view_lookup.py
b/paimon-python/pypaimon/utils/blob_view_lookup.py
index 2e6fd1e8be..03867b9885 100644
--- a/paimon-python/pypaimon/utils/blob_view_lookup.py
+++ b/paimon-python/pypaimon/utils/blob_view_lookup.py
@@ -290,7 +290,15 @@ class BlobViewLookup:
return max(_MIN_ROWS_PER_TASK, (total_rows + _PRELOAD_THREAD_NUM - 1)
// _PRELOAD_THREAD_NUM)
def _load_table(self, identifier: Identifier):
- catalog = self._table.catalog_environment.catalog_loader.load()
+ catalog_environment = self._table.catalog_environment
+ catalog_loader = catalog_environment.catalog_loader
+ dependency_context = catalog_environment.dependency_read_context()
+ if dependency_context is catalog_environment.catalog_context():
+ catalog = catalog_loader.load()
+ else:
+ from pypaimon.catalog.catalog_factory import CatalogFactory
+ catalog = CatalogFactory.create_from_context(
+ dependency_context, config_required=False)
return catalog.get_table(identifier)
@staticmethod