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 3d25b62e5b [python] Honor data-file.path-directory when resolving data 
files (#10158)
3d25b62e5b is described below

commit 3d25b62e5b8ae7209a0b954ca9580263a47aac47
Author: jackylee <[email protected]>
AuthorDate: Fri Sep 25 21:10:52 2026 +0800

    [python] Honor data-file.path-directory when resolving data files (#10158)
---
 .../pypaimon/common/options/core_options.py        |  10 ++
 .../pypaimon/manifest/schema/data_file_meta.py     |   5 +-
 .../read/scanner/chunk_shuffle_split_generator.py  |   1 +
 .../read/scanner/data_evolution_split_generator.py |   3 +-
 .../pypaimon/read/scanner/split_generator.py       |   3 +-
 paimon-python/pypaimon/read/table_read.py          |   5 +
 paimon-python/pypaimon/read/table_scan.py          |   9 +
 paimon-python/pypaimon/table/file_store_table.py   |   2 +-
 .../tests/data_evolution_group_stats_test.py       |   3 +
 .../tests/data_evolution_split_generator_test.py   |   3 +
 .../tests/data_file_path_directory_test.py         | 192 +++++++++++++++++++++
 paimon-python/pypaimon/tests/native_plan_test.py   |   1 +
 paimon-python/pypaimon/tests/native_read_test.py   |   1 +
 paimon-python/pypaimon/write/table_commit.py       |   5 +
 paimon-python/pypaimon/write/write_builder.py      |   6 +
 15 files changed, 245 insertions(+), 4 deletions(-)

diff --git a/paimon-python/pypaimon/common/options/core_options.py 
b/paimon-python/pypaimon/common/options/core_options.py
index 84fd675349..a1fb6b2abf 100644
--- a/paimon-python/pypaimon/common/options/core_options.py
+++ b/paimon-python/pypaimon/common/options/core_options.py
@@ -520,6 +520,13 @@ class CoreOptions:
         .default_value("data-")
         .with_description("Specify the file name prefix of data files.")
     )
+
+    DATA_FILE_PATH_DIRECTORY: ConfigOption[str] = (
+        ConfigOptions.key("data-file.path-directory")
+        .string_type()
+        .no_default_value()
+        .with_description("Specify the path directory of data files.")
+    )
     # Scan options
     SCAN_MODE: ConfigOption[StartupMode] = (
         ConfigOptions.key("scan.mode")
@@ -1479,6 +1486,9 @@ class CoreOptions:
     def data_file_prefix(self, default=None):
         return self.options.get(CoreOptions.DATA_FILE_PREFIX, default)
 
+    def data_file_path_directory(self, default=None):
+        return self.options.get(CoreOptions.DATA_FILE_PATH_DIRECTORY, default)
+
     def scan_mode(self, default=None):
         return self.options.get(CoreOptions.SCAN_MODE, default)
 
diff --git a/paimon-python/pypaimon/manifest/schema/data_file_meta.py 
b/paimon-python/pypaimon/manifest/schema/data_file_meta.py
index a315c5157f..6fd0a6a4f2 100644
--- a/paimon-python/pypaimon/manifest/schema/data_file_meta.py
+++ b/paimon-python/pypaimon/manifest/schema/data_file_meta.py
@@ -143,8 +143,11 @@ class DataFileMeta:
 
     def set_file_path(
             self, table_path: str, partition: GenericRow, bucket: int,
-            default_part_value: str = "__DEFAULT_PARTITION__"):
+            default_part_value: str = "__DEFAULT_PARTITION__",
+            data_file_path_directory: Optional[str] = None):
         path_builder = table_path.rstrip('/')
+        if data_file_path_directory:
+            path_builder = f"{path_builder}/{data_file_path_directory}"
         partition_dict = partition.to_dict()
         for field_name, field_value in partition_dict.items():
             part_value = default_part_value if 
_is_null_or_whitespace_only(field_value) else str(field_value)
diff --git 
a/paimon-python/pypaimon/read/scanner/chunk_shuffle_split_generator.py 
b/paimon-python/pypaimon/read/scanner/chunk_shuffle_split_generator.py
index 6ed8846afe..eb007ae8e3 100644
--- a/paimon-python/pypaimon/read/scanner/chunk_shuffle_split_generator.py
+++ b/paimon-python/pypaimon/read/scanner/chunk_shuffle_split_generator.py
@@ -251,6 +251,7 @@ class 
ChunkShuffleSplitGeneratorBase(AbstractSplitGenerator):
                         partition_row,
                         bucket,
                         self.default_part_value,
+                        self.table.options.data_file_path_directory(),
                     )
             for segments in self._slice_group_into_chunks(entries_in_group):
                 all_chunks.append(_Chunk(partition_row, bucket, segments))
diff --git 
a/paimon-python/pypaimon/read/scanner/data_evolution_split_generator.py 
b/paimon-python/pypaimon/read/scanner/data_evolution_split_generator.py
index 8da8b57a6d..089f404f48 100644
--- a/paimon-python/pypaimon/read/scanner/data_evolution_split_generator.py
+++ b/paimon-python/pypaimon/read/scanner/data_evolution_split_generator.py
@@ -139,7 +139,8 @@ class DataEvolutionSplitGenerator(AbstractSplitGenerator):
                     self.table.table_path,
                     file_entries[0].partition,
                     file_entries[0].bucket,
-                    self.default_part_value
+                    self.default_part_value,
+                    self.table.options.data_file_path_directory()
                 )
 
             if file_group:
diff --git a/paimon-python/pypaimon/read/scanner/split_generator.py 
b/paimon-python/pypaimon/read/scanner/split_generator.py
index c491fef39a..9607fec6f9 100644
--- a/paimon-python/pypaimon/read/scanner/split_generator.py
+++ b/paimon-python/pypaimon/read/scanner/split_generator.py
@@ -117,7 +117,8 @@ class AbstractSplitGenerator(ABC):
                     self.table.table_path,
                     file_entries[0].partition,
                     file_entries[0].bucket,
-                    self.default_part_value
+                    self.default_part_value,
+                    self.table.options.data_file_path_directory()
                 )
                 if escaped_partition and not data_file.external_path:
                     canonical_path = canonical_data_file_path(
diff --git a/paimon-python/pypaimon/read/table_read.py 
b/paimon-python/pypaimon/read/table_read.py
index eb5e64752d..d8fbf2ea52 100644
--- a/paimon-python/pypaimon/read/table_read.py
+++ b/paimon-python/pypaimon/read/table_read.py
@@ -426,6 +426,11 @@ class TableRead:
         """Return Rust-read batches, or ``None`` when this read must fall 
back."""
         if not self.table.options.native_read_enabled():
             return None
+        # data-file.path-directory relocates data files under a sub-directory
+        # the native reader resolves at the bucket root -- it would 404. The
+        # Python reader honors the directory, matching the write/plan fallback.
+        if self.table.options.data_file_path_directory() is not None:
+            return None
         if self.table.options.file_format() not in _NATIVE_READ_FILE_FORMATS:
             return None
         if not splits:
diff --git a/paimon-python/pypaimon/read/table_scan.py 
b/paimon-python/pypaimon/read/table_scan.py
index 344e1a1703..4bddf3345b 100755
--- a/paimon-python/pypaimon/read/table_scan.py
+++ b/paimon-python/pypaimon/read/table_scan.py
@@ -106,6 +106,15 @@ class TableScan:
         )
         if not native_runtime_available():
             return False
+        # ``data-file.path-directory`` relocates data files under a
+        # sub-directory that only the Python write/plan paths resolve (see
+        # FileStoreTable and the split generators). The native planner still
+        # resolves files at the bucket root, so a native plan would read the
+        # wrong location and fail with NotFound. Fall back to the Python
+        # scanner -- which honors the directory -- until the native runtime
+        # learns this option.
+        if self.table.options.data_file_path_directory() is not None:
+            return False
         fs = self.file_scanner
         if not self._native_global_index_result_supported():
             return False
diff --git a/paimon-python/pypaimon/table/file_store_table.py 
b/paimon-python/pypaimon/table/file_store_table.py
index c939a34586..d0ec91f960 100644
--- a/paimon-python/pypaimon/table/file_store_table.py
+++ b/paimon-python/pypaimon/table/file_store_table.py
@@ -392,7 +392,7 @@ class FileStoreTable(Table):
             
legacy_partition_name=self.options.options.get(CoreOptions.PARTITION_GENERATE_LEGACY_NAME),
             file_suffix_include_compression=False,
             file_compression=file_compression,
-            data_file_path_directory=None,
+            data_file_path_directory=self.options.data_file_path_directory(),
             external_paths=external_paths,
             
external_path_strategy=self.options.data_file_external_paths_strategy(),
             
external_path_weights=self.options.data_file_external_paths_weights(),
diff --git a/paimon-python/pypaimon/tests/data_evolution_group_stats_test.py 
b/paimon-python/pypaimon/tests/data_evolution_group_stats_test.py
index 4b9fbd0495..80826a7f0d 100644
--- a/paimon-python/pypaimon/tests/data_evolution_group_stats_test.py
+++ b/paimon-python/pypaimon/tests/data_evolution_group_stats_test.py
@@ -392,6 +392,9 @@ class 
DataEvolutionGroupStatsPlanningTest(unittest.TestCase):
         class _Options:
             options = {}
 
+            def data_file_path_directory(self, default=None):
+                return default
+
         class _Table:
             table_path = '/tmp/table'
             options = _Options()
diff --git 
a/paimon-python/pypaimon/tests/data_evolution_split_generator_test.py 
b/paimon-python/pypaimon/tests/data_evolution_split_generator_test.py
index 88fc69fb31..4b3c739070 100644
--- a/paimon-python/pypaimon/tests/data_evolution_split_generator_test.py
+++ b/paimon-python/pypaimon/tests/data_evolution_split_generator_test.py
@@ -120,6 +120,9 @@ class SplitOrderTest(unittest.TestCase):
     class _Options:
         options = {}
 
+        def data_file_path_directory(self, default=None):
+            return default
+
     class _Table:
         table_path = '/table'
         options = None
diff --git a/paimon-python/pypaimon/tests/data_file_path_directory_test.py 
b/paimon-python/pypaimon/tests/data_file_path_directory_test.py
new file mode 100644
index 0000000000..c2fbccd925
--- /dev/null
+++ b/paimon-python/pypaimon/tests/data_file_path_directory_test.py
@@ -0,0 +1,192 @@
+# 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.
+
+"""``data-file.path-directory`` places data files under a sub-directory."""
+
+import os
+import shutil
+import tempfile
+import unittest
+from unittest import mock
+
+import pyarrow as pa
+
+from pypaimon import CatalogFactory, Schema
+
+
+class DataFilePathDirectoryTest(unittest.TestCase):
+
+    def setUp(self):
+        self.tmp = tempfile.mkdtemp(prefix="data_file_dir_")
+        warehouse = os.path.join(self.tmp, "warehouse")
+        self.catalog = CatalogFactory.create({"warehouse": warehouse})
+        self.catalog.create_database("db", False)
+
+    def tearDown(self):
+        shutil.rmtree(self.tmp, ignore_errors=True)
+
+    def _write(self, identifier):
+        table = self.catalog.get_table(identifier)
+        wb = table.new_batch_write_builder()
+        writer = wb.new_write()
+        commit = wb.new_commit()
+        writer.write_arrow(pa.table({
+            "id": pa.array([1, 2, 3], type=pa.int32()),
+            "v": ["a", "b", "c"],
+        }))
+        commit.commit(writer.prepare_commit())
+        writer.close()
+        commit.close()
+        return table
+
+    @staticmethod
+    def _data_files(root):
+        found = []
+        for dirpath, _, filenames in os.walk(root):
+            for name in filenames:
+                if name.startswith("data-"):
+                    found.append(os.path.relpath(
+                        os.path.join(dirpath, name), root))
+        return found
+
+    def _create_with_options(self, identifier, options):
+        self.catalog.create_table(
+            identifier,
+            Schema(fields=Schema.from_pyarrow_schema(pa.schema([
+                ("id", pa.int32()), ("v", pa.string())])).fields,
+                options=options),
+            False,
+        )
+        return self.catalog.get_table(identifier)
+
+    def test_data_files_written_under_configured_directory(self):
+        self.catalog.create_table(
+            "db.t",
+            Schema(fields=Schema.from_pyarrow_schema(pa.schema([
+                ("id", pa.int32()), ("v", pa.string())])).fields,
+                options={"data-file.path-directory": "data"}),
+            False,
+        )
+        table = self._write("db.t")
+
+        data_files = self._data_files(str(table.table_path))
+        self.assertTrue(data_files, "no data files were written")
+        for rel in data_files:
+            self.assertEqual(
+                "data", rel.split(os.sep)[0],
+                "data file not under the configured directory: {}".format(rel))
+
+        # The write is readable back through the same path factory.
+        read_builder = table.new_read_builder()
+        splits = read_builder.new_scan().plan().splits()
+        result = read_builder.new_read().to_arrow(splits)
+        self.assertEqual(sorted(result.column("id").to_pylist()), [1, 2, 3])
+
+    def test_native_plan_falls_back_when_directory_configured(self):
+        # ``data-file.path-directory`` only resolves under Python planning;
+        # the native planner still looks under the bucket root and would 404.
+        # Even with scan.native-plan.enabled the scan must refuse native
+        # planning and fall back to Python, so the relocated files are still
+        # found. Regression for the Python/Rust Plan divergence raised in
+        # review (native read failed with NotFound before the guard).
+        self.catalog.create_table(
+            "db.t_native",
+            Schema(fields=Schema.from_pyarrow_schema(pa.schema([
+                ("id", pa.int32()), ("v", pa.string())])).fields,
+                options={"data-file.path-directory": "data",
+                         "scan.native-plan.enabled": "true"}),
+            False,
+        )
+        table = self._write("db.t_native")
+
+        scan = table.new_read_builder().new_scan()
+        # Force the runtime probe to succeed so the assertion isolates the
+        # directory guard (not merely a missing pypaimon-rust): the gate must
+        # still refuse native planning because the directory is configured.
+        with mock.patch(
+                "pypaimon.read.native_plan.native_runtime_available",
+                return_value=True):
+            self.assertFalse(scan._native_plan_supported())
+
+        # End-to-end: a native-requested scan returns the rows through the
+        # Python fallback rather than failing to locate the relocated files.
+        splits = scan.plan().splits()
+        result = table.new_read_builder().new_read().to_arrow(splits)
+        self.assertEqual(sorted(result.column("id").to_pylist()), [1, 2, 3])
+
+    def test_native_write_falls_back_when_directory_configured(self):
+        # write.native.enabled + data-file.path-directory: the write builder
+        # must return the Python writer (native probe returns None) so the
+        # relocated directory is honored; the native writer writes at the root.
+        # Patch the native constructor so the guard, not a missing runtime, is
+        # what forces the fallback (without the guard this returns the 
sentinel).
+        table = self._create_with_options(
+            "db.t_native_write",
+            {"data-file.path-directory": "data",
+             "write.native.enabled": "true"})
+        with mock.patch(
+                "pypaimon.write.native_write.create_native_write",
+                return_value=object()):
+            self.assertIsNone(table.new_batch_write_builder()._native_write())
+
+    def test_native_read_falls_back_when_directory_configured(self):
+        # read.native.enabled + data-file.path-directory: the native read
+        # probe short-circuits to the Python reader before loading the Rust
+        # runtime, so the relocated files are resolved.
+        table = self._create_with_options(
+            "db.t_native_read",
+            {"data-file.path-directory": "data",
+             "read.native.enabled": "true"})
+        self._write("db.t_native_read")
+        rb = table.new_read_builder()
+        splits = rb.new_scan().plan().splits()
+        self.assertIsNone(rb.new_read()._try_native_batches(
+            splits, pa.schema([("id", pa.int32()), ("v", pa.string())])))
+
+    def test_native_commit_falls_back_when_directory_configured(self):
+        # commit.native.enabled + data-file.path-directory: the native commit
+        # probe returns None so the Python committer records the relocated
+        # paths.
+        table = self._create_with_options(
+            "db.t_native_commit",
+            {"data-file.path-directory": "data",
+             "commit.native.enabled": "true"})
+        commit = table.new_batch_write_builder().new_commit()
+        try:
+            self.assertIsNone(commit._prepare_native_commit([]))
+        finally:
+            commit.close()
+
+    def test_default_keeps_data_files_at_bucket_root(self):
+        self.catalog.create_table(
+            "db.t_default",
+            Schema(fields=Schema.from_pyarrow_schema(pa.schema([
+                ("id", pa.int32()), ("v", pa.string())])).fields),
+            False,
+        )
+        table = self._write("db.t_default")
+
+        data_files = self._data_files(str(table.table_path))
+        self.assertTrue(data_files, "no data files were written")
+        for rel in data_files:
+            self.assertTrue(
+                rel.split(os.sep)[0].startswith("bucket-"),
+                "unexpected data file location: {}".format(rel))
+
+
+if __name__ == "__main__":
+    unittest.main()
diff --git a/paimon-python/pypaimon/tests/native_plan_test.py 
b/paimon-python/pypaimon/tests/native_plan_test.py
index 285c0a40fa..5c6c6078d6 100644
--- a/paimon-python/pypaimon/tests/native_plan_test.py
+++ b/paimon-python/pypaimon/tests/native_plan_test.py
@@ -59,6 +59,7 @@ def _scan(native_enabled, file_scanner):
     scan.table.options.options.contains.return_value = False       # no 
incremental
     scan.table.options.merge_engine.return_value = None            # not 
first-row
     scan.table.options.query_auth_enabled = False
+    scan.table.options.data_file_path_directory.return_value = None  # no 
relocated dir
     scan.table.current_branch.return_value = 'main'
     scan.table.is_primary_key_table = False        # not a pk table
     scan.table.trimmed_primary_keys = ['k']        # non-empty trimmed pk
diff --git a/paimon-python/pypaimon/tests/native_read_test.py 
b/paimon-python/pypaimon/tests/native_read_test.py
index 8aa85566f4..e078b640f9 100644
--- a/paimon-python/pypaimon/tests/native_read_test.py
+++ b/paimon-python/pypaimon/tests/native_read_test.py
@@ -36,6 +36,7 @@ def _table_read(limit=None):
     read.table = Mock()
     read.table.options.native_read_enabled.return_value = True
     read.table.options.file_format.return_value = 'parquet'
+    read.table.options.data_file_path_directory.return_value = None
     read.table.options.blob_as_descriptor.return_value = False
     read.table.options.blob_descriptor_fields.return_value = set()
     read.table.options.blob_view_fields.return_value = set()
diff --git a/paimon-python/pypaimon/write/table_commit.py 
b/paimon-python/pypaimon/write/table_commit.py
index ef2595d93c..41472bf358 100644
--- a/paimon-python/pypaimon/write/table_commit.py
+++ b/paimon-python/pypaimon/write/table_commit.py
@@ -117,6 +117,11 @@ class TableCommit:
         if (not self.table.options.native_commit_enabled()
                 or self._commit_callbacks):
             return None
+        # data-file.path-directory keeps the whole pipeline on the Python
+        # path (which resolves the relocated directory); see the matching
+        # write / read / plan fallbacks.
+        if self.table.options.data_file_path_directory() is not None:
+            return None
         try:
             from pypaimon.write.native_commit import (
                 create_native_commit, native_messages_supported,
diff --git a/paimon-python/pypaimon/write/write_builder.py 
b/paimon-python/pypaimon/write/write_builder.py
index 5180b8b08f..0586ff2d73 100644
--- a/paimon-python/pypaimon/write/write_builder.py
+++ b/paimon-python/pypaimon/write/write_builder.py
@@ -56,6 +56,12 @@ class WriteBuilder(ABC):
     def _native_write(self, static_partition=None, stream=False):
         if not self.table.options.native_write_enabled():
             return None
+        # data-file.path-directory relocates data files under a sub-directory
+        # that the native writer does not honor (it writes at the bucket
+        # root). Use the Python writer, which resolves the directory, so
+        # write / read / plan / commit stay consistent for this option.
+        if self.table.options.data_file_path_directory() is not None:
+            return None
         try:
             from pypaimon.write.native_write import create_native_write
             return create_native_write(self.table, self.commit_user,

Reply via email to