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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
26 changes: 22 additions & 4 deletions paimon-python/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.

Expand All @@ -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
Expand Down
7 changes: 5 additions & 2 deletions paimon-python/pypaimon/tests/blob_table_test.py
Original file line number Diff line number Diff line change
Expand Up @@ -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()
Expand Down Expand Up @@ -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")
Expand Down
15 changes: 6 additions & 9 deletions paimon-python/pypaimon/tests/blob_test.py
Original file line number Diff line number Diff line change
Expand Up @@ -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)

Expand All @@ -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']))
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down
10 changes: 8 additions & 2 deletions paimon-python/pypaimon/tests/deferred_blob_resolve_test.py
Original file line number Diff line number Diff line change
Expand Up @@ -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))

Expand All @@ -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 = [
Expand Down
Loading
Loading