From c566c4d592b065bb6920e4bb9b5a4e2e458dcef0 Mon Sep 17 00:00:00 2001 From: JingsongLi Date: Mon, 28 Sep 2026 22:59:48 +0800 Subject: [PATCH 1/2] [python] Enable native scalar Blob Arrow writes --- paimon-python/README.md | 18 +- .../pypaimon/tests/blob_table_test.py | 7 +- paimon-python/pypaimon/tests/blob_test.py | 15 +- .../tests/data_evolution_row_rolling_test.py | 1 + .../tests/deferred_blob_resolve_test.py | 10 +- .../pypaimon/tests/native_blob_write_test.py | 368 ++++++++++++++++++ .../tests/native_external_write_test.py | 5 +- paimon-python/pypaimon/write/native_write.py | 11 +- .../write/writer/dedicated_format_writer.py | 15 +- 9 files changed, 416 insertions(+), 34 deletions(-) create mode 100644 paimon-python/pypaimon/tests/native_blob_write_test.py diff --git a/paimon-python/README.md b/paimon-python/README.md index a32becf69596..69f7fb709540 100644 --- a/paimon-python/README.md +++ b/paimon-python/README.md @@ -161,10 +161,11 @@ 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 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 +write supports Parquet append, primary-key and data-evolution tables, including +top-level scalar BLOB Arrow columns, on the same filesystem/JDBC publication +route as native commit. ARRAY/MAP BLOB, video and optional data-evolution row +sidecars select the Python writer. 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. @@ -173,6 +174,15 @@ Native writes honor `data-file.path-directory` and the configured when the write destinations change. Python and native readers and committers can exchange these files, including external data files and their index sidecars. +BLOB Arrow values may contain payload bytes or serialized descriptors. Fields +listed in `blob-descriptor-field` remain inline and require descriptors. HTTP(S) +references use decoded response streams, including gzip and deflate. Blob +files roll by payload size, independently of the normal Parquet files. Python +`Blob` row objects can provide custom streams or URI readers; `write_row` on a +BLOB table selects the Python writer before any native data is written. Switching +to row writes after native Arrow writes is rejected. Use `write.native.enabled=false` when +mixing Arrow batches and Python `Blob` objects in one writer. + Both native options are disabled by default. # Native commit diff --git a/paimon-python/pypaimon/tests/blob_table_test.py b/paimon-python/pypaimon/tests/blob_table_test.py index 9a4172f0f9a8..ae4c5aa0320c 100755 --- a/paimon-python/pypaimon/tests/blob_table_test.py +++ b/paimon-python/pypaimon/tests/blob_table_test.py @@ -295,7 +295,9 @@ def test_dedicated_format_writer_schema_detection(self): } ) self.catalog.create_table('test_db.blob_detection_test', schema, False) - table = self.catalog.get_table('test_db.blob_detection_test') + # This test inspects the Python writer's internal column routing. + table = self.catalog.get_table('test_db.blob_detection_test').copy( + {'write.native.enabled': 'false'}) # Use proper table API to create writer write_builder = table.new_batch_write_builder() @@ -2494,7 +2496,8 @@ def test_blob_write_read_partition(self): table_scan = read_builder.new_scan() table_read = read_builder.new_read() splits = table_scan.plan().splits() - result = table_read.to_arrow(splits) + # Scans do not promise an ordering across partitions. + result = table_read.to_arrow(splits).sort_by('id') # Verify the data was read back correctly self.assertEqual(result.num_rows, 5, "Should have 5 rows") diff --git a/paimon-python/pypaimon/tests/blob_test.py b/paimon-python/pypaimon/tests/blob_test.py index 21f2dc14e780..9beb2829cb79 100644 --- a/paimon-python/pypaimon/tests/blob_test.py +++ b/paimon-python/pypaimon/tests/blob_test.py @@ -1618,7 +1618,7 @@ def test_blob_descriptor_fields_ignores_legacy_stored_key(self): })) self.assertEqual(set(), blank_canonical.blob_descriptor_fields()) - def test_dedicated_writer_accepts_exact_v1_descriptor_bytes(self): + def test_dedicated_writer_accepts_java_descriptor_prefixes(self): from pypaimon.write.writer.dedicated_format_writer import ( DedicatedFormatWriter) @@ -1642,19 +1642,16 @@ def test_dedicated_writer_accepts_exact_v1_descriptor_bytes(self): pa.RecordBatch.from_arrays( [pa.array([v1], type=pa.large_binary())], names=['payload'])) - padded = pa.RecordBatch.from_arrays( - [pa.array([v1 + b"x"], type=pa.large_binary())], names=['payload']) - with self.assertRaisesRegex(ValueError, "trailing bytes"): - writer._validate_inline_stored_fields_input(padded) - + # Java's known-descriptor parser accepts padding and the older + # no-magic layout. Do not apply heuristic exact-length detection here. v0 = bytes([0]) + v1[1:] - with self.assertRaisesRegex(ValueError, r"in \[1, 2\], but found 0"): + for value in (v0, v1 + b"x", v2 + b"padding"): writer._validate_inline_stored_fields_input( pa.RecordBatch.from_arrays( - [pa.array([v0], type=pa.large_binary())], names=['payload'])) + [pa.array([value], type=pa.large_binary())], names=['payload'])) v3 = bytes([3]) + v2[1:] - with self.assertRaisesRegex(ValueError, r"in \[1, 2\], but found 3"): + with self.assertRaisesRegex(ValueError, "serialized BlobDescriptor"): writer._validate_inline_stored_fields_input( pa.RecordBatch.from_arrays( [pa.array([v3], type=pa.large_binary())], names=['payload'])) diff --git a/paimon-python/pypaimon/tests/data_evolution_row_rolling_test.py b/paimon-python/pypaimon/tests/data_evolution_row_rolling_test.py index 6f7fb20421d8..d36da758a4ab 100644 --- a/paimon-python/pypaimon/tests/data_evolution_row_rolling_test.py +++ b/paimon-python/pypaimon/tests/data_evolution_row_rolling_test.py @@ -208,6 +208,7 @@ def test_non_de_table_still_fails_fast(self): NotImplementedError, 'row-count based file rolling'): tw.write_arrow(self._rows(4)) + @pytest.mark.python_write def test_blob_writer_supports_target_file_row_num(self): table = self._create_with_schema( self.blob_schema, diff --git a/paimon-python/pypaimon/tests/deferred_blob_resolve_test.py b/paimon-python/pypaimon/tests/deferred_blob_resolve_test.py index ee013017d6a4..3c8a1a071a6c 100644 --- a/paimon-python/pypaimon/tests/deferred_blob_resolve_test.py +++ b/paimon-python/pypaimon/tests/deferred_blob_resolve_test.py @@ -285,7 +285,10 @@ def test_iterator_passes_remaining_limit_across_splits(self): read_builder = table.new_read_builder().with_projection( ["sample_id", "payload", "score"] ).with_limit(2) - splits = read_builder.new_scan().plan().splits() + # The reader's remaining-limit behavior needs both partitions, even + # when a planner could satisfy the limit with a single partition. + splits = sorted(table.new_read_builder().new_scan().plan().splits(), + key=lambda split: split.partition.values) rows = list(read_builder.new_read().to_iterator(splits)) @@ -308,7 +311,10 @@ def test_iterator_applies_limit_after_auth_filter(self): auth_result = _RejectScoreOneAuthResult() splits = [ QueryAuthSplit(split, auth_result) - for split in read_builder.new_scan().plan().splits() + # Attach authorization after planning; do not push the reader's + # limit into this unauthenticated scan. Order partitions explicitly. + for split in sorted(table.new_read_builder().new_scan().plan().splits(), + key=lambda split: split.partition.values) ] scores = [ diff --git a/paimon-python/pypaimon/tests/native_blob_write_test.py b/paimon-python/pypaimon/tests/native_blob_write_test.py new file mode 100644 index 000000000000..74aa82cd612b --- /dev/null +++ b/paimon-python/pypaimon/tests/native_blob_write_test.py @@ -0,0 +1,368 @@ +# 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. + +"""Blob Arrow writes must retain row identity across physical file rolls.""" + +import io +from contextlib import ExitStack +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.table.row.blob import Blob, BlobDescriptor +from pypaimon.table.row.generic_row import GenericRow +from pypaimon.write.native_write import NativeTableWrite + +pytestmark = pytest.mark.native_plan +_SCHEMA = pa.schema([('id', pa.int32()), ('large', pa.large_binary()), ('small', pa.large_binary())]) + + +def _table(tmp_path, options=None): + catalog = CatalogFactory.create({'warehouse': str(tmp_path / 'warehouse')}) + catalog.create_database('db', True) + settings = {'row-tracking.enabled': 'true', 'data-evolution.enabled': 'true', + 'write.native.enabled': 'true', 'target-file-row-num': '2', + 'blob.target-file-size': '32 B', 'file-index.bloom-filter.columns': 'id', + 'file-index.bloom-filter.id.items': '10', 'file-index.in-manifest-threshold': '0 B'} + settings.update(options or {}) + fields = [DataField(0, 'id', AtomicType('INT')), + DataField(1, 'large', AtomicType('BLOB')), DataField(2, 'small', AtomicType('BLOB'))] + catalog.create_table('db.t', Schema(fields=fields, options=settings), False) + return catalog.get_table('db.t') + + +def _data(start=0): + return pa.table({'id': [start, start + 1], + 'large': [bytes([start + 1]) * 40, None], + 'small': [b'', b'abc']}, schema=_SCHEMA) + + +def _read(table, planner, reader): + table = table.copy({'scan.native-plan.enabled': str(planner).lower(), + 'read.native.enabled': str(reader).lower()}) + 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'))) + return sorted(read.to_arrow(plan.splits()).to_pylist(), key=lambda row: row['id']) + + +def _physical_files(tmp_path): + return {path for path in tmp_path.rglob('*') + if path.is_file() and path.suffix in ('.parquet', '.blob', '.index')} + + +@pytest.mark.parametrize('external', [False, True]) +@pytest.mark.parametrize('optimize', [False, True]) +@pytest.mark.parametrize('native_commit', [False, True]) +def test_native_blob_stream_rolls_and_reuses_writer(tmp_path, external, optimize, native_commit): + options = {'data-evolution.write-cols-optimization.enabled': str(optimize).lower(), + 'commit.native.enabled': str(native_commit).lower()} + if external: + options.update({'data-file.external-paths': (tmp_path / 'external').as_uri(), + 'data-file.external-paths.strategy': 'entropy-inject'}) + table = _table(tmp_path, options) + builder = table.new_stream_write_builder() + writer, commit = builder.new_write(), builder.new_commit() + expected = [] + try: + assert isinstance(writer, NativeTableWrite) + for identifier in (1, 2): + for start in range((identifier - 1) * 6, identifier * 6, 2): + data = _data(start) + writer.write_arrow(data) + expected.extend(data.to_pylist()) + messages = writer.prepare_commit(identifier) + files = [file for message in messages for file in message.new_files] + normal = [file for file in files if file.file_name.endswith('.parquet')] + assert len(normal) == 3 + assert all(file.row_count == 2 and file.extra_files for file in normal) + assert all(file.write_cols == (None if optimize else ['id']) for file in normal) + assert all((file.min_sequence_number, file.max_sequence_number) == (0, file.row_count - 1) + for file in files) + assert all(bool(file.external_path) == external for file in files) + # Each normal file owns the following Blob files until the next normal file. + groups = [] + for file in files: + if file.file_name.endswith('.parquet'): + groups.append({'large': [], 'small': []}) + else: + groups[-1][file.write_cols[0]].append(file.row_count) + assert groups == [{'large': [1, 1], 'small': [2]}] * 3 + commit.commit(messages, identifier) + assert writer._python_writer is None + for planner in (False, True): + for reader in (False, True): + assert _read(table, planner, reader) == expected + finally: + writer.close() + commit.close() + + +@pytest.mark.parametrize('external', [False, True]) +@pytest.mark.parametrize('action', ['abort', 'close', 'failed-write']) +def test_native_blob_cleanup_includes_completed_groups(tmp_path, external, action): + options = {} + if external: + options.update({'data-file.external-paths': (tmp_path / 'external').as_uri(), + 'data-file.external-paths.strategy': 'entropy-inject'}) + table = _table(tmp_path, options) + writer = table.new_batch_write_builder().new_write() + try: + writer.write_arrow(_data()) + completed = _physical_files(tmp_path) + assert {path.suffix for path in completed} == {'.parquet', '.blob', '.index'} + if action == 'failed-write': + missing = BlobDescriptor(str(tmp_path / 'missing'), 0, 40).serialize() + with pytest.raises(Exception, match='missing'): + writer.write_arrow(pa.table({'id': [2], 'large': [missing], 'small': [b'ok']}, + schema=_SCHEMA)) + else: + # Leave another normal group open when aborting/closing. + writer.write_arrow(_data(2).slice(0, 1)) + getattr(writer, action)() + assert not _physical_files(tmp_path) + assert table.snapshot_manager().get_latest_snapshot() is None + finally: + writer.close() + + +class _StreamingBlob(Blob): + def __init__(self, data): + self.data = data + self.opened = False + + def to_data(self): + raise AssertionError('Blob must stay streaming') + + def to_descriptor(self): + raise RuntimeError('No descriptor') + + def new_input_stream(self): + self.opened = True + return io.BytesIO(self.data) + + +def test_blob_object_selects_python_before_writing(tmp_path): + table = _table(tmp_path) + builder = table.new_batch_write_builder() + writer, commit = builder.new_write(), builder.new_commit() + blob = _StreamingBlob(b'stream') + try: + assert isinstance(writer, NativeTableWrite) + writer.write_row(GenericRow([0, b'bytes', None], table.fields)) + writer.write_row(GenericRow([1, blob, None], table.fields)) + writer.write_arrow(_data(2)) + assert writer._python_writer is not None + commit.commit(writer.prepare_commit()) + assert blob.opened + expected = [{'id': 0, 'large': b'bytes', 'small': None}, + {'id': 1, 'large': b'stream', 'small': None}] + _data(2).to_pylist() + assert _read(table, True, True) == expected + finally: + writer.close() + commit.close() + + +def test_blob_object_cannot_switch_after_native_write(tmp_path): + table = _table(tmp_path) + builder = table.new_batch_write_builder() + writer, commit = builder.new_write(), builder.new_commit() + blob = _StreamingBlob(b'stream') + try: + writer.write_arrow(_data()) + with pytest.raises(RuntimeError, match='after native data was written'): + writer.write_row(GenericRow([2, blob, None], table.fields)) + assert not blob.opened + commit.commit(writer.prepare_commit()) + assert _read(table, True, True) == _data().to_pylist() + finally: + writer.close() + commit.close() + + +@pytest.mark.parametrize('native', [False, True]) +@pytest.mark.parametrize('trailing', [b'', b'padding']) +def test_inline_descriptors_keep_v1_v2_and_null_semantics(tmp_path, native, trailing): + source = tmp_path / 'source' + source.write_bytes(b'prefix-payload-suffix') + v2 = BlobDescriptor(str(source), 7, 7).serialize() + v1 = bytes([1]) + v2[9:] + trailing + v2 += trailing + table = _table(tmp_path, {'blob-descriptor-field': 'large', + 'write.native.enabled': str(native).lower()}) + data = pa.table({'id': [0, 1, 2], 'large': [v1, v2, None], 'small': [b'a', None, b'']}, schema=_SCHEMA) + builder = table.new_batch_write_builder() + writer, commit = builder.new_write(), builder.new_commit() + try: + assert isinstance(writer, NativeTableWrite) == native + writer.write_arrow(data) + commit.commit(writer.prepare_commit()) + finally: + writer.close() + commit.close() + expected = data.to_pylist() + expected[0]['large'] = expected[1]['large'] = b'payload' + for planner in (False, True): + for reader in (False, True): + assert _read(table, planner, reader) == expected + assert _read(table.copy({'blob-as-descriptor': 'true'}), planner, reader)[0]['large'] == v1 + + +@pytest.fixture +def http_blob_uri(): + from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer + from threading import Thread + + requests = [] + + class Handler(BaseHTTPRequestHandler): + def do_GET(self): + requests.append(self.path) + if self.path == '/large': + self.send_response(200) + self.send_header('Content-Length', str(7 + 18 * 1024 * 1024)) + self.end_headers() + self.wfile.write(b'prefix-') + for _ in range(18): + self.wfile.write(b'x' * 1024 * 1024) + return + if self.path in ('/gzip', '/deflate'): + import gzip + import zlib + payload = b'prefix-payload' + body = gzip.compress(payload) if self.path == '/gzip' else zlib.compress(payload) + self.send_response(200) + self.send_header('Content-Encoding', self.path[1:]) + self.send_header('Content-Length', str(len(body))) + self.end_headers() + self.wfile.write(body) + return + if self.path == '/missing': + self.send_error(404) + return + # Deliberately ignore Range and omit Content-Length. Native must + # stream past the prefix and determine EOF for length=-1. + self.send_response(200) + self.end_headers() + self.wfile.write(b'prefix-payload') + + def log_message(self, *args): + pass + + server = ThreadingHTTPServer(('127.0.0.1', 0), Handler) + thread = Thread(target=server.serve_forever, daemon=True) + thread.start() + try: + yield 'http://127.0.0.1:{}'.format(server.server_port), requests + finally: + server.shutdown() + thread.join() + server.server_close() + + +@pytest.mark.parametrize('encoding', ['payload', 'gzip', 'deflate']) +@pytest.mark.parametrize('inline', [False, True]) +@pytest.mark.parametrize('length', [3, -1]) +@pytest.mark.parametrize('native', [False, True]) +def test_http_descriptors_follow_uri_scheme(tmp_path, http_blob_uri, inline, length, native, encoding): + uri, _ = http_blob_uri + options = {'write.native.enabled': str(native).lower()} + if inline: + options['blob-descriptor-field'] = 'large' + table = _table(tmp_path, options) + descriptor = BlobDescriptor(uri + '/' + encoding, 7, length).serialize() + data = pa.table({'id': [0], 'large': [descriptor], 'small': [None]}, schema=_SCHEMA) + builder = table.new_batch_write_builder() + writer, commit = builder.new_write(), builder.new_commit() + try: + assert isinstance(writer, NativeTableWrite) == native + # Write a preceding local payload to ensure HTTP input works after + # native ownership is established, without any late fallback. + prior = data if inline else _data(1) + writer.write_arrow(prior) + writer.write_arrow(data) + commit.commit(writer.prepare_commit()) + if native: + assert writer._python_writer is None + expected = [{'id': 0, 'large': b'pay' if length == 3 else b'payload', 'small': None}] + expected += expected if inline else prior.to_pylist() + for planner in (False, True): + for reader in (False, True): + assert _read(table, planner, reader) == expected + finally: + writer.close() + commit.close() + + +def test_http_descriptor_failure_cleans_native_files(tmp_path, http_blob_uri): + uri, _ = http_blob_uri + table = _table(tmp_path) + writer = table.new_batch_write_builder().new_write() + try: + writer.write_arrow(_data()) + descriptor = BlobDescriptor(uri + '/missing', 0, 5).serialize() + with pytest.raises(Exception, match='404'): + writer.write_arrow(pa.table({'id': [2], 'large': [descriptor], 'small': [None]}, schema=_SCHEMA)) + assert not _physical_files(tmp_path) + finally: + writer.close() + + +def test_large_http_descriptor_uses_one_stream_for_all_copy_chunks(tmp_path, http_blob_uri): + uri, requests = http_blob_uri + table = _table(tmp_path) + builder = table.new_batch_write_builder() + writer, commit = builder.new_write(), builder.new_commit() + length = 18 * 1024 * 1024 # More than two native 8 MiB copy buffers. + descriptor = BlobDescriptor(uri + '/large', 7, length).serialize() + try: + writer.write_arrow(pa.table({'id': [0], 'large': [descriptor], 'small': [None]}, schema=_SCHEMA)) + commit.commit(writer.prepare_commit()) + assert requests == ['/large'] + assert _read(table, True, True) == [{'id': 0, 'large': b'x' * length, 'small': None}] + finally: + writer.close() + commit.close() + + +def test_native_blob_row_target_rolls_normal_files_at_batch_boundaries(tmp_path): + table = _table(tmp_path, {'target-file-row-num': '3', 'blob.target-file-size': '10 MB'}) + data = pa.table({'id': list(range(7)), 'large': [b'a'] * 7, 'small': [None] * 7}, 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() + files = [file for message in messages for file in message.new_files] + assert [file.row_count for file in files if file.file_name.endswith('.parquet')] == [7] + for name in ('large', 'small'): + assert [file.row_count for file in files if file.write_cols == [name]] == [3, 3, 1] + commit.commit(messages) + for planner in (False, True): + for reader in (False, True): + assert _read(table, planner, reader) == data.to_pylist() + finally: + writer.close() + commit.close() diff --git a/paimon-python/pypaimon/tests/native_external_write_test.py b/paimon-python/pypaimon/tests/native_external_write_test.py index f1bed9ddd59b..9c45d6fc4c79 100644 --- a/paimon-python/pypaimon/tests/native_external_write_test.py +++ b/paimon-python/pypaimon/tests/native_external_write_test.py @@ -184,7 +184,7 @@ def test_external_updates_upserts_deletes_history_and_abort(tmp_path, strategy, @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): +def test_external_blob_writer_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() @@ -196,8 +196,7 @@ def test_external_blob_fallback_and_readers(tmp_path, strategy, native): 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) + _write(table, data, native) for planner in (False, True): for reader in (False, True): assert _read(table, planner, reader) == data.to_pylist() diff --git a/paimon-python/pypaimon/write/native_write.py b/paimon-python/pypaimon/write/native_write.py index f01c63688761..ae74f831f56b 100644 --- a/paimon-python/pypaimon/write/native_write.py +++ b/paimon-python/pypaimon/write/native_write.py @@ -22,7 +22,7 @@ from pypaimon.common.options.core_options import MergeEngine from pypaimon.schema.arrow_schema import arrow_schemas_compatible, normalize_arrow_strings -from pypaimon.schema.data_types import PyarrowFieldParser, is_blob_file_field +from pypaimon.schema.data_types import PyarrowFieldParser, is_blob_file_field, is_blob_type from pypaimon.table.bucket_mode import BucketMode from pypaimon.write.native_commit import ( create_native_write_table, from_native_commit_messages, @@ -81,7 +81,10 @@ def create_native_write(table, commit_user, static_partition=None, stream=False) or not _native_map_layouts_supported(table, schema) # Rust cannot encode these partition keys yet. or not _native_partition_types_supported(schema, table.partition_keys) - or any(is_blob_file_field(field) for field in table.table_schema.fields)): + # Native dedicated files currently support top-level scalar Blob fields. + or table.options.video_frame_fields() + or any(is_blob_file_field(field) and not is_blob_type(field.type) + for field in table.table_schema.fields)): return None native_table = create_native_write_table(table) if native_table is None: @@ -169,6 +172,10 @@ def write_pandas(self, dataframe): def write_row(self, row): if self._python_writer is not None: return self._python_writer.write_row(row) + if any(is_blob_file_field(field) for field in self.table.table_schema.fields): + # Select the row-aware writer from the first row, including byte + # values: later rows may provide custom Blob streams or URI readers. + return self._switch_to_python().write_row(row) values = row_to_named_values(row, self.table.table_schema.fields) names = list(self.table.field_names) self.write_arrow_batch(row_values_to_arrow_table( diff --git a/paimon-python/pypaimon/write/writer/dedicated_format_writer.py b/paimon-python/pypaimon/write/writer/dedicated_format_writer.py index 98fc208bef36..dd0bb4096aa0 100644 --- a/paimon-python/pypaimon/write/writer/dedicated_format_writer.py +++ b/paimon-python/pypaimon/write/writer/dedicated_format_writer.py @@ -403,14 +403,9 @@ def _validate_inline_stored_fields_input(self, data: pa.RecordBatch): "BlobDescriptor." ) descriptor_bytes = bytes(value) - if descriptor_bytes: - version = descriptor_bytes[0] - if version < 1 or version > BlobDescriptor.CURRENT_VERSION: - raise ValueError( - f"blob-descriptor-field requires BlobDescriptor version " - f"in [1, {BlobDescriptor.CURRENT_VERSION}], but found " - f"{version}." - ) + # Like Java, schema-declared descriptors use the prefix parser. + # Exact wire length is only needed when detecting descriptors + # among arbitrary payload bytes. try: BlobDescriptor.deserialize(descriptor_bytes) except Exception as e: @@ -418,10 +413,6 @@ def _validate_inline_stored_fields_input(self, data: pa.RecordBatch): "blob-descriptor-field requires blob field value to be a serialized " "BlobDescriptor." ) from e - # serialize() always emits CURRENT_VERSION, so a round-trip - # would reject exact v1 bytes. Check exact wire length instead. - if BlobDescriptor.parse_if_serialized(descriptor_bytes) is None: - raise ValueError("Descriptor payload contains trailing bytes.") for field_name in self.blob_view_fields: if field_name not in data.schema.names: From d234756442b51f2b327e0a32fe52a1dd08ec3f66 Mon Sep 17 00:00:00 2001 From: YeJunHao <41894543+leaves12138@users.noreply.github.com> Date: Tue, 29 Sep 2026 11:15:31 +0800 Subject: [PATCH 2/2] [python] Preserve native prepared-file abort ownership Track prepared native messages until handoff to the committer or close. Clean up prepared files on explicit abort without deleting files from successful or uncertain commit attempts. Cover Blob rolling, external paths, streaming reuse, and native REST publication failures; keep Ray Python fault injection on the Python writer. --- paimon-python/README.md | 8 ++ .../pypaimon/tests/native_blob_write_test.py | 104 ++++++++++++++++++ .../pypaimon/tests/native_write_test.py | 49 ++++++++- paimon-python/pypaimon/tests/ray_sink_test.py | 2 + .../pypaimon/write/commit_message.py | 1 + paimon-python/pypaimon/write/native_write.py | 18 ++- paimon-python/pypaimon/write/table_commit.py | 3 + 7 files changed, 182 insertions(+), 3 deletions(-) diff --git a/paimon-python/README.md b/paimon-python/README.md index 69f7fb709540..246af0693a72 100644 --- a/paimon-python/README.md +++ b/paimon-python/README.md @@ -183,6 +183,14 @@ BLOB table selects the Python writer before any native data is written. Switchin to row writes after native Arrow writes is rejected. Use `write.native.enabled=false` when mixing Arrow batches and Python `Blob` objects in one writer. +An explicit native writer `abort()` also deletes prepared files that have not +been passed to a PyPaimon committer. Calling `close()` instead releases those +files to the caller without deleting them; use `commit.abort(messages)` to +discard them after closing the writer. Once a commit attempt starts, writer +abort preserves its files even if the attempt raises, because a snapshot may +already reference them. Stream writers retain cleanup ownership only for +messages that have not been submitted to a committer. + Both native options are disabled by default. # Native commit diff --git a/paimon-python/pypaimon/tests/native_blob_write_test.py b/paimon-python/pypaimon/tests/native_blob_write_test.py index 74aa82cd612b..d78167e33a07 100644 --- a/paimon-python/pypaimon/tests/native_blob_write_test.py +++ b/paimon-python/pypaimon/tests/native_blob_write_test.py @@ -148,6 +148,110 @@ def test_native_blob_cleanup_includes_completed_groups(tmp_path, external, actio writer.close() +@pytest.mark.parametrize('external', [False, True]) +@pytest.mark.parametrize('stream', [False, True]) +def test_native_blob_abort_removes_prepared_and_outstanding_files(tmp_path, external, stream): + options = {'data-file.path-directory': 'data/nested'} + if external: + options['data-file.external-paths'] = (tmp_path / 'external').as_uri() + table = _table(tmp_path, options) + builder = table.new_stream_write_builder() if stream else table.new_batch_write_builder() + writer = builder.new_write() + try: + assert isinstance(writer, NativeTableWrite) + for identifier in range(1, 3 if stream else 2): + writer.write_arrow(_data(identifier * 2)) + messages = writer.prepare_commit(identifier) if stream else writer.prepare_commit() + assert messages + assert {path.suffix for path in _physical_files(tmp_path)} == {'.parquet', '.blob', '.index'} + writer.write_arrow(_data(6)) + writer.abort() + writer.abort() + assert not _physical_files(tmp_path) + assert table.snapshot_manager().get_latest_snapshot() is None + finally: + writer.close() + + +@pytest.mark.parametrize('external', [False, True]) +def test_native_blob_close_releases_prepared_files_to_committer(tmp_path, external): + options = {'data-file.external-paths': (tmp_path / 'external').as_uri()} if external else {} + table = _table(tmp_path, options) + builder = table.new_batch_write_builder() + writer, commit = builder.new_write(), builder.new_commit() + try: + writer.write_arrow(_data()) + messages = writer.prepare_commit() + prepared = _physical_files(tmp_path) + assert prepared + writer.close() + writer.abort() + assert _physical_files(tmp_path) == prepared + commit.commit(messages) + for planner in (False, True): + for reader in (False, True): + assert _read(table, planner, reader) == _data().to_pylist() + finally: + writer.close() + commit.close() + + +@pytest.mark.parametrize('external', [False, True]) +@pytest.mark.parametrize('commit_native_option', [False, True]) +def test_native_blob_stream_abort_preserves_submitted_files(tmp_path, external, commit_native_option): + options = {'commit.native.enabled': str(commit_native_option).lower()} + if external: + options['data-file.external-paths'] = (tmp_path / 'external').as_uri() + table = _table(tmp_path, options) + builder = table.new_stream_write_builder() + writer, commit = builder.new_write(), builder.new_commit() + try: + writer.write_arrow(_data()) + commit.commit(writer.prepare_commit(1), 1) + submitted = _physical_files(tmp_path) + writer.write_arrow(_data(2)) + assert writer.prepare_commit(2) + writer.write_arrow(_data(4)) + assert _physical_files(tmp_path) > submitted + writer.abort() + assert _physical_files(tmp_path) == submitted + for planner in (False, True): + for reader in (False, True): + assert _read(table, planner, reader) == _data().to_pylist() + finally: + writer.close() + commit.close() + + +@pytest.mark.parametrize('published', [False, True]) +def test_native_blob_abort_preserves_files_after_commit_exception(tmp_path, published): + table = _table(tmp_path) + builder = table.new_batch_write_builder() + writer, commit = builder.new_write(), builder.new_commit() + try: + writer.write_arrow(_data()) + messages = writer.prepare_commit() + prepared = _physical_files(tmp_path) + original_commit = commit.file_store_commit.commit + + def fail_commit(**kwargs): + if published: + original_commit(**kwargs) + raise RuntimeError('Commit outcome is unknown') + + with patch.object(commit.file_store_commit, 'commit', side_effect=fail_commit): + with pytest.raises(RuntimeError, match='Commit outcome is unknown'): + commit.commit(messages) + writer.abort() + assert _physical_files(tmp_path) == prepared + assert (table.snapshot_manager().get_latest_snapshot() is not None) == published + if published: + assert _read(table, True, True) == _data().to_pylist() + finally: + writer.close() + commit.close() + + class _StreamingBlob(Blob): def __init__(self, data): self.data = data diff --git a/paimon-python/pypaimon/tests/native_write_test.py b/paimon-python/pypaimon/tests/native_write_test.py index bfe8b638b483..ea6646e45789 100644 --- a/paimon-python/pypaimon/tests/native_write_test.py +++ b/paimon-python/pypaimon/tests/native_write_test.py @@ -16,7 +16,7 @@ """End-to-end coverage of the optional native data writer bridge.""" -from unittest.mock import patch +from unittest.mock import Mock, patch import pyarrow as pa import pytest @@ -257,6 +257,53 @@ def test_native_overwrite_and_empty_overwrite(tmp_path): assert _rows(table) == [{'id': 2, 'pt': 'b'}] +@requires_native +@pytest.mark.parametrize('failure', [None, 'before-publication', 'after-publication']) +def test_writer_abort_preserves_native_commit_files(native_rest_catalog, failure): + native_rest_catalog.create_table('default.abort_handoff', Schema.from_pyarrow_schema( + pa.schema([('id', pa.int64()), ('pt', pa.string())]), options={ + 'write.native.enabled': 'true', 'commit.native.enabled': 'true', + }), False) + table = native_rest_catalog.get_table('default.abort_handoff') + builder = table.new_batch_write_builder() + writer, commit = builder.new_write(), builder.new_commit() + try: + assert isinstance(writer, NativeTableWrite) + writer.write_arrow_batch(_batch([1, 2], ['a', 'b'])) + messages = writer.prepare_commit() + prepared = commit._prepare_native_commit(messages) + assert prepared is not None + native, native_messages = prepared + proxy = Mock(wraps=native) + + def publish(*args): + if failure != 'before-publication': + native.commit(*args) + if failure: + raise RuntimeError('Native commit outcome is unknown') + + proxy.commit.side_effect = publish + with patch.object(commit, '_prepare_native_commit', return_value=(proxy, native_messages)), \ + patch.object(commit.file_store_commit, 'commit', + side_effect=AssertionError('Python commit fallback')): + if failure: + with pytest.raises(RuntimeError, match='Native commit outcome is unknown'): + commit.commit(messages) + else: + commit.commit(messages) + proxy.commit.assert_called_once() + writer.abort() + assert all(table.file_io.exists(file.file_path) + for message in messages for file in message.new_files) + if failure != 'before-publication': + assert _rows(table) == [{'id': 1, 'pt': 'a'}, {'id': 2, 'pt': 'b'}] + else: + assert table.snapshot_manager().get_latest_snapshot() is None + finally: + writer.close() + commit.close() + + @requires_native def test_advanced_api_switches_before_write_and_rejects_late_switch(tmp_path): table = _table(tmp_path) diff --git a/paimon-python/pypaimon/tests/ray_sink_test.py b/paimon-python/pypaimon/tests/ray_sink_test.py index 72c565828e41..6153b37014ee 100644 --- a/paimon-python/pypaimon/tests/ray_sink_test.py +++ b/paimon-python/pypaimon/tests/ray_sink_test.py @@ -21,6 +21,7 @@ from unittest.mock import Mock, patch import pyarrow as pa +import pytest from ray.data._internal.execution.interfaces import TaskContext from pypaimon import CatalogFactory, Schema @@ -361,6 +362,7 @@ def test_postpone_worker_uses_driver_bucket_plan_without_manifest_scan(self): self.assertEqual({0, 1}, {message.bucket for message in messages}) self.assertEqual({2}, {message.total_buckets for message in messages}) + @pytest.mark.python_write def test_write_does_not_return_prepared_messages_when_dedicated_close_aborts(self): from pypaimon.write.writer.dedicated_format_writer import DedicatedFormatWriter diff --git a/paimon-python/pypaimon/write/commit_message.py b/paimon-python/pypaimon/write/commit_message.py index 1325260cd213..0a25ba5cf4ca 100644 --- a/paimon-python/pypaimon/write/commit_message.py +++ b/paimon-python/pypaimon/write/commit_message.py @@ -41,6 +41,7 @@ class CommitMessage: compact_changelog_files: List[DataFileMeta] = field(default_factory=list) compact_index_adds: List['IndexManifestEntry'] = field(default_factory=list) compact_index_deletes: List['IndexManifestEntry'] = field(default_factory=list) + _native_write_pending: bool = field(default=False, init=False, repr=False, compare=False) def is_empty(self): return ( diff --git a/paimon-python/pypaimon/write/native_write.py b/paimon-python/pypaimon/write/native_write.py index ae74f831f56b..8e00e1317198 100644 --- a/paimon-python/pypaimon/write/native_write.py +++ b/paimon-python/pypaimon/write/native_write.py @@ -24,6 +24,7 @@ from pypaimon.schema.arrow_schema import arrow_schemas_compatible, normalize_arrow_strings from pypaimon.schema.data_types import PyarrowFieldParser, is_blob_file_field, is_blob_type from pypaimon.table.bucket_mode import BucketMode +from pypaimon.write.file_store_commit import _abort_commit_messages from pypaimon.write.native_commit import ( create_native_write_table, from_native_commit_messages, ) @@ -115,6 +116,7 @@ def __init__(self, table, commit_user, static_partition, stream, native_writer): self._native_writer = native_writer self._python_writer = None self._written = False + self._prepared_messages = [] self._schema = PyarrowFieldParser.from_paimon_schema(table.table_schema.fields) def _switch_to_python(self): @@ -194,7 +196,13 @@ def prepare_commit(self, commit_identifier=None): if commit_identifier is not None: raise TypeError('BatchTableWrite.prepare_commit accepts no identifier') messages = self._native_writer.prepare_commit() - return from_native_commit_messages(self.table, messages) + messages = from_native_commit_messages(self.table, messages) + self._prepared_messages = [message for message in self._prepared_messages + if message._native_write_pending] + for message in messages: + message._native_write_pending = True + self._prepared_messages.extend(messages) + return messages def close(self): if self._python_writer is not None: @@ -202,9 +210,15 @@ def close(self): elif self._native_writer is not None: self._native_writer.close() self._native_writer = None + self._prepared_messages.clear() def abort(self): if self._python_writer is not None: self._python_writer.abort() else: - self.close() + messages = [message for message in self._prepared_messages + if message._native_write_pending] + try: + self.close() + finally: + _abort_commit_messages(self.table, messages) diff --git a/paimon-python/pypaimon/write/table_commit.py b/paimon-python/pypaimon/write/table_commit.py index ef2595d93c40..486dcbcff7b5 100644 --- a/paimon-python/pypaimon/write/table_commit.py +++ b/paimon-python/pypaimon/write/table_commit.py @@ -66,6 +66,9 @@ def _commit( commit_messages: List[CommitMessage], commit_identifier: int = BATCH_COMMIT_IDENTIFIER, snapshot_properties: Optional[Dict[str, str]] = None): + """Release native-writer cleanup ownership before publication can start.""" + for message in commit_messages: + message._native_write_pending = False non_empty_messages = [msg for msg in commit_messages if not msg.is_empty()] commit_kwargs = { "commit_messages": non_empty_messages,