Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
10 changes: 6 additions & 4 deletions paimon-python/pypaimon/common/external_path_provider.py
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,8 @@
from abc import ABC, abstractmethod
from typing import List

from pypaimon.utils.path import resolve_path, to_file_io_path


class ExternalPathProvider(ABC):
"""Provider for external data paths."""
Expand Down Expand Up @@ -71,7 +73,7 @@ def get_next_external_data_path(self, file_name: str) -> str:

external_base = self._external_table_paths[self._position]
if self._relative_bucket_path:
return f"{external_base.rstrip('/')}/{self._relative_bucket_path.strip('/')}/{file_name}"
return to_file_io_path(resolve_path(resolve_path(external_base, self._relative_bucket_path), file_name))
else:
return f"{external_base.rstrip('/')}/{file_name}"

Expand All @@ -97,7 +99,7 @@ def __init__(self, external_table_paths: List[str], relative_bucket_path: str =
def get_next_external_data_path(self, file_name: str) -> str:
hash_dirs = self._compute_hash(file_name)
if self._relative_bucket_path:
file_path_with_hash = f"{self._relative_bucket_path.strip('/')}/{hash_dirs}/{file_name}"
file_path_with_hash = resolve_path(self._relative_bucket_path, f"{hash_dirs}/{file_name}")
else:
file_path_with_hash = f"{hash_dirs}/{file_name}"

Expand All @@ -106,7 +108,7 @@ def get_next_external_data_path(self, file_name: str) -> str:
self._position = 0

external_base = self._external_table_paths[self._position]
return f"{external_base.rstrip('/')}/{file_path_with_hash}"
return to_file_io_path(resolve_path(external_base, file_path_with_hash))

def _compute_hash(self, file_name: str) -> str:
hash_int = _murmur3_32(file_name.encode('utf-8'))
Expand Down Expand Up @@ -151,7 +153,7 @@ def get_next_external_data_path(self, file_name: str) -> str:
index = len(self._external_table_paths) - 1
selected_base = self._external_table_paths[index]
if self._relative_bucket_path:
return f"{selected_base.rstrip('/')}/{self._relative_bucket_path.strip('/')}/{file_name}"
return to_file_io_path(resolve_path(resolve_path(selected_base, self._relative_bucket_path), file_name))
else:
return f"{selected_base.rstrip('/')}/{file_name}"

Expand Down
7 changes: 4 additions & 3 deletions paimon-python/pypaimon/manifest/schema/data_file_meta.py
Original file line number Diff line number Diff line change
Expand Up @@ -145,9 +145,10 @@ def set_file_path(
self, table_path: str, partition: GenericRow, bucket: int,
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}"
from pypaimon.utils.path import resolve_path, to_file_io_path
path_builder = resolve_path(table_path, data_file_path_directory)
if data_file_path_directory is not None:
path_builder = to_file_io_path(path_builder)
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)
Expand Down
6 changes: 4 additions & 2 deletions paimon-python/pypaimon/read/split_serializer.py
Original file line number Diff line number Diff line change
Expand Up @@ -36,6 +36,7 @@
from pypaimon.table.row.generic_row import (
GenericRow, GenericRowDeserializer, GenericRowSerializer)
from pypaimon.table.source.deletion_file import DeletionFile
from pypaimon.utils.path import to_file_io_path
from pypaimon.utils.range import Range

# Frame magics/versions, mirroring the Rust/Java constants.
Expand Down Expand Up @@ -654,7 +655,8 @@ def _datafilemeta_from_row(row_bytes: bytes, bucket_path: str, arity: int,
write_cols_sequences=(
_decode_non_null_long_array(g(20)) if arity >= 21 else None),
)
meta.file_path = external_path if external_path else "%s/%s" % (bucket_path.rstrip('/'), file_name)
meta.file_path = external_path if external_path else to_file_io_path(
"%s/%s" % (bucket_path.rstrip('/'), file_name))
return meta


Expand Down Expand Up @@ -691,7 +693,7 @@ def _read_deletion_list(r: _Reader) -> Optional[List[Optional[DeletionFile]]]:
path = r.java_utf()
offset, length, cardinality = r.i64(), r.i64(), r.i64()
result.append(DeletionFile(
dv_index_path=path,
dv_index_path=to_file_io_path(path),
offset=offset,
length=length,
cardinality=None if cardinality == -1 else cardinality,
Expand Down
5 changes: 0 additions & 5 deletions paimon-python/pypaimon/read/table_read.py
Original file line number Diff line number Diff line change
Expand Up @@ -433,11 +433,6 @@ def _try_native_batches(
"""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:
Expand Down
9 changes: 0 additions & 9 deletions paimon-python/pypaimon/read/table_scan.py
Original file line number Diff line number Diff line change
Expand Up @@ -106,15 +106,6 @@ def _native_plan_supported_impl(self) -> bool:
)
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
Expand Down
5 changes: 4 additions & 1 deletion paimon-python/pypaimon/table/file_store_table.py
Original file line number Diff line number Diff line change
Expand Up @@ -534,12 +534,15 @@ def _copy_with_snapshot(self, snapshot):
return table

def _copy(self, options: dict, resolve_time_travel: bool) -> 'FileStoreTable':
directory_key = CoreOptions.DATA_FILE_PATH_DIRECTORY.key()
if directory_key in options and options[directory_key] != self.options.data_file_path_directory():
raise ValueError('Cannot change immutable option ' + directory_key)
if CoreOptions.BUCKET.key() in options and int(options.get(CoreOptions.BUCKET.key())) != self.options.bucket():
raise ValueError("Cannot change bucket number")
new_options = CoreOptions.copy(self.options).options.to_map()
for k, v in options.items():
if v is None:
new_options.pop(k)
new_options.pop(k, None)
else:
new_options[k] = v

Expand Down
194 changes: 128 additions & 66 deletions paimon-python/pypaimon/tests/data_file_path_directory_test.py
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,8 @@
import shutil
import tempfile
import unittest

import pytest
from unittest import mock

import pyarrow as pa
Expand Down Expand Up @@ -96,79 +98,35 @@ def test_data_files_written_under_configured_directory(self):
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")

def test_native_plan_accepts_configured_directory(self):
table = self._create_with_options(
"db.t_native", {"data-file.path-directory": "data"})
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])
with mock.patch("pypaimon.read.native_plan.native_runtime_available", return_value=True):
self.assertTrue(scan._native_plan_supported())

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.
def test_native_write_dispatches_with_configured_directory(self):
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.
"db.t_native_write", {"data-file.path-directory": "data",
"write.native.enabled": "true"})
sentinel = object()
with mock.patch("pypaimon.write.native_write.create_native_write", return_value=sentinel):
self.assertIs(table.new_batch_write_builder()._native_write(), sentinel)

def test_native_commit_dispatches_with_configured_directory(self):
table = self._create_with_options(
"db.t_native_commit",
{"data-file.path-directory": "data",
"commit.native.enabled": "true"})
"db.t_native_commit", {"data-file.path-directory": "data",
"commit.native.enabled": "true"})
commit = table.new_batch_write_builder().new_commit()
sentinel = object()
try:
self.assertIsNone(commit._prepare_native_commit([]))
with mock.patch("pypaimon.write.native_commit.native_messages_supported", return_value=True), \
mock.patch("pypaimon.write.native_commit.create_native_commit", return_value=sentinel), \
mock.patch("pypaimon.write.native_commit.to_native_commit_messages", return_value=[]):
self.assertEqual(commit._prepare_native_commit([]), (sentinel, []))
finally:
# The sentinel represents a prepared committer, not a real resource.
commit._native_commit = None
commit.close()

def test_default_keeps_data_files_at_bucket_root(self):
Expand All @@ -190,3 +148,107 @@ def test_default_keeps_data_files_at_bucket_root(self):

if __name__ == "__main__":
unittest.main()


# Expected strings were produced by org.apache.paimon.fs.Path from this tree.
@pytest.mark.parametrize('parent,child,expected', [
('/warehouse/t', 'data/nested', '/warehouse/t/data/nested'),
('/warehouse/t/', 'data//discard/../nested/.', '/warehouse/t/data/nested'),
('s3://bucket/warehouse/t', '/shared/data', 's3://bucket/shared/data'),
('s3://bucket/t', 's3://other/data', 's3://other/data'),
('file:/warehouse/t', 'file:///other/data', 'file:/other/data'),
('memory:/t', '../other', 'memory:/other'),
('warehouse/t', '../data', 'warehouse/data'),
('warehouse/t', 'data', 'warehouse/t/data'),
('s3://bucket/t', '//other/data', 's3://other/data'),
('s3://bucket/t', 'data 100%?#/你好', 's3://bucket/t/data 100%?#/你好'),
('file:/t', './data', 'file:/t/data'),
('/', 'data', '/data'),
('/t', '.', '/t'),
('/t', '..', '/'),
('/t', '../../data', '/../data'),
('s3://bucket/', 'data', 's3://bucket/data'),
('s3://bucket', '.', 's3://bucket'),
('s3://bucket', 'x/..', 's3://bucket'),
('warehouse', 'a/../b:c', 'warehouse/b:c'),
('memory:/t', 'a/../../b', 'memory:/b'),
('/t', '///data', '/data'),
('/t', '////data', '/data'),
('file:/t', 'file:////data', 'file:/data'),
('/t', '//', '/'),
('file:/t', '//', 'file:/'),
('s3://bucket/t', '//other:9000/data', 's3://other:9000/data'),
('//host:8020/table', '//other/data', '//other/data'),
('relative', '../b:c', './b:c'),
('relative', '../a:bb', './a:bb'),
('relative', '../b:c/x', './b:c/x'),
('file:/t', 'file:/data', 'file:/data'),
])
@pytest.mark.parametrize('windows', [False, True])
def test_path_resolution_matches_java(monkeypatch, parent, child, expected, windows):
import pypaimon.utils.path as paths
monkeypatch.setattr(paths, '_WINDOWS', windows)
assert paths.resolve_path(parent, child) == expected


def test_empty_path_is_rejected():
from pypaimon.utils.path import resolve_path
with pytest.raises(ValueError, match='empty string'):
resolve_path('/warehouse/t', '')


@pytest.mark.parametrize('stored', [None, 'data'])
@pytest.mark.parametrize('requested', [None, 'data', 'other'])
def test_data_directory_cannot_change_on_table_copy(tmp_path, stored, requested):
catalog = CatalogFactory.create({'warehouse': str(tmp_path)})
catalog.create_database('db', True)
options = {} if stored is None else {'data-file.path-directory': stored}
catalog.create_table('db.t', Schema.from_pyarrow_schema(pa.schema([('id', pa.int32())]), options=options), False)
table = catalog.get_table('db.t')
for copy in (table.copy, table.copy_without_time_travel):
if stored == requested:
assert copy({'data-file.path-directory': requested}).options.data_file_path_directory() == stored
else:
with pytest.raises(ValueError, match='immutable option'):
copy({'data-file.path-directory': requested})


@pytest.mark.parametrize('parent,child,expected', [
('relative', '../b:c', './b:c'),
('relative', '../a:bb', './a:bb'),
('relative', '../b:c/x', './b:c/x'),
('warehouse/t', 'data/../../../b:c', './b:c'),
('.', 'a/../b:c', './b:c'),
('s3://bucket', '.', 's3://bucket'),
(r'C:\warehouse\table', '../data', 'C:/warehouse/data'),
(r'C:\warehouse\table', r'..\data', 'C:/warehouse/data'),
('C:/warehouse/table', '/data', '/data'),
('C:/warehouse/table', '../..', 'C:/'),
('file:/C:/warehouse/table', '../..', 'file:/C:/'),
(r'\\server\share\table', 'data', '//server/share/table/data'),
('s3://bucket/table', r'\\other\share', 's3://other/share'),
('file:/C:/warehouse/table', r'..\data', 'file:/C:/warehouse/data'),
('s3://bucket/t', r'data\literal', 's3://bucket/t/data/literal'),
])
def test_windows_paths_match_java(monkeypatch, parent, child, expected):
import pypaimon.utils.path as paths
monkeypatch.setattr(paths, '_WINDOWS', True)
assert paths.resolve_path(parent, child) == expected


@pytest.mark.parametrize('windows,path,expected', [
(False, 'file:/tmp/data%2Fwith space?#part', '/tmp/data%2Fwith space?#part'),
(False, 'file:///tmp/data%20', '/tmp/data%20'),
(False, 'file://localhost/tmp/data%20', '/tmp/data%20'),
(False, 'file://host/share/data%20', '//host/share/data%20'),
(True, 'file:/C:/data%20', 'C:/data%20'),
(True, 'file://C:/data%20', 'C:/data%20'),
(True, 'file://host/share/data%20', '//host/share/data%20'),
(False, '/tmp/data%2F?#', '/tmp/data%2F?#'),
(False, 's3://bucket/data%2F?#', 's3://bucket/data%2F?#'),
(False, 'file:/tmp/plain', 'file:/tmp/plain'),
])
def test_file_io_paths_keep_literal_characters(monkeypatch, windows, path, expected):
import pypaimon.utils.path as paths
monkeypatch.setattr(paths, '_WINDOWS', windows)
assert paths.to_file_io_path(path) == expected
3 changes: 2 additions & 1 deletion paimon-python/pypaimon/tests/deletion_vector_path_test.py
Original file line number Diff line number Diff line change
Expand Up @@ -128,7 +128,8 @@ def test_delete_paths_preserve_repeated_deletes_and_historical_reads(tmp_path, p
else:
expected = Path(table.table_path) / 'index' / file.file_name
assert expected.is_file()
assert file.external_path == ('file://' + str(expected) if 'external' in layout else None)
assert file.external_path == (('file:' if layout == 'bucket-external' else 'file://')
+ str(expected) if 'external' in layout else None)
# Obsolete locations with the same name must never shadow canonical or
# explicit paths. Invalid bytes make a wrong-path read fail observably.
if layout != 'table':
Expand Down
2 changes: 1 addition & 1 deletion paimon-python/pypaimon/tests/external_paths_test.py
Original file line number Diff line number Diff line change
Expand Up @@ -579,7 +579,7 @@ def test_write_with_external_paths(self):
for file_meta in commit_msg.new_files:
# External path should be set
self.assertIsNotNone(file_meta.external_path)
self.assertTrue(file_meta.external_path.startswith("file://"))
self.assertTrue(file_meta.external_path.startswith("file:"))
self.assertIn(self.external_dir, file_meta.external_path)

table_commit.commit(commit_messages)
Expand Down
Loading
Loading