From 6702c813e314b70b6725dbf6817aa375c69047bc Mon Sep 17 00:00:00 2001 From: JingsongLi Date: Sun, 27 Sep 2026 23:47:53 +0800 Subject: [PATCH 1/2] [python] Align data directories with Java and enable Native IO --- .../pypaimon/common/external_path_provider.py | 10 +- .../manifest/schema/data_file_meta.py | 5 +- paimon-python/pypaimon/read/table_read.py | 5 - paimon-python/pypaimon/read/table_scan.py | 9 - .../pypaimon/table/file_store_table.py | 5 +- .../tests/data_file_path_directory_test.py | 164 +++++++++------ .../tests/deletion_vector_path_test.py | 3 +- .../pypaimon/tests/external_paths_test.py | 2 +- .../tests/native_data_directory_test.py | 192 ++++++++++++++++++ .../pypaimon/tests/native_write_test.py | 13 +- .../pypaimon/utils/file_store_path_factory.py | 21 +- paimon-python/pypaimon/utils/path.py | 100 +++++++++ paimon-python/pypaimon/write/native_update.py | 1 - paimon-python/pypaimon/write/table_commit.py | 5 - paimon-python/pypaimon/write/table_update.py | 1 - paimon-python/pypaimon/write/write_builder.py | 6 - 16 files changed, 428 insertions(+), 114 deletions(-) create mode 100644 paimon-python/pypaimon/tests/native_data_directory_test.py create mode 100644 paimon-python/pypaimon/utils/path.py diff --git a/paimon-python/pypaimon/common/external_path_provider.py b/paimon-python/pypaimon/common/external_path_provider.py index e039b1d86ed2..21c87db7ac77 100644 --- a/paimon-python/pypaimon/common/external_path_provider.py +++ b/paimon-python/pypaimon/common/external_path_provider.py @@ -22,6 +22,8 @@ from abc import ABC, abstractmethod from typing import List +from pypaimon.utils.path import resolve_path + class ExternalPathProvider(ABC): """Provider for external data paths.""" @@ -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 resolve_path(resolve_path(external_base, self._relative_bucket_path), file_name) else: return f"{external_base.rstrip('/')}/{file_name}" @@ -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}" @@ -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 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')) @@ -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 resolve_path(resolve_path(selected_base, self._relative_bucket_path), file_name) else: return f"{selected_base.rstrip('/')}/{file_name}" diff --git a/paimon-python/pypaimon/manifest/schema/data_file_meta.py b/paimon-python/pypaimon/manifest/schema/data_file_meta.py index 6fd0a6a4f236..58b57ff6f26e 100644 --- a/paimon-python/pypaimon/manifest/schema/data_file_meta.py +++ b/paimon-python/pypaimon/manifest/schema/data_file_meta.py @@ -145,9 +145,8 @@ 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 + path_builder = resolve_path(table_path, 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/table_read.py b/paimon-python/pypaimon/read/table_read.py index 781d05e90c4e..3ec796c68ea2 100644 --- a/paimon-python/pypaimon/read/table_read.py +++ b/paimon-python/pypaimon/read/table_read.py @@ -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: diff --git a/paimon-python/pypaimon/read/table_scan.py b/paimon-python/pypaimon/read/table_scan.py index 4bddf3345be5..344e1a1703b5 100755 --- a/paimon-python/pypaimon/read/table_scan.py +++ b/paimon-python/pypaimon/read/table_scan.py @@ -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 diff --git a/paimon-python/pypaimon/table/file_store_table.py b/paimon-python/pypaimon/table/file_store_table.py index d0ec91f960b9..85c9c5775b4d 100644 --- a/paimon-python/pypaimon/table/file_store_table.py +++ b/paimon-python/pypaimon/table/file_store_table.py @@ -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 diff --git a/paimon-python/pypaimon/tests/data_file_path_directory_test.py b/paimon-python/pypaimon/tests/data_file_path_directory_test.py index c2fbccd925c5..784ec930a39a 100644 --- a/paimon-python/pypaimon/tests/data_file_path_directory_test.py +++ b/paimon-python/pypaimon/tests/data_file_path_directory_test.py @@ -21,6 +21,8 @@ import shutil import tempfile import unittest + +import pytest from unittest import mock import pyarrow as pa @@ -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): @@ -190,3 +148,77 @@ 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'), + ('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'), + ('file:/t', 'file:/data', 'file:/data'), +]) +def test_path_resolution_matches_java(parent, child, expected): + from pypaimon.utils.path import resolve_path + assert 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', [ + (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 diff --git a/paimon-python/pypaimon/tests/deletion_vector_path_test.py b/paimon-python/pypaimon/tests/deletion_vector_path_test.py index f6d4673004a1..bfa21004480f 100644 --- a/paimon-python/pypaimon/tests/deletion_vector_path_test.py +++ b/paimon-python/pypaimon/tests/deletion_vector_path_test.py @@ -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': diff --git a/paimon-python/pypaimon/tests/external_paths_test.py b/paimon-python/pypaimon/tests/external_paths_test.py index aab4fa38a599..c249ba4db849 100644 --- a/paimon-python/pypaimon/tests/external_paths_test.py +++ b/paimon-python/pypaimon/tests/external_paths_test.py @@ -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) diff --git a/paimon-python/pypaimon/tests/native_data_directory_test.py b/paimon-python/pypaimon/tests/native_data_directory_test.py new file mode 100644 index 000000000000..505494e65c14 --- /dev/null +++ b/paimon-python/pypaimon/tests/native_data_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. + +"""Java-compatible data directories across Python and Rust execution.""" + +from contextlib import ExitStack +from pathlib import Path +from unittest.mock import patch + +import pyarrow as pa +import pytest + +from pypaimon import CatalogFactory, Schema +from pypaimon.read.native_plan import native_plan +from pypaimon.schema.data_types import AtomicType, DataField +from pypaimon.write.native_write import NativeTableWrite +from pypaimon.write.table_upsert_by_key import TableUpsertByKey + + +pytestmark = pytest.mark.native_plan +_SCHEMA = pa.schema([('id', pa.int32()), ('pt', pa.string()), ('value', pa.int32())]) + + +def _table(tmp_path, mode='append', directory='relative', native=True, partitioned=True, extra=None): + catalog = CatalogFactory.create({'warehouse': str(tmp_path / 'warehouse')}) + catalog.create_database('db', True) + directories = {'relative': 'data/nested', 'normalized': 'data//discard/../nested/.', + 'absolute': str(tmp_path / 'relocated'), 'uri': (tmp_path / 'relocated').as_uri()} + options = {'data-file.path-directory': directories[directory], + 'write.native.enabled': str(native).lower()} + if mode == 'pk': + options.update({'bucket': '1', 'changelog-producer': 'input'}) + if mode == 'evolution': + options.update({'row-tracking.enabled': 'true', 'data-evolution.enabled': 'true', + 'target-file-row-num': '2', + 'deletion-vectors.enabled': 'true'}) + options.update(extra or {}) + partitions = ['pt'] if partitioned else [] + keys = ['id'] + partitions if mode == 'pk' else [] + catalog.create_table('db.t', Schema.from_pyarrow_schema( + _SCHEMA, primary_keys=keys, partition_keys=partitions, options=options), False) + return catalog.get_table('db.t') + + +def _write(table, data, native, streaming=False, identifier=1): + builder = table.new_stream_write_builder() if streaming else table.new_batch_write_builder() + writer, commit = builder.new_write(), builder.new_commit() + try: + if native: + assert isinstance(writer, NativeTableWrite) + writer.write_arrow(data) + messages = writer.prepare_commit(identifier) if streaming else writer.prepare_commit() + factory = table.path_factory() + for message in messages: + expected = factory.bucket_path(tuple(message.partition), message.bucket, canonical_partition=True) + for file in message.new_files + message.changelog_files: + assert file.file_path.startswith(expected + '/') + assert table.file_io.exists(file.file_path) + if streaming: + commit.commit(messages, identifier) + else: + commit.commit(messages) + if native: + assert writer._python_writer is None + finally: + writer.close() + commit.close() + + +def _read(table, native_scan, native_read, streaming=False, snapshot=None): + options = {'read.native.enabled': str(native_read).lower(), + 'scan.native-plan.enabled': str(native_scan).lower()} + if snapshot is not None: + options['scan.snapshot-id'] = str(snapshot) + table = table.copy(options) + builder = table.new_read_builder() + plan = native_plan(table) if native_scan else builder.new_scan().plan() + read = builder.new_read() + with ExitStack() as stack: + if native_read: + stack.enter_context(patch.object(read, '_create_split_read', + side_effect=AssertionError('Python read fallback'))) + if streaming: + reader = read.to_arrow_batch_reader(plan.splits()) + try: + result = reader.read_all() + finally: + reader.close() + else: + result = read.to_arrow(plan.splits()) + return result.sort_by('id').to_pylist() + + +@pytest.mark.parametrize('directory', ['relative', 'normalized', 'absolute', 'uri']) +@pytest.mark.parametrize('mode', ['append', 'pk', 'evolution']) +@pytest.mark.parametrize('native', [False, True]) +@pytest.mark.parametrize('partitioned', [False, True]) +def test_data_directory_interoperability(tmp_path, directory, mode, native, partitioned): + table = _table(tmp_path, mode, directory, native, partitioned) + data = pa.table({'id': [1, 2, 3], 'pt': ['a', 'a', 'b'], 'value': [10, 20, 30]}, schema=_SCHEMA) + _write(table, data, native) + _write(table, pa.table({'id': [4], 'pt': ['b'], 'value': [40]}, schema=_SCHEMA), native, + streaming=True, identifier=2) + expected = data.to_pylist() + [{'id': 4, 'pt': 'b', 'value': 40}] + for planner in (False, True): + for reader in (False, True): + assert _read(table, planner, reader, streaming=reader) == expected + assert (Path(table.table_path) / 'snapshot' / 'snapshot-2').is_file() + assert (Path(table.table_path) / 'schema' / 'schema-0').is_file() + assert not (Path(table.table_path) / 'bucket-0').exists() + + +@pytest.mark.parametrize('directory', ['relative', 'normalized', 'absolute', 'uri']) +@pytest.mark.parametrize('native', [False, True]) +@pytest.mark.parametrize('bucket_indexes', [False, True]) +def test_data_directory_updates_deletes_history_and_abort(tmp_path, directory, native, bucket_indexes): + table = _table(tmp_path, 'evolution', directory, native, extra={ + 'index-file-in-data-file-dir': str(bucket_indexes).lower()}) + _write(table, pa.table({'id': [1, 2, 3], 'pt': ['a'] * 3, 'value': [10, 20, 30]}, schema=_SCHEMA), native) + builder = table.new_batch_write_builder() + update = builder.new_update().with_update_type(['value']) + with ExitStack() as stack: + if native: + stack.enter_context(patch.object(TableUpsertByKey, '_upsert_partition', + side_effect=AssertionError('Python upsert fallback'))) + messages = update.upsert_by_arrow_with_key( + pa.table({'id': [1, 4], 'pt': ['a', 'b'], 'value': [11, 40]}, schema=_SCHEMA), ['id']) + commit = builder.new_commit() + try: + commit.commit(messages) + finally: + commit.close() + for row_id, expected in [(0, [20, 30, 40]), (1, [30, 40])]: + messages = builder.new_update().delete_by_row_id([row_id]) + commit = builder.new_commit() + try: + commit.commit(messages) + finally: + commit.close() + for planner in (False, True): + for reader in (False, True): + assert sorted(row['value'] for row in _read(table, planner, reader)) == expected + messages = builder.new_update().delete_by_row_id([2]) + staged = {entry.index_file.file_name for msg in messages for entry in msg.index_adds} + assert staged + before = {p for p in tmp_path.rglob('index-*') if p.name in staged} + assert before + commit = builder.new_commit() + try: + commit.abort(messages) + finally: + commit.close() + assert not any(path.exists() for path in before) + assert [row['value'] for row in _read(table, True, True, snapshot=1)] == [10, 20, 30] + assert [row['value'] for row in _read(table, True, True)] == [30, 40] + + +@pytest.mark.parametrize('directory', ['relative', 'absolute', 'uri']) +def test_data_directory_blob_native_read(tmp_path, directory): + table = _table(tmp_path, 'evolution', directory, False) + catalog = table.catalog_environment.catalog_loader.load() + catalog.drop_table('db.t') + fields = [DataField(0, 'id', AtomicType('INT')), DataField(1, 'payload', AtomicType('BLOB'))] + catalog.create_table('db.t', Schema(fields=fields, options=table.table_schema.options), False) + table = catalog.get_table('db.t') + builder = table.new_batch_write_builder() + writer, commit = builder.new_write(), builder.new_commit() + try: + writer.write_arrow(pa.table({'id': pa.array([1, 2, 3], pa.int32()), + 'payload': pa.array([b'hello', None, b''], pa.large_binary())})) + commit.commit(writer.prepare_commit()) + finally: + writer.close() + commit.close() + for planner in (False, True): + for reader in (False, True): + assert _read(table, planner, reader) == [ + {'id': 1, 'payload': b'hello'}, {'id': 2, 'payload': None}, {'id': 3, 'payload': b''}] diff --git a/paimon-python/pypaimon/tests/native_write_test.py b/paimon-python/pypaimon/tests/native_write_test.py index ba853082c01d..bfe8b638b483 100644 --- a/paimon-python/pypaimon/tests/native_write_test.py +++ b/paimon-python/pypaimon/tests/native_write_test.py @@ -130,12 +130,19 @@ def test_batch_native_write_commits_through_both_committers( @requires_native -def test_escaped_partition_file_path_and_abort(tmp_path, native_rest_catalog): +@pytest.mark.parametrize('directory', [None, 'relative', 'absolute', 'uri']) +def test_escaped_partition_file_path_and_abort(tmp_path, native_rest_catalog, directory): catalog = native_rest_catalog + options = {'file.format': 'parquet', 'write.native.enabled': 'true', + 'commit.native.enabled': 'true'} + if directory is not None: + options['data-file.path-directory'] = { + 'relative': 'data/nested', 'absolute': str(tmp_path / 'relocated'), + 'uri': (tmp_path / 'relocated').as_uri(), + }[directory] catalog.create_table('default.t', Schema.from_pyarrow_schema( pa.schema([('id', pa.int64()), ('pt', pa.string())]), - options={'file.format': 'parquet', 'write.native.enabled': 'true', - 'commit.native.enabled': 'true'}, partition_keys=['pt']), False) + options=options, partition_keys=['pt']), False) table = catalog.get_table('default.t') builder = table.new_batch_write_builder() writer = builder.new_write() diff --git a/paimon-python/pypaimon/utils/file_store_path_factory.py b/paimon-python/pypaimon/utils/file_store_path_factory.py index 35594c2963e1..67377a139652 100644 --- a/paimon-python/pypaimon/utils/file_store_path_factory.py +++ b/paimon-python/pypaimon/utils/file_store_path_factory.py @@ -22,6 +22,7 @@ from pypaimon.casting.row_to_string import cast_value_to_string, _format_timestamp, _is_unsupported from pypaimon.common.external_path_provider import ExternalPathProvider +from pypaimon.utils.path import resolve_path from pypaimon.schema.data_types import DataType from pypaimon.table.bucket_mode import BucketMode from pypaimon.table.row.generic_row import _is_ltz_type, _normalize_ltz, _parse_type_precision_scale @@ -47,7 +48,7 @@ def canonical_data_file_path(table, partition, bucket, file_name): bucket_path = table.path_factory().bucket_path( tuple(partition), bucket, canonical_partition=True) path = f"{bucket_path.rstrip('/')}/{file_name}" - root = table.table_path.rstrip('/') + root = resolve_path(table.path_factory().data_file_path(), '.').rstrip('/') if (root.startswith('file:') and path.startswith(root + '/') and '%' in path[len(root) + 1:]): from pypaimon.filesystem.local_file_io import LocalFileIO @@ -133,6 +134,8 @@ def __init__( self.changelog_file_prefix = changelog_file_prefix self.file_suffix_include_compression = file_suffix_include_compression self.file_compression = file_compression + if data_file_path_directory == '': + raise ValueError('data-file.path-directory must not be empty') self.data_file_path_directory = data_file_path_directory self.external_paths = external_paths or [] self.external_path_strategy = external_path_strategy @@ -159,9 +162,7 @@ def statistics_path(self) -> str: return f"{self._root}/{self.STATISTICS_PATH}" def data_file_path(self) -> str: - if self.data_file_path_directory: - return f"{self._root}/{self.data_file_path_directory}" - return self._root + return resolve_path(self._root, self.data_file_path_directory) def relative_bucket_path(self, partition: Tuple, bucket: int, canonical_partition: bool = False) -> str: if canonical_partition and partition: @@ -228,14 +229,18 @@ def _relative_bucket_path(self, partition: Tuple, bucket: int, canonical_partiti relative_parts = partition_parts + relative_parts # Add data file path directory if specified + relative = "/".join(relative_parts) + # Legacy Python timestamp partitions contain literal colons. Prefix a + # relative path so URI resolution cannot mistake the first field for a scheme. + if ':' in relative.split('/', 1)[0]: + relative = './' + relative if self.data_file_path_directory: - relative_parts = [self.data_file_path_directory] + relative_parts - - return "/".join(relative_parts) + return resolve_path(self.data_file_path_directory, relative) + return relative def bucket_path(self, partition: Tuple, bucket: int, canonical_partition: bool = False) -> str: relative_path = self.relative_bucket_path(partition, bucket, canonical_partition) - return f"{self._root}/{relative_path}" + return resolve_path(self._root, relative_path) def create_external_path_provider( self, partition: Tuple, bucket: int diff --git a/paimon-python/pypaimon/utils/path.py b/paimon-python/pypaimon/utils/path.py new file mode 100644 index 000000000000..4816bf98e2ac --- /dev/null +++ b/paimon-python/pypaimon/utils/path.py @@ -0,0 +1,100 @@ +# 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. + +"""Resolve storage paths using Paimon's Java Path semantics. + +Path text contains literal percent signs, spaces, query and fragment characters; +URL joining would incorrectly interpret or escape these characters. +""" + + +import os +import re + + +_WINDOWS = os.name == 'nt' + + +def _has_drive(path): + return bool(re.match(r'^/?[A-Za-z]:', path)) + + +def _normalize(path): + components = [] + for component in path.split('/'): + if component in ('', '.'): + continue + if component == '..' and components and components[-1] != '..': + components.pop() + else: + components.append(component) + normalized = ('/' if path.startswith('/') else '') + '/'.join(components) + if _WINDOWS and _has_drive(normalized) and len(normalized) == 3 and path != normalized: + normalized += '/' + return normalized + + +def _parts(path): + if _WINDOWS and _has_drive(path) and not path.startswith('/'): + path = '/' + path + colon = path.find(':') + scheme_end = colon + 1 if colon >= 0 and '/' not in path[:colon] else 0 + scheme, authority = path[:scheme_end], '' + if path[scheme_end:].startswith('//') and len(path) - scheme_end > 2: + slash = path.find('/', scheme_end + 2) + end = len(path) if slash < 0 else slash + if end > scheme_end + 2: + authority = path[scheme_end:end] + path = path[end:] + else: + path = path[scheme_end:] + path = re.sub(r'/+', '/', path) + if _WINDOWS and (_has_drive(path) or scheme in ('', 'file:')): + path = path.replace('\\', '/') + # Java URI recognizes a UNC authority after Windows separator conversion. + if not authority and path.startswith('//') and len(path) > 2: + end = path.find('/', 2) + end = len(path) if end < 0 else end + authority, path = path[:end], path[end:] + return scheme, authority, _normalize(path) + + +def resolve_path(parent, child): + """Resolve child against a directory, retaining its scheme and authority.""" + if child is None: + return parent.rstrip('/') or '/' + if child == '': + raise ValueError('Can not create a Path from an empty string') + parent_scheme, parent_authority, parent_path = _parts(parent) + child_scheme, child_authority, child_path = _parts(child) + if child_scheme: + scheme, authority, path = child_scheme, child_authority, child_path + elif child_authority: + scheme, authority, path = parent_scheme, child_authority, child_path + elif child_path.startswith('/'): + scheme, authority, path = parent_scheme, parent_authority, child_path + else: + scheme, authority = parent_scheme, parent_authority + path = (parent_path.rstrip('/') + '/' + child_path + if parent_path or scheme or authority else child_path) + path = _normalize(path) + if not scheme and not authority: + if _WINDOWS and _has_drive(path): + path = path.lstrip('/') + elif not path.startswith('/') and ':' in path.split('/')[0]: + path = './' + path + return scheme + authority + path diff --git a/paimon-python/pypaimon/write/native_update.py b/paimon-python/pypaimon/write/native_update.py index 8ecab848026e..3bd97ea32e09 100644 --- a/paimon-python/pypaimon/write/native_update.py +++ b/paimon-python/pypaimon/write/native_update.py @@ -37,7 +37,6 @@ def _native_row_id_table(table): or not table.options.native_write_enabled() or not table.options.data_evolution_enabled() or not table.options.row_tracking_enabled() - or table.options.data_file_path_directory() is not None or not _RowIdUpdateFileWriter.supports_table(table) or any(table.options.options.contains_key(key) for key in SCAN_KEYS) or not native_write_available()): diff --git a/paimon-python/pypaimon/write/table_commit.py b/paimon-python/pypaimon/write/table_commit.py index 41472bf3585a..ef2595d93c40 100644 --- a/paimon-python/pypaimon/write/table_commit.py +++ b/paimon-python/pypaimon/write/table_commit.py @@ -117,11 +117,6 @@ def _prepare_native_commit(self, messages): 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/table_update.py b/paimon-python/pypaimon/write/table_update.py index 81e3717ed1ed..38589a92158e 100644 --- a/paimon-python/pypaimon/write/table_update.py +++ b/paimon-python/pypaimon/write/table_update.py @@ -683,7 +683,6 @@ def _matched_delete_row_ids( splits = scan.plan_for_write().splits() if (splits and self.table.options.native_write_enabled() - and self.table.options.data_file_path_directory() is None and not any(isinstance(split, QueryAuthSplit) for split in splits)): try: diff --git a/paimon-python/pypaimon/write/write_builder.py b/paimon-python/pypaimon/write/write_builder.py index bfe76659f9a4..bd518d826da7 100644 --- a/paimon-python/pypaimon/write/write_builder.py +++ b/paimon-python/pypaimon/write/write_builder.py @@ -66,12 +66,6 @@ def _native_write(self, static_partition=None, stream=False): check_sequence_field_supported(self.table) 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, From 03420e9c8b78e7e6d8e63e43e81b6dd0f8252ec6 Mon Sep 17 00:00:00 2001 From: leaves12138 <41894543+leaves12138@users.noreply.github.com> Date: Mon, 28 Sep 2026 08:53:36 +0800 Subject: [PATCH 2/2] fix(python): preserve literal data-directory paths across native IO --- .../pypaimon/common/external_path_provider.py | 8 ++-- .../manifest/schema/data_file_meta.py | 4 +- .../pypaimon/read/split_serializer.py | 6 ++- .../tests/data_file_path_directory_test.py | 36 +++++++++++++++-- .../tests/native_data_directory_test.py | 40 +++++++++++++++++-- .../tests/table_upsert_by_key_test.py | 2 + .../pypaimon/utils/file_store_path_factory.py | 8 ++-- paimon-python/pypaimon/utils/path.py | 21 +++++++++- 8 files changed, 107 insertions(+), 18 deletions(-) diff --git a/paimon-python/pypaimon/common/external_path_provider.py b/paimon-python/pypaimon/common/external_path_provider.py index 21c87db7ac77..b32538d34586 100644 --- a/paimon-python/pypaimon/common/external_path_provider.py +++ b/paimon-python/pypaimon/common/external_path_provider.py @@ -22,7 +22,7 @@ from abc import ABC, abstractmethod from typing import List -from pypaimon.utils.path import resolve_path +from pypaimon.utils.path import resolve_path, to_file_io_path class ExternalPathProvider(ABC): @@ -73,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 resolve_path(resolve_path(external_base, self._relative_bucket_path), 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}" @@ -108,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 resolve_path(external_base, 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')) @@ -153,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 resolve_path(resolve_path(selected_base, self._relative_bucket_path), 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}" diff --git a/paimon-python/pypaimon/manifest/schema/data_file_meta.py b/paimon-python/pypaimon/manifest/schema/data_file_meta.py index 58b57ff6f26e..4cefe7048ccb 100644 --- a/paimon-python/pypaimon/manifest/schema/data_file_meta.py +++ b/paimon-python/pypaimon/manifest/schema/data_file_meta.py @@ -145,8 +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): - from pypaimon.utils.path import resolve_path + 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) diff --git a/paimon-python/pypaimon/read/split_serializer.py b/paimon-python/pypaimon/read/split_serializer.py index c329a34e2f44..d7360df7c055 100644 --- a/paimon-python/pypaimon/read/split_serializer.py +++ b/paimon-python/pypaimon/read/split_serializer.py @@ -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. @@ -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 @@ -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, diff --git a/paimon-python/pypaimon/tests/data_file_path_directory_test.py b/paimon-python/pypaimon/tests/data_file_path_directory_test.py index 784ec930a39a..3920a3c34acc 100644 --- a/paimon-python/pypaimon/tests/data_file_path_directory_test.py +++ b/paimon-python/pypaimon/tests/data_file_path_directory_test.py @@ -168,6 +168,8 @@ def test_default_keeps_data_files_at_bucket_root(self): ('/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'), @@ -178,11 +180,15 @@ def test_default_keeps_data_files_at_bucket_root(self): ('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'), ]) -def test_path_resolution_matches_java(parent, child, expected): - from pypaimon.utils.path import resolve_path - assert resolve_path(parent, child) == expected +@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(): @@ -208,6 +214,12 @@ def test_data_directory_cannot_change_on_table_copy(tmp_path, stored, 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'), @@ -222,3 +234,21 @@ 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 diff --git a/paimon-python/pypaimon/tests/native_data_directory_test.py b/paimon-python/pypaimon/tests/native_data_directory_test.py index 505494e65c14..f24fb0dbd5ae 100644 --- a/paimon-python/pypaimon/tests/native_data_directory_test.py +++ b/paimon-python/pypaimon/tests/native_data_directory_test.py @@ -18,6 +18,7 @@ """Java-compatible data directories across Python and Rust execution.""" from contextlib import ExitStack +import os from pathlib import Path from unittest.mock import patch @@ -39,7 +40,8 @@ def _table(tmp_path, mode='append', directory='relative', native=True, partition catalog = CatalogFactory.create({'warehouse': str(tmp_path / 'warehouse')}) catalog.create_database('db', True) directories = {'relative': 'data/nested', 'normalized': 'data//discard/../nested/.', - 'absolute': str(tmp_path / 'relocated'), 'uri': (tmp_path / 'relocated').as_uri()} + 'absolute': str(tmp_path / 'relocated'), 'uri': (tmp_path / 'relocated').as_uri(), + 'literal_uri': 'file:' + str(tmp_path / 'data%2Fwith space?#fragment')} options = {'data-file.path-directory': directories[directory], 'write.native.enabled': str(native).lower()} if mode == 'pk': @@ -105,7 +107,9 @@ def _read(table, native_scan, native_read, streaming=False, snapshot=None): return result.sort_by('id').to_pylist() -@pytest.mark.parametrize('directory', ['relative', 'normalized', 'absolute', 'uri']) +@pytest.mark.parametrize('directory', [ + 'relative', 'normalized', 'absolute', 'uri', + pytest.param('literal_uri', marks=pytest.mark.skipif(os.name == 'nt', reason='POSIX file names'))]) @pytest.mark.parametrize('mode', ['append', 'pk', 'evolution']) @pytest.mark.parametrize('native', [False, True]) @pytest.mark.parametrize('partitioned', [False, True]) @@ -115,6 +119,9 @@ def test_data_directory_interoperability(tmp_path, directory, mode, native, part _write(table, data, native) _write(table, pa.table({'id': [4], 'pt': ['b'], 'value': [40]}, schema=_SCHEMA), native, streaming=True, identifier=2) + if directory == 'literal_uri': + assert (tmp_path / 'data%2Fwith space?#fragment').is_dir() + assert not (tmp_path / 'data').exists() expected = data.to_pylist() + [{'id': 4, 'pt': 'b', 'value': 40}] for planner in (False, True): for reader in (False, True): @@ -124,7 +131,34 @@ def test_data_directory_interoperability(tmp_path, directory, mode, native, part assert not (Path(table.table_path) / 'bucket-0').exists() -@pytest.mark.parametrize('directory', ['relative', 'normalized', 'absolute', 'uri']) +@pytest.mark.skipif(os.name == 'nt', reason='POSIX file names') +@pytest.mark.parametrize('strategy', ['round-robin', 'entropy-inject', 'weight-robin']) +def test_literal_data_directory_with_external_paths(tmp_path, strategy): + table = _table(tmp_path, directory='literal_uri', native=False, extra={ + 'data-file.external-paths': ','.join((tmp_path / name).as_uri() for name in ['first', 'second']), + 'data-file.external-paths.strategy': strategy, + 'data-file.external-paths.weights': '1,2', + }) + data = pa.table({'id': [1], 'pt': ['a'], 'value': [10]}, schema=_SCHEMA) + builder = table.new_batch_write_builder() + writer, commit = builder.new_write(), builder.new_commit() + try: + writer.write_arrow(data) + messages = writer.prepare_commit() + assert all(file.external_path for message in messages for file in message.new_files) + commit.commit(messages) + finally: + writer.close() + commit.close() + assert (tmp_path / 'data%2Fwith space?#fragment').is_dir() + for planner in (False, True): + for reader in (False, True): + assert _read(table, planner, reader) == data.to_pylist() + + +@pytest.mark.parametrize('directory', [ + 'relative', 'normalized', 'absolute', 'uri', + pytest.param('literal_uri', marks=pytest.mark.skipif(os.name == 'nt', reason='POSIX file names'))]) @pytest.mark.parametrize('native', [False, True]) @pytest.mark.parametrize('bucket_indexes', [False, True]) def test_data_directory_updates_deletes_history_and_abort(tmp_path, directory, native, bucket_indexes): diff --git a/paimon-python/pypaimon/tests/table_upsert_by_key_test.py b/paimon-python/pypaimon/tests/table_upsert_by_key_test.py index 68d26e86e8db..8833b1fc8855 100644 --- a/paimon-python/pypaimon/tests/table_upsert_by_key_test.py +++ b/paimon-python/pypaimon/tests/table_upsert_by_key_test.py @@ -22,6 +22,7 @@ import pyarrow as pa import pyarrow.parquet as pq +import pytest from pypaimon.read.table_read import TableRead from pypaimon.table.special_fields import SpecialFields @@ -140,6 +141,7 @@ def test_invalid_row_group_size_fails_before_opening_output(self): _RowIdUpdateFileWriter(table, (), ['id']) output_stream.assert_not_called() + @pytest.mark.python_write @mock.patch.object(_RowIdUpdateFileWriter, '_ROW_GROUP_MAX_ROWS', 2) def test_partial_upsert_streams_original_file_group(self): schema = pa.schema([ diff --git a/paimon-python/pypaimon/utils/file_store_path_factory.py b/paimon-python/pypaimon/utils/file_store_path_factory.py index 67377a139652..147bf3d70ccb 100644 --- a/paimon-python/pypaimon/utils/file_store_path_factory.py +++ b/paimon-python/pypaimon/utils/file_store_path_factory.py @@ -22,7 +22,7 @@ from pypaimon.casting.row_to_string import cast_value_to_string, _format_timestamp, _is_unsupported from pypaimon.common.external_path_provider import ExternalPathProvider -from pypaimon.utils.path import resolve_path +from pypaimon.utils.path import resolve_path, to_file_io_path from pypaimon.schema.data_types import DataType from pypaimon.table.bucket_mode import BucketMode from pypaimon.table.row.generic_row import _is_ltz_type, _normalize_ltz, _parse_type_precision_scale @@ -162,7 +162,8 @@ def statistics_path(self) -> str: return f"{self._root}/{self.STATISTICS_PATH}" def data_file_path(self) -> str: - return resolve_path(self._root, self.data_file_path_directory) + path = resolve_path(self._root, self.data_file_path_directory) + return to_file_io_path(path) if self.data_file_path_directory is not None else path def relative_bucket_path(self, partition: Tuple, bucket: int, canonical_partition: bool = False) -> str: if canonical_partition and partition: @@ -240,7 +241,8 @@ def _relative_bucket_path(self, partition: Tuple, bucket: int, canonical_partiti def bucket_path(self, partition: Tuple, bucket: int, canonical_partition: bool = False) -> str: relative_path = self.relative_bucket_path(partition, bucket, canonical_partition) - return resolve_path(self._root, relative_path) + path = resolve_path(self._root, relative_path) + return to_file_io_path(path) if self.data_file_path_directory is not None else path def create_external_path_provider( self, partition: Tuple, bucket: int diff --git a/paimon-python/pypaimon/utils/path.py b/paimon-python/pypaimon/utils/path.py index 4816bf98e2ac..6987353bb093 100644 --- a/paimon-python/pypaimon/utils/path.py +++ b/paimon-python/pypaimon/utils/path.py @@ -43,8 +43,11 @@ def _normalize(path): else: components.append(component) normalized = ('/' if path.startswith('/') else '') + '/'.join(components) - if _WINDOWS and _has_drive(normalized) and len(normalized) == 3 and path != normalized: + if (_WINDOWS and normalized.startswith('/') and _has_drive(normalized) + and len(normalized) == 3 and path != normalized): normalized += '/' + elif not normalized.startswith('/') and ':' in normalized.split('/')[0]: + normalized = './' + normalized return normalized @@ -90,7 +93,7 @@ def resolve_path(parent, child): else: scheme, authority = parent_scheme, parent_authority path = (parent_path.rstrip('/') + '/' + child_path - if parent_path or scheme or authority else child_path) + if parent_path or (child_path and (scheme or authority)) else child_path) path = _normalize(path) if not scheme and not authority: if _WINDOWS and _has_drive(path): @@ -98,3 +101,17 @@ def resolve_path(parent, child): elif not path.startswith('/') and ':' in path.split('/')[0]: path = './' + path return scheme + authority + path + + +def to_file_io_path(path): + """Keep literal Java file-path characters from being decoded as a URL.""" + if not path.startswith('file:') or not any(char in path for char in '%?#'): + return path + _, authority, local_path = _parts(path) + if authority and authority != '//localhost': + if not (_WINDOWS and authority.endswith(':')): + return authority + local_path + local_path = authority[1:] + local_path + if _WINDOWS and _has_drive(local_path): + local_path = local_path.lstrip('/') + return local_path