diff --git a/paimon-python/README.md b/paimon-python/README.md index a32becf69596..246af0693a72 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,23 @@ 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. + +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/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..d78167e33a07 --- /dev/null +++ b/paimon-python/pypaimon/tests/native_blob_write_test.py @@ -0,0 +1,472 @@ +# 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() + + +@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 + 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/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 f01c63688761..8e00e1317198 100644 --- a/paimon-python/pypaimon/write/native_write.py +++ b/paimon-python/pypaimon/write/native_write.py @@ -22,8 +22,9 @@ 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.file_store_commit import _abort_commit_messages from pypaimon.write.native_commit import ( create_native_write_table, from_native_commit_messages, ) @@ -81,7 +82,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: @@ -112,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): @@ -169,6 +174,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( @@ -187,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: @@ -195,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, 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: