From 3600adbaae635e49ad325a6e93f6d2104c833cdf Mon Sep 17 00:00:00 2001 From: JingsongLi Date: Mon, 28 Sep 2026 20:11:56 +0800 Subject: [PATCH 1/3] [python] Enable native writes with external data paths --- paimon-python/README.md | 11 +- .../pypaimon/common/options/core_options.py | 6 +- .../manifest/schema/data_file_meta.py | 22 ++ paimon-python/pypaimon/read/split_read.py | 11 +- .../pypaimon/read/split_serializer.py | 2 +- .../pypaimon/tests/external_paths_test.py | 8 + .../tests/native_external_write_test.py | 226 ++++++++++++++++++ .../pypaimon/utils/file_store_path_factory.py | 2 +- .../write/commit/row_id_conflict_rewriter.py | 4 +- .../pypaimon/write/file_store_commit.py | 13 +- paimon-python/pypaimon/write/native_commit.py | 2 +- paimon-python/pypaimon/write/native_write.py | 1 - .../pypaimon/write/writer/data_writer.py | 17 +- 13 files changed, 284 insertions(+), 41 deletions(-) create mode 100644 paimon-python/pypaimon/tests/native_external_write_test.py diff --git a/paimon-python/README.md b/paimon-python/README.md index f1654fa0e796..a32becf69596 100644 --- a/paimon-python/README.md +++ b/paimon-python/README.md @@ -161,13 +161,18 @@ finally: The native writer returns ordinary PyPaimon commit messages, so the Python committer also works when `commit.native.enabled` is false. Batch overwrite and reusable stream writers retain the builder's commit user and identifier. Native -write is currently limited to Parquet tables without BLOB fields or -data-evolution mode, on the same filesystem/JDBC publication route as native -commit. Writer methods requiring Python's specialized path select the Python +write supports Parquet append, primary-key and data-evolution tables without +BLOB fields or optional data-evolution row sidecars, on the same filesystem/JDBC +publication route as native commit. Writer methods requiring Python's specialized path select the Python writer before native data is written. If the runtime or table route is unavailable, write uses Python. Once Rust starts writing a batch, errors propagate without retrying that batch through Python. +Native writes honor `data-file.path-directory` and the configured +`data-file.external-paths` strategy. Existing files keep their recorded locations +when the write destinations change. Python and native readers and committers +can exchange these files, including external data files and their index sidecars. + Both native options are disabled by default. # Native commit diff --git a/paimon-python/pypaimon/common/options/core_options.py b/paimon-python/pypaimon/common/options/core_options.py index a1fb6b2abfbf..0a97a6d2a7f8 100644 --- a/paimon-python/pypaimon/common/options/core_options.py +++ b/paimon-python/pypaimon/common/options/core_options.py @@ -1692,13 +1692,13 @@ def data_file_external_paths_weights(self, default=None): ) if value is None: return None - parts = value.split(",") + parts = value.rstrip(",").split(",") weights = [] for part in parts: parsed = int(part.strip()) - if parsed <= 0: + if parsed <= 0 or parsed > 2147483647: raise ValueError( - f"Weight must be positive, got: {parsed}" + f"Weight must be a positive 32-bit integer, got: {parsed}" ) weights.append(parsed) return weights diff --git a/paimon-python/pypaimon/manifest/schema/data_file_meta.py b/paimon-python/pypaimon/manifest/schema/data_file_meta.py index 4cefe7048ccb..8e454bc84a15 100644 --- a/paimon-python/pypaimon/manifest/schema/data_file_meta.py +++ b/paimon-python/pypaimon/manifest/schema/data_file_meta.py @@ -141,6 +141,28 @@ def create( file_path=file_path, ) + def physical_path(self): + """Resolve persisted Java path text for Python file I/O without URL decoding.""" + from pypaimon.utils.path import to_file_io_path + path = self.external_path or self.file_path + return to_file_io_path(str(path)) if path is not None else None + + def aligned_file_path(self, file_name, bucket_path=None): + """Place a sidecar beside its data file, including external locations.""" + from pypaimon.utils.path import resolve_path, to_file_io_path + data_path = self.physical_path() + parent = data_path.rsplit('/', 1)[0] if data_path and '/' in data_path else bucket_path + return to_file_io_path(resolve_path(parent, file_name) if parent else file_name) + + def collect_files(self, bucket_path=None): + """Return the physical data file and its aligned sidecars for cleanup.""" + from pypaimon.utils.path import resolve_path, to_file_io_path + path = self.physical_path() + if not path and bucket_path: + path = to_file_io_path(resolve_path(bucket_path, self.file_name)) + return ([path] if path else []) + [ + self.aligned_file_path(name, bucket_path) for name in self.extra_files] + def set_file_path( self, table_path: str, partition: GenericRow, bucket: int, default_part_value: str = "__DEFAULT_PARTITION__", diff --git a/paimon-python/pypaimon/read/split_read.py b/paimon-python/pypaimon/read/split_read.py index 2a044c4045b9..a26b6e6803f4 100644 --- a/paimon-python/pypaimon/read/split_read.py +++ b/paimon-python/pypaimon/read/split_read.py @@ -287,7 +287,7 @@ def file_reader_supplier(self, file: DataFileMeta, for_merge_read: bool, read_paimon_predicate = None # Use external_path if available, otherwise use file_path - file_path = file.external_path if file.external_path else file.file_path + file_path = file.physical_path() file_format = format_identifier(os.path.basename(file_path)) batch_size = self.table.options.read_batch_size() @@ -607,12 +607,7 @@ def _is_row_sidecar_file(file_name: str) -> bool: @staticmethod def _aligned_extra_file_path(file: DataFileMeta, extra_file: str) -> str: - if "://" in extra_file or extra_file.startswith("/"): - return extra_file - file_path = file.external_path if file.external_path else file.file_path - if not file_path or "/" not in file_path: - return extra_file - return f"{file_path.rsplit('/', 1)[0]}/{extra_file}" + return file.aligned_file_path(extra_file) def _get_fields_and_predicate(self, schema_id: int, read_fields): key = (schema_id, tuple(read_fields)) @@ -1751,7 +1746,7 @@ def _create_raw_blob_file_reader( if not row_indices: return None - file_path = file.external_path if file.external_path else file.file_path + file_path = file.physical_path() blob_parallelism = self._blob_parallelism return FormatBlobReader( self.table.file_io, diff --git a/paimon-python/pypaimon/read/split_serializer.py b/paimon-python/pypaimon/read/split_serializer.py index d7360df7c055..0dbf0484b261 100644 --- a/paimon-python/pypaimon/read/split_serializer.py +++ b/paimon-python/pypaimon/read/split_serializer.py @@ -655,7 +655,7 @@ 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 to_file_io_path( + meta.file_path = meta.physical_path() if external_path else to_file_io_path( "%s/%s" % (bucket_path.rstrip('/'), file_name)) return meta diff --git a/paimon-python/pypaimon/tests/external_paths_test.py b/paimon-python/pypaimon/tests/external_paths_test.py index c249ba4db849..a77d3506a9b4 100644 --- a/paimon-python/pypaimon/tests/external_paths_test.py +++ b/paimon-python/pypaimon/tests/external_paths_test.py @@ -265,6 +265,14 @@ def test_mismatched_lengths_raises(self): class WeightsParsingTest(unittest.TestCase): """Test CoreOptions.data_file_external_paths_weights() parsing and validation.""" + def test_weights_use_java_integer_range_and_trailing_comma_rules(self): + self.assertEqual(CoreOptions.from_dict({ + 'data-file.external-paths.weights': '1,2147483647,,' + }).data_file_external_paths_weights(), [1, 2147483647]) + for value in ['2147483648', '1,9999999999', '', '1,,2']: + with self.subTest(value=value), self.assertRaises(ValueError): + CoreOptions.from_dict({'data-file.external-paths.weights': value}).data_file_external_paths_weights() + def test_valid_weights(self): """Normal comma-separated positive integers.""" from pypaimon.common.options.core_options import CoreOptions diff --git a/paimon-python/pypaimon/tests/native_external_write_test.py b/paimon-python/pypaimon/tests/native_external_write_test.py new file mode 100644 index 000000000000..f1bed9ddd59b --- /dev/null +++ b/paimon-python/pypaimon/tests/native_external_write_test.py @@ -0,0 +1,226 @@ +# 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. + +"""Native external data files interoperate with Python readers and committers.""" + +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.write.native_write import NativeTableWrite + +pytestmark = pytest.mark.native_plan +_SCHEMA = pa.schema([('id', pa.int32()), ('pt', pa.string()), ('value', pa.int32())]) + + +def _table(tmp_path, mode, strategy, native, partitioned=True, extra=None): + catalog = CatalogFactory.create({'warehouse': str(tmp_path / 'warehouse')}) + catalog.create_database('db', True) + paths = [(tmp_path / name).as_uri() for name in ['external-a', 'external-b']] + if strategy == 'specific-fs': + paths.insert(0, 's3://unused-bucket/data') + options = {'write.native.enabled': str(native).lower(), + 'data-file.external-paths': ','.join(paths), + 'data-file.external-paths.strategy': strategy, + 'data-file.external-paths.weights': '1,2', + 'data-file.external-paths.specific-fs': 'FILE', + 'data-file.path-directory': 'data/nested'} + partitions = ['pt'] if partitioned else [] + keys = [] + if mode == 'pk': + keys = ['id'] + partitions + options.update({'bucket': '1', 'changelog-producer': 'input'}) + elif mode == 'evolution': + options.update({'row-tracking.enabled': 'true', 'data-evolution.enabled': 'true', + 'deletion-vectors.enabled': 'true', 'target-file-row-num': '1', + 'index-file-in-data-file-dir': 'true'}) + options.update(extra or {}) + 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: + assert isinstance(writer, NativeTableWrite) == native + writer.write_arrow(data) + messages = writer.prepare_commit(identifier) if streaming else writer.prepare_commit() + for message in messages: + for file in message.new_files + message.changelog_files: + assert bool(file.external_path) == (table.options.data_file_external_paths_strategy() != 'none') + assert table.file_io.exists(file.file_path), file.file_path + if file.external_path: + assert '/external-' in file.external_path + assert not Path(table.table_path, file.file_name).exists() + if streaming: + commit.commit(messages, identifier) + else: + commit.commit(messages) + if native: + assert writer._python_writer is None + return messages + finally: + writer.close() + commit.close() + + +def _read(table, planner, reader, snapshot=None): + options = {'scan.native-plan.enabled': str(planner).lower(), + 'read.native.enabled': str(reader).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 planner else builder.new_scan().plan() + read = builder.new_read() + with ExitStack() as stack: + if reader: + stack.enter_context(patch.object(read, '_create_split_read', + side_effect=AssertionError('Python read fallback'))) + result = read.to_arrow(plan.splits()) + return sorted(result.to_pylist(), key=lambda row: row['id']) + + +@pytest.mark.parametrize('mode', ['append', 'pk', 'evolution']) +@pytest.mark.parametrize('strategy', ['none', 'round-robin', 'weight-robin', 'entropy-inject', 'specific-fs']) +@pytest.mark.parametrize('native', [False, True]) +@pytest.mark.parametrize('partitioned', [False, True]) +def test_external_batch_and_stream_interoperability(tmp_path, mode, strategy, native, partitioned): + table = _table(tmp_path, mode, strategy, native, partitioned) + first = pa.table({'id': [1, 2], 'pt': ['a', 'b'], 'value': [10, 20]}, schema=_SCHEMA) + second = pa.table({'id': [3], 'pt': ['a'], 'value': [30]}, schema=_SCHEMA) + _write(table, first, native) + _write(table, second, native, streaming=True, identifier=2) + # Existing files must use their recorded paths after changing write destinations. + table = table.copy({'data-file.external-paths': (tmp_path / 'new-location').as_uri()}) + for planner in (False, True): + for reader in (False, True): + assert _read(table, planner, reader) == first.to_pylist() + second.to_pylist() + + +@pytest.mark.parametrize('strategy', ['round-robin', 'entropy-inject']) +def test_native_external_escaped_partition_and_abort(tmp_path, strategy): + table = _table(tmp_path, 'append', strategy, True) + data = pa.table({'id': [1], 'pt': ['a/b%?#'], 'value': [10]}, schema=_SCHEMA) + _write(table, data, True) + for planner in (False, True): + for reader in (False, True): + assert _read(table, planner, reader) == data.to_pylist() + builder = table.new_batch_write_builder() + writer, commit = builder.new_write(), builder.new_commit() + try: + writer.write_arrow(data) + messages = writer.prepare_commit() + files = [file for message in messages for file in message.new_files] + assert files and all(table.file_io.exists(file.file_path) for file in files) + commit.abort(messages) + assert not any(table.file_io.exists(file.file_path) for file in files) + finally: + writer.close() + commit.close() + assert _read(table, True, True) == data.to_pylist() + + +@pytest.mark.parametrize('strategy', ['round-robin', 'entropy-inject']) +@pytest.mark.parametrize('native', [False, True]) +def test_external_updates_upserts_deletes_history_and_abort(tmp_path, strategy, native): + table = _table(tmp_path, 'evolution', strategy, native) + original = pa.table({'id': [1, 2, 3], 'pt': ['a'] * 3, 'value': [10, 20, 30]}, schema=_SCHEMA) + _write(table, original, native) + builder = table.new_batch_write_builder() + update = builder.new_update().with_update_type(['value']) + 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 [row['value'] for row in _read(table, planner, reader)] == expected + messages = builder.new_update().delete_by_row_id([2]) + staged = [table.path_factory().bucket_index_path( + tuple(entry.partition.values), entry.bucket, entry.index_file) + for message in messages for entry in message.index_adds] + assert staged and all(table.file_io.exists(path) for path in staged) + commit = builder.new_commit() + try: + commit.abort(messages) + finally: + commit.close() + assert not any(table.file_io.exists(path) for path in staged) + assert _read(table, True, True, snapshot=1) == original.to_pylist() + assert [row['value'] for row in _read(table, True, True)] == [30, 40] + + +@pytest.mark.parametrize('strategy', ['none', 'round-robin', 'entropy-inject']) +@pytest.mark.parametrize('native', [False, True]) +def test_external_blob_fallback_and_readers(tmp_path, strategy, native): + from pypaimon.schema.data_types import AtomicType, DataField + table = _table(tmp_path, 'evolution', strategy, native, partitioned=False) + catalog = table.catalog_environment.catalog_loader.load() + catalog.drop_table('db.t') + fields = [DataField(0, 'id', AtomicType('INT')), + DataField(1, 'first', AtomicType('BLOB')), DataField(2, 'second', AtomicType('BLOB'))] + catalog.create_table('db.t', Schema(fields=fields, options=table.table_schema.options), False) + table = catalog.get_table('db.t') + data = pa.table({'id': pa.array([1, 2, 3], pa.int32()), + 'first': pa.array([b'hello', None, b''], pa.large_binary()), + 'second': pa.array([None, b'world', b'!'], pa.large_binary())}) + # Blob row streams still require the Python writer, including when opted in. + _write(table, data, False) + for planner in (False, True): + for reader in (False, True): + assert _read(table, planner, reader) == data.to_pylist() + + +@pytest.mark.parametrize('mode', ['append', 'pk', 'evolution']) +def test_python_abort_removes_native_external_sidecars(tmp_path, mode): + table = _table(tmp_path, mode, 'entropy-inject', True, extra={ + 'commit.native.enabled': 'false', 'file-index.bloom-filter.columns': 'value', + 'file-index.bloom-filter.value.items': '10', 'file-index.in-manifest-threshold': '0 B'}) + builder = table.new_batch_write_builder() + writer, commit = builder.new_write(), builder.new_commit() + try: + assert isinstance(writer, NativeTableWrite) + writer.write_arrow(pa.table({'id': [1], 'pt': ['a/b%?#'], 'value': [10]}, schema=_SCHEMA)) + messages = writer.prepare_commit() + files = [file for message in messages for file in message.new_files] + assert files and all(file.extra_files for file in files) + paths = [path for message in messages for file in message.new_files + message.changelog_files + for path in file.collect_files()] + assert all(table.file_io.exists(path) for path in paths) + commit.abort(messages) + assert not any(table.file_io.exists(path) for path in paths) + finally: + writer.close() + commit.close() diff --git a/paimon-python/pypaimon/utils/file_store_path_factory.py b/paimon-python/pypaimon/utils/file_store_path_factory.py index 147bf3d70ccb..15cfcceb221b 100644 --- a/paimon-python/pypaimon/utils/file_store_path_factory.py +++ b/paimon-python/pypaimon/utils/file_store_path_factory.py @@ -288,7 +288,7 @@ def _partition_path_requires_explicit_location(self, partition: Tuple) -> bool: def bucket_index_path(self, partition: Tuple, bucket: int, index_file, file_io=None) -> str: """Resolve an existing bucket index, including the legacy Python DV layout.""" if index_file.external_path: - return index_file.external_path + return to_file_io_path(index_file.external_path) legacy_path = f"{self.index_path()}/{index_file.file_name}" if not self.index_file_in_data_file_dir: return legacy_path diff --git a/paimon-python/pypaimon/write/commit/row_id_conflict_rewriter.py b/paimon-python/pypaimon/write/commit/row_id_conflict_rewriter.py index 898fe909dfb1..3fdaf916c612 100644 --- a/paimon-python/pypaimon/write/commit/row_id_conflict_rewriter.py +++ b/paimon-python/pypaimon/write/commit/row_id_conflict_rewriter.py @@ -202,7 +202,7 @@ def _affected_current_files(self, base_entries, candidates): ): key = ( base_key, - base.file.external_path or base.file.file_path + base.file.physical_path() or base.file.file_name, ) affected[key] = base.file @@ -274,7 +274,7 @@ def _to_manifest_entries(self, messages: List[CommitMessage]): def _abort(self, messages): for message in messages: for file in message.new_files: - path = file.external_path or file.file_path + path = file.physical_path() if path: self.table.file_io.delete_quietly(path) diff --git a/paimon-python/pypaimon/write/file_store_commit.py b/paimon-python/pypaimon/write/file_store_commit.py index 5c0b6f742749..1b70cf8a5318 100644 --- a/paimon-python/pypaimon/write/file_store_commit.py +++ b/paimon-python/pypaimon/write/file_store_commit.py @@ -93,13 +93,10 @@ def _abort_commit_messages(table, commit_messages: List[CommitMessage]): + list(message.compact_changelog_files)): path = None try: - path = file.external_path or file.file_path - if not path: - bucket_path = table.path_factory().bucket_path( - tuple(message.partition), message.bucket) - path = '%s/%s' % (bucket_path.rstrip('/'), file.file_name) - if path: - table.file_io.delete_quietly(str(path)) + bucket_path = None if file.physical_path() else table.path_factory().bucket_path( + tuple(message.partition), message.bucket) + for path in file.collect_files(bucket_path): + table.file_io.delete_quietly(path) except Exception as error: logger.warning( "Failed to clean up file %s during abort: %s", @@ -897,7 +894,7 @@ def _is_duplicate_commit( path_factory = self.table.path_factory() for entry in entries: file = entry.file - file.file_path = file.external_path or "%s/%s" % ( + file.file_path = file.physical_path() if file.external_path else "%s/%s" % ( path_factory.bucket_path( tuple(entry.partition.values), entry.bucket, diff --git a/paimon-python/pypaimon/write/native_commit.py b/paimon-python/pypaimon/write/native_commit.py index 39b86cabe030..46de93bcb6c9 100644 --- a/paimon-python/pypaimon/write/native_commit.py +++ b/paimon-python/pypaimon/write/native_commit.py @@ -194,6 +194,6 @@ def from_native_commit_messages(table, messages): table.trimmed_primary_keys_fields) for message in messages] for message in decoded: for file in message.new_files + message.changelog_files: - file.file_path = file.external_path or canonical_data_file_path( + file.file_path = file.physical_path() if file.external_path else canonical_data_file_path( table, message.partition, message.bucket, file.file_name) return decoded diff --git a/paimon-python/pypaimon/write/native_write.py b/paimon-python/pypaimon/write/native_write.py index 227a1b905c3d..f01c63688761 100644 --- a/paimon-python/pypaimon/write/native_write.py +++ b/paimon-python/pypaimon/write/native_write.py @@ -70,7 +70,6 @@ def create_native_write(table, commit_user, static_partition=None, stream=False) # Rust does not produce the optional random-access .row sidecars. or (table.options.data_evolution_enabled() and table.options.data_evolution_row_sidecar_enabled()) - or table.options.data_file_external_paths() or table.bucket_mode() not in (BucketMode.HASH_FIXED, BucketMode.BUCKET_UNAWARE) or (table.options.deletion_vectors_enabled() diff --git a/paimon-python/pypaimon/write/writer/data_writer.py b/paimon-python/pypaimon/write/writer/data_writer.py index ad336b0b39f7..731727fbc25c 100644 --- a/paimon-python/pypaimon/write/writer/data_writer.py +++ b/paimon-python/pypaimon/write/writer/data_writer.py @@ -198,16 +198,12 @@ def abort(self): def _delete_committed_files(self, file_metas: List[DataFileMeta]): for file_meta in file_metas: try: - path_to_delete = file_meta.external_path if file_meta.external_path else file_meta.file_path - if path_to_delete: - path_str = str(path_to_delete) - self.file_io.delete_quietly(path_str) - for extra_file in file_meta.extra_files: - self.file_io.delete_quietly(self._aligned_extra_file_path(file_meta, extra_file)) + for path_to_delete in file_meta.collect_files(): + self.file_io.delete_quietly(path_to_delete) except Exception as e: import logging logger = logging.getLogger(__name__) - path_to_delete = file_meta.external_path if file_meta.external_path else file_meta.file_path + path_to_delete = file_meta.physical_path() logger.warning(f"Failed to delete file {path_to_delete} during abort: {e}") @abstractmethod @@ -511,12 +507,7 @@ def _row_sidecar_fields(self, data: pa.Table) -> List: @staticmethod def _aligned_extra_file_path(file_meta: DataFileMeta, extra_file: str) -> str: - if "://" in extra_file or extra_file.startswith("/"): - return extra_file - file_path = file_meta.external_path if file_meta.external_path else file_meta.file_path - if not file_path or "/" not in file_path: - return extra_file - return f"{file_path.rsplit('/', 1)[0]}/{extra_file}" + return file_meta.aligned_file_path(extra_file) @staticmethod def _find_optimal_split_point(data: pa.RecordBatch, target_size: int) -> int: From 9e79d1ffe9d8c4f212b4c666c2854b9a0b43af3c Mon Sep 17 00:00:00 2001 From: leaves12138 <41894543+leaves12138@users.noreply.github.com> Date: Mon, 28 Sep 2026 20:33:54 +0800 Subject: [PATCH 2/3] test(python): exercise real file metadata in abort cleanup regressions --- .../pypaimon/tests/file_store_commit_test.py | 39 ++++++++++++++++++- 1 file changed, 38 insertions(+), 1 deletion(-) diff --git a/paimon-python/pypaimon/tests/file_store_commit_test.py b/paimon-python/pypaimon/tests/file_store_commit_test.py index 4f8778a23aa5..559b76e56785 100644 --- a/paimon-python/pypaimon/tests/file_store_commit_test.py +++ b/paimon-python/pypaimon/tests/file_store_commit_test.py @@ -19,10 +19,13 @@ import uuid from dataclasses import replace from datetime import datetime +from pathlib import Path +from tempfile import TemporaryDirectory from unittest.mock import MagicMock, Mock, patch from pypaimon.common.options.core_options import CoreOptions from pypaimon.common.options.options import Options +from pypaimon.filesystem.local_file_io import LocalFileIO from pypaimon.manifest.schema.data_file_meta import DataFileMeta from pypaimon.manifest.schema.manifest_entry import ManifestEntry from pypaimon.manifest.schema.manifest_file_meta import ManifestFileMeta @@ -78,15 +81,49 @@ def test_overwrite_validates_message_baseline(self): class TestAbortCommitMessages(unittest.TestCase): + @staticmethod + def _file_meta(**kwargs): + return DataFileMeta( + file_name='data.parquet', file_size=1, row_count=1, + min_key=None, max_key=None, key_stats=None, value_stats=None, + min_sequence_number=0, max_sequence_number=0, schema_id=0, + level=0, extra_files=kwargs.pop('extra_files', []), **kwargs) + def test_reconstructs_local_path_after_wire_decode(self): table = Mock() table.path_factory.return_value.bucket_path.return_value = '/table/p=1/bucket-0' - file = Mock(file_name='data.parquet', external_path=None, file_path=None) + file = self._file_meta() message = CommitMessage((1,), 0, [file]) _abort_commit_messages(table, [message]) table.file_io.delete_quietly.assert_called_once_with( '/table/p=1/bucket-0/data.parquet') + def test_reconstructs_aligned_sidecars_after_wire_decode(self): + table = Mock() + table.path_factory.return_value.bucket_path.return_value = '/table/p=1/bucket-0' + file = self._file_meta(extra_files=['data.parquet.index']) + _abort_commit_messages(table, [CommitMessage((1,), 0, [file])]) + self.assertEqual(table.file_io.delete_quietly.call_args_list, [ + unittest.mock.call('/table/p=1/bucket-0/data.parquet'), + unittest.mock.call('/table/p=1/bucket-0/data.parquet.index'), + ]) + + def test_deletes_literal_external_paths_and_preserves_metadata(self): + with TemporaryDirectory() as directory: + parent = Path(directory) / 'pt=a%2Fb%25%3F%23' + parent.mkdir() + paths = [parent / 'data.parquet', parent / 'data.parquet.index'] + for path in paths: + path.touch() + external_path = 'file:' + paths[0].as_posix() + file = self._file_meta(external_path=external_path, extra_files=[paths[1].name]) + table = Mock(file_io=LocalFileIO()) + _abort_commit_messages(table, [CommitMessage(('a/b%?#',), 0, [file])]) + self.assertFalse(any(path.exists() for path in paths)) + self.assertEqual(file.external_path, external_path) + self.assertEqual(file.extra_files, [paths[1].name]) + table.path_factory.assert_not_called() + def test_index_path_failure_does_not_escape_abort(self): table = Mock() table.path_factory.side_effect = RuntimeError("path lookup failed") From 2ad0b9a9d24c2a881b88adda5e7f683aeeac56a5 Mon Sep 17 00:00:00 2001 From: JingsongLi Date: Mon, 28 Sep 2026 21:39:49 +0800 Subject: [PATCH 3/3] [python][test] Use real file metadata for buffered writer cleanup --- .../pypaimon/tests/write/write_buffer_test.py | 49 ++++++++++++++----- 1 file changed, 37 insertions(+), 12 deletions(-) diff --git a/paimon-python/pypaimon/tests/write/write_buffer_test.py b/paimon-python/pypaimon/tests/write/write_buffer_test.py index 0c539d62827a..942798b75a21 100644 --- a/paimon-python/pypaimon/tests/write/write_buffer_test.py +++ b/paimon-python/pypaimon/tests/write/write_buffer_test.py @@ -31,6 +31,7 @@ import pyarrow as pa +from pypaimon.manifest.schema.data_file_meta import DataFileMeta from pypaimon.write.writer.append_only_data_writer import AppendOnlyDataWriter from pypaimon.write.writer.data_vector_writer import DataVectorWriter from pypaimon.write.writer.data_writer import DataWriter @@ -482,15 +483,14 @@ def test_failed_normal_flush_keeps_the_rows_for_the_retry(self): self.assertEqual(writer._normal_buffer.num_rows, 0) -class _StubMeta: - """The handful of ``DataFileMeta`` fields the flush and abort paths read.""" - - def __init__(self, row_count: int, file_name: str): - self.row_count = row_count - self.file_name = file_name - self.file_path = '/warehouse/%s' % file_name - self.external_path = None - self.extra_files = [] +def _file_meta(row_count: int, file_name: str) -> DataFileMeta: + """Use real metadata so flush and abort exercise the shared path resolver.""" + return DataFileMeta( + file_name=file_name, file_size=1, row_count=row_count, + min_key=None, max_key=None, key_stats=None, value_stats=None, + min_sequence_number=0, max_sequence_number=row_count - 1, + schema_id=0, level=0, extra_files=[], + file_path='/warehouse/%s' % file_name) class _RecordingFileIO: @@ -528,7 +528,7 @@ def prepare_commit(self): raise IOError('transient sidecar failure') if not self.committed_files: self.committed_files.append( - _StubMeta(self._row_count, self._file_name)) + _file_meta(self._row_count, self._file_name)) return self.committed_files.copy() def _release_prepared_files(self): @@ -581,7 +581,7 @@ def _write_normal_data_to_file(self, data: pa.Table): self._fail_normal_times -= 1 raise IOError('transient storage failure') self.written.append(data) - return _StubMeta(data.num_rows, 'data-%d' % len(self.written)) + return _file_meta(data.num_rows, 'data-%d' % len(self.written)) class _DedicatedHarness(DedicatedFormatWriter): def __init__(self, blob_writers, vector_writer=None): @@ -603,7 +603,7 @@ def __init__(self, blob_writers, vector_writer=None): def _write_normal_data_to_file(self, data: pa.Table): self.written.append(data) - return _StubMeta(data.num_rows, 'data-%d' % len(self.written)) + return _file_meta(data.num_rows, 'data-%d' % len(self.written)) def test_failed_sidecar_publishes_nothing_and_the_retry_resumes(self): vector = _StubSidecarWriter(3, 'vector-0', fail_times=1) @@ -690,6 +690,31 @@ def test_abort_deletes_the_unpublished_normal_file(self): self.assertIsNone(writer._pending_normal_meta) self.assertTrue(vector.aborted) + def test_abort_deletes_external_unpublished_file_and_sidecars(self): + for kind in ('vector', 'blob'): + with self.subTest(kind=kind): + sidecar = _StubSidecarWriter(3, 'sidecar-0', fail_times=1) + writer = (self._VectorHarness(sidecar) if kind == 'vector' + else self._DedicatedHarness({'payload': sidecar})) + writer._normal_buffer.append(_table(0, 3)) + with self.assertRaises(IOError): + writer._close_current_writers() + + meta = writer._pending_normal_meta + external_path = 'file:/external/pt=a%2Fb%25/data-1' + meta.external_path = external_path + meta.extra_files = ['data-1.index'] + writer.abort() + + self.assertEqual(writer.file_io.deleted, [ + '/external/pt=a%2Fb%25/data-1', + '/external/pt=a%2Fb%25/data-1.index', + ]) + self.assertEqual(meta.external_path, external_path) + self.assertEqual(writer.committed_files, []) + self.assertIsNone(writer._pending_normal_meta) + self.assertTrue(sidecar.aborted) + def test_dedicated_writer_failed_blob_phase_publishes_nothing(self): blob = _StubSidecarWriter(3, 'blob-0', fail_times=1) writer = self._DedicatedHarness({'payload': blob})