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,