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 e8c2178644 [python][daft] Restore IOConfig after source 
deserialization (#8715)
e8c2178644 is described below

commit e8c21786449e7ad7a3d0e929c0dec0c3c92137d0
Author: QuakeWang <[email protected]>
AuthorDate: Mon Jul 20 12:01:47 2026 +0800

    [python][daft] Restore IOConfig after source deserialization (#8715)
---
 paimon-python/pypaimon/daft/daft_datasource.py     | 23 +++++--
 paimon-python/pypaimon/daft/daft_paimon.py         |  2 +-
 .../pypaimon/tests/daft/daft_catalog_rest_test.py  | 59 ++++++++++++++++
 .../pypaimon/tests/daft/daft_data_test.py          | 79 ++++++++++++++++++++++
 4 files changed, 157 insertions(+), 6 deletions(-)

diff --git a/paimon-python/pypaimon/daft/daft_datasource.py 
b/paimon-python/pypaimon/daft/daft_datasource.py
index dff73f02bc..f7232de0d1 100644
--- a/paimon-python/pypaimon/daft/daft_datasource.py
+++ b/paimon-python/pypaimon/daft/daft_datasource.py
@@ -269,15 +269,26 @@ def _load_table(
 
 
 def _build_storage_config(
+    table: FileStoreTable,
     catalog_options: dict[str, Any],
     multithreaded_io: bool,
+    explicit_io_config_bytes: bytes | None,
 ) -> StorageConfig:
     from daft import context
     from daft.daft import StorageConfig
 
     from pypaimon.daft.daft_io_config import 
_convert_paimon_catalog_options_to_io_config
 
-    io_config = _convert_paimon_catalog_options_to_io_config(catalog_options)
+    if explicit_io_config_bytes is not None:
+        from daft.io import IOConfig
+
+        io_config = IOConfig._from_serialized(explicit_io_config_bytes)
+    else:
+        from pypaimon.daft.daft_paimon import _enrich_options_with_rest_token
+
+        io_config = _convert_paimon_catalog_options_to_io_config(
+            _enrich_options_with_rest_token(catalog_options, table)
+        )
     io_config = io_config or 
context.get_context().daft_planning_config.default_io_config
     return StorageConfig(multithreaded_io, io_config)
 
@@ -577,10 +588,6 @@ class PaimonDataSource(DataSource):
         self._pushed_filters = state.get("_pushed_filters")
         self._paimon_predicate = state["_paimon_predicate"]
         self._remaining_filters = state["_remaining_filters"]
-        self._storage_config = _build_storage_config(
-            self._table_catalog_options,
-            state["_multithreaded_io"],
-        )
 
         table = _load_table(
             self._table_catalog_options,
@@ -588,6 +595,12 @@ class PaimonDataSource(DataSource):
             self._table_path,
             self._table_options,
         )
+        self._storage_config = _build_storage_config(
+            table,
+            self._table_catalog_options,
+            state["_multithreaded_io"],
+            self._explicit_io_config_bytes,
+        )
         self._init_table(table)
 
     def _init_table(self, table: FileStoreTable) -> None:
diff --git a/paimon-python/pypaimon/daft/daft_paimon.py 
b/paimon-python/pypaimon/daft/daft_paimon.py
index 996a5228ff..c4bfa7af17 100644
--- a/paimon-python/pypaimon/daft/daft_paimon.py
+++ b/paimon-python/pypaimon/daft/daft_paimon.py
@@ -144,7 +144,7 @@ def _source_for_table(
     if catalog_options is None:
         catalog_options = {}
 
-    # Keep the caller's io_config as a blob File fallback when nothing else is 
derivable.
+    # Keep the caller's io_config for source restoration and blob File 
fallback.
     explicit_io_config_bytes = serialize_io_config(io_config) if io_config is 
not None else None
 
     io_config = io_config or _convert_paimon_catalog_options_to_io_config(
diff --git a/paimon-python/pypaimon/tests/daft/daft_catalog_rest_test.py 
b/paimon-python/pypaimon/tests/daft/daft_catalog_rest_test.py
index afcecfbd1c..bcd7686901 100644
--- a/paimon-python/pypaimon/tests/daft/daft_catalog_rest_test.py
+++ b/paimon-python/pypaimon/tests/daft/daft_catalog_rest_test.py
@@ -294,6 +294,65 @@ class DaftRestReadTest(RESTBaseTest):
         self.assertEqual(captured["opts"].get("warehouse"), "oss://my-bucket", 
captured["opts"])
         fake_file_io.try_to_refresh_token.assert_called()
 
+    def test_source_deserialization_refreshes_rest_io_config(self):
+        from daft.pickle import dumps, loads
+
+        from pypaimon.daft.daft_paimon import _source_for_table
+
+        initial_token = {
+            "fs.oss.accessKeyId": "initial-key",
+            "fs.oss.accessKeySecret": "initial-secret",
+            "fs.oss.securityToken": "initial-token",
+        }
+        refreshed_token = {
+            "fs.oss.accessKeyId": "refreshed-key",
+            "fs.oss.accessKeySecret": "refreshed-secret",
+            "fs.oss.securityToken": "refreshed-token",
+        }
+        fake_token = MagicMock()
+        fake_token.token = initial_token
+        fake_file_io = MagicMock()
+        fake_file_io.token = fake_token
+
+        catalog_options = {**self.options, "warehouse": "morax_test"}
+        fake_file_io.properties = catalog_options
+        table_path = "oss://my-bucket/db.db/tbl-abc"
+
+        with patch.object(self.table, "file_io", fake_file_io), patch.object(
+            self.table, "table_path", table_path
+        ):
+            source = _source_for_table(self.table, 
catalog_options=catalog_options)
+            serialized = dumps(source)
+
+            fake_token.token = refreshed_token
+            with patch(
+                "pypaimon.daft.daft_datasource._load_table",
+                return_value=self.table,
+            ):
+                restored = loads(serialized)
+
+        initial_oss = source._storage_config.io_config.opendal_backends["oss"]
+        restored_oss = (restored._storage_config.io_config.opendal_backends or 
{}).get(
+            "oss", {}
+        )
+        self.assertEqual(
+            (
+                initial_oss.get("access_key_id"),
+                initial_oss.get("access_key_secret"),
+                initial_oss.get("security_token"),
+            ),
+            ("initial-key", "initial-secret", "initial-token"),
+        )
+        self.assertEqual(
+            (
+                restored_oss.get("access_key_id"),
+                restored_oss.get("access_key_secret"),
+                restored_oss.get("security_token"),
+            ),
+            ("refreshed-key", "refreshed-secret", "refreshed-token"),
+        )
+        self.assertEqual(fake_file_io.try_to_refresh_token.call_count, 2)
+
     def test_enrich_is_noop_when_not_rest_metastore(self):
         from pypaimon.daft.daft_paimon import _enrich_options_with_rest_token
         opts = {"warehouse": "/tmp/x", "metastore": "filesystem"}
diff --git a/paimon-python/pypaimon/tests/daft/daft_data_test.py 
b/paimon-python/pypaimon/tests/daft/daft_data_test.py
index a7a1842ba2..eeb0ecdc73 100644
--- a/paimon-python/pypaimon/tests/daft/daft_data_test.py
+++ b/paimon-python/pypaimon/tests/daft/daft_data_test.py
@@ -249,6 +249,85 @@ def 
test_read_paimon_source_is_serializable(append_only_table):
     assert restored._storage_config.multithreaded_io is False
 
 
+def test_source_serialization_preserves_explicit_io_config_for_native_task(
+    append_only_table, monkeypatch
+):
+    from daft.io import IOConfig, S3Config
+    from daft.io.pushdowns import Pushdowns
+    from daft.io.source import DataSourceTask
+    from daft.pickle import dumps, loads
+
+    from pypaimon.daft.daft_paimon import _source_for_table
+
+    table, warehouse = append_only_table
+    _write_to_paimon(
+        table,
+        pa.table(
+            {
+                "id": [1],
+                "name": ["a"],
+                "value": [1.0],
+                "dt": ["2024-01-01"],
+            }
+        ),
+    )
+
+    explicit_values = (
+        "https://explicit.example";,
+        "explicit-key",
+        "explicit-secret",
+        "explicit-token",
+    )
+    explicit_io_config = IOConfig(
+        s3=S3Config(
+            endpoint_url=explicit_values[0],
+            key_id=explicit_values[1],
+            access_key=explicit_values[2],
+            session_token=explicit_values[3],
+        )
+    )
+    catalog_options = {
+        "warehouse": str(warehouse),
+        "fs.s3.endpoint": "https://catalog.example";,
+        "fs.s3.accessKeyId": "catalog-key",
+        "fs.s3.accessKeySecret": "catalog-secret",
+        "fs.s3.securityToken": "catalog-token",
+    }
+
+    source = _source_for_table(
+        table,
+        catalog_options=catalog_options,
+        io_config=explicit_io_config,
+    )
+    restored = loads(dumps(source))
+
+    captured = {}
+    sentinel = object()
+
+    def capture_parquet(**kwargs):
+        captured.update(kwargs)
+        return sentinel
+
+    monkeypatch.setattr(DataSourceTask, "parquet", capture_parquet)
+
+    async def first_task():
+        async for task in restored.get_tasks(Pushdowns()):
+            return task
+        raise AssertionError("Expected at least one native task")
+
+    assert asyncio.run(first_task()) is sentinel
+    assert captured["storage_config"] is restored._storage_config
+
+    for storage_config in (source._storage_config, restored._storage_config):
+        s3 = storage_config.io_config.s3
+        assert (
+            s3.endpoint_url,
+            s3.key_id,
+            s3.access_key,
+            s3.session_token,
+        ) == explicit_values
+
+
 def 
test_read_paimon_source_serialization_preserves_pushed_filter_for_fallback(local_paimon_catalog):
     """A serialized source must keep filters accepted by 
SupportsPushdownFilters."""
     from daft import context, runners

Reply via email to