This is an automated email from the ASF dual-hosted git repository.
JingsongLi pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/paimon.git
The following commit(s) were added to refs/heads/master by this push:
new 701aee173d [Python] Align overwrite changelog with Java streaming
defaults (#10046)
701aee173d is described below
commit 701aee173d9b3b296687c4ab283c6cb2d52c38a0
Author: Jingsong Lee <[email protected]>
AuthorDate: Mon Sep 21 14:17:57 2026 +0800
[Python] Align overwrite changelog with Java streaming defaults (#10046)
---
.../pypaimon/tests/native_plan_incremental_test.py | 41 +++++--------
.../tests/write/changelog_producer_test.py | 70 ++++++++++++++++++++++
paimon-python/pypaimon/write/file_store_commit.py | 3 +-
paimon-python/pypaimon/write/table_write.py | 6 ++
4 files changed, 93 insertions(+), 27 deletions(-)
diff --git a/paimon-python/pypaimon/tests/native_plan_incremental_test.py
b/paimon-python/pypaimon/tests/native_plan_incremental_test.py
index 791041b753..17b4627b1a 100644
--- a/paimon-python/pypaimon/tests/native_plan_incremental_test.py
+++ b/paimon-python/pypaimon/tests/native_plan_incremental_test.py
@@ -493,43 +493,34 @@ def
test_streaming_changelog_frames_use_native_plan_and_read(catalog):
]
[email protected]_plan
-def test_streaming_overwrite_changelog_uses_java_follow_up_semantics(catalog):
+def test_streaming_overwrite_is_skipped_by_default_like_java(catalog):
+ import asyncio
+
table = _table(catalog, 'native_overwrite_changelog', True, {
'bucket': '1',
'changelog-producer': 'input',
})
_write(table, 100, [{'k': 1, 'v': 'before'}])
_write(table, 200, [{'k': 2, 'v': 'after'}], overwrite=True)
- snapshot = table.snapshot_manager().get_latest_snapshot()
- assert snapshot.commit_kind == 'OVERWRITE'
- assert snapshot.changelog_manifest_list is not None
+ overwrite = table.snapshot_manager().get_latest_snapshot()
+ assert overwrite.commit_kind == 'OVERWRITE'
+ assert overwrite.changelog_manifest_list is None
+ _write(table, 300, [{'k': 3, 'v': 'next'}])
- native_table = table.copy({
- 'scan.native-plan.enabled': 'true',
- 'read.native.enabled': 'true',
- })
- builder = (native_table.new_stream_read_builder()
+ builder = (table.new_stream_read_builder()
.with_projection(['v'])
.with_include_row_kind())
scan = builder.new_streaming_scan()
+ scan.next_snapshot_id = overwrite.id
- with patch.object(
- scan, '_try_native_plan',
- side_effect=AssertionError(
- 'OVERWRITE must not use range-based native changelog
planning')):
- plan = scan._create_changelog_plan(snapshot)
+ async def next_plan():
+ async for plan in scan.stream():
+ return plan
- assert plan.snapshot_id == snapshot.id
- assert all(
- file.file_name.startswith('changelog-')
- for split in plan.splits() for file in split.files)
- with patch(
- 'pypaimon.read.table_read.TableRead._create_split_read',
- side_effect=AssertionError(
- 'OVERWRITE changelog native read fell back to Python')):
- rows = builder.new_read().to_arrow(plan.splits()).to_pylist()
- assert rows == [{'_row_kind': '+I', 'v': 'after'}]
+ plan = asyncio.run(next_plan())
+ assert plan.snapshot_id == 3
+ assert builder.new_read().to_arrow(plan.splits()).to_pylist() == [
+ {'_row_kind': '+I', 'v': 'next'}]
def test_streaming_reader_honors_explicit_split_deletion_vector(catalog,
native, tmp_path):
diff --git a/paimon-python/pypaimon/tests/write/changelog_producer_test.py
b/paimon-python/pypaimon/tests/write/changelog_producer_test.py
index 2e2b2e20da..88faa792ac 100644
--- a/paimon-python/pypaimon/tests/write/changelog_producer_test.py
+++ b/paimon-python/pypaimon/tests/write/changelog_producer_test.py
@@ -175,6 +175,76 @@ class ChangelogProducerTest(unittest.TestCase):
table_write.close()
table_commit.close()
+ def test_input_mode_overwrite_has_no_changelog(self):
+ table_name = 'test_input_overwrite'
+ table = self._create_table(
+ table_name,
+ options={'changelog-producer': 'input', 'bucket': '1'}
+ )
+ append = table.new_batch_write_builder()
+ writer, commit = append.new_write(), append.new_commit()
+ try:
+ writer.write_arrow(self._sample_data())
+ commit.commit(writer.prepare_commit())
+ finally:
+ writer.close()
+ commit.close()
+
+ bucket_dir = os.path.join(
+ self.warehouse, 'default.db', table_name, 'dt=p1', 'bucket-0')
+ before_files = set(glob.glob(os.path.join(bucket_dir, 'changelog-*')))
+ self.assertTrue(before_files)
+
+ overwrite = table.new_batch_write_builder().overwrite()
+ writer, commit = overwrite.new_write(), overwrite.new_commit()
+ try:
+ writer.write_arrow(self._sample_data())
+ messages = writer.prepare_commit()
+ self.assertTrue(messages)
+ self.assertTrue(all(not message.changelog_files for message in
messages))
+ commit.commit(messages)
+ finally:
+ writer.close()
+ commit.close()
+
+ snapshot = table.snapshot_manager().get_latest_snapshot()
+ self.assertEqual(snapshot.commit_kind, 'OVERWRITE')
+ self.assertIsNone(snapshot.changelog_manifest_list)
+ self.assertEqual(
+ set(glob.glob(os.path.join(bucket_dir, 'changelog-*'))),
before_files)
+
+ append = table.new_batch_write_builder()
+ writer, commit = append.new_write(), append.new_commit()
+ try:
+ writer.write_arrow(self._sample_data())
+ commit.commit(writer.prepare_commit())
+ finally:
+ writer.close()
+ commit.close()
+ self.assertIsNotNone(
+
table.snapshot_manager().get_latest_snapshot().changelog_manifest_list)
+
+ def test_overwrite_commit_discards_supplied_changelog(self):
+ table = self._create_table(
+ 'test_overwrite_supplied_changelog',
+ options={'changelog-producer': 'input', 'bucket': '1'}
+ )
+ writer_builder = table.new_batch_write_builder()
+ writer = writer_builder.new_write()
+ overwrite_commit =
table.new_batch_write_builder().overwrite().new_commit()
+ try:
+ writer.write_arrow(self._sample_data())
+ messages = writer.prepare_commit()
+ self.assertTrue(any(message.changelog_files for message in
messages))
+ overwrite_commit.commit(messages)
+ finally:
+ writer.close()
+ overwrite_commit.close()
+
+ snapshot = table.snapshot_manager().get_latest_snapshot()
+ self.assertEqual(snapshot.commit_kind, 'OVERWRITE')
+ self.assertIsNone(snapshot.changelog_manifest_list)
+
def test_input_mode_changelog_manifest_readable(self):
table = self._create_table(
'test_input_readable',
diff --git a/paimon-python/pypaimon/write/file_store_commit.py
b/paimon-python/pypaimon/write/file_store_commit.py
index 27907f8741..d56fe5acf2 100644
--- a/paimon-python/pypaimon/write/file_store_commit.py
+++ b/paimon-python/pypaimon/write/file_store_commit.py
@@ -380,7 +380,6 @@ class FileStoreCommit:
else:
partition_filter =
self._create_static_partition_filter(overwrite_partition, commit_messages)
- changelog_entries = self._collect_changelog_entries(commit_messages)
index_adds = [
entry for message in commit_messages for entry in
message.index_adds
]
@@ -400,7 +399,7 @@ class FileStoreCommit:
commit_kind="OVERWRITE",
commit_identifier=commit_identifier,
commit_entries_plan=provider.provide,
- changelog_entries=changelog_entries,
+ changelog_entries=[],
detect_conflicts=True,
allow_rollback=False,
index_deletes=index_deletes,
diff --git a/paimon-python/pypaimon/write/table_write.py
b/paimon-python/pypaimon/write/table_write.py
index eea1319e12..d9952726e0 100644
--- a/paimon-python/pypaimon/write/table_write.py
+++ b/paimon-python/pypaimon/write/table_write.py
@@ -19,6 +19,7 @@ from typing import TYPE_CHECKING, Any, Dict, List, Optional
import pyarrow as pa
+from pypaimon.common.options.core_options import ChangelogProducer
from pypaimon.schema.arrow_schema import arrow_schemas_compatible,
normalize_arrow_strings
from pypaimon.schema.data_types import PyarrowFieldParser
from pypaimon.snapshot.snapshot import BATCH_COMMIT_IDENTIFIER
@@ -43,6 +44,11 @@ class TableWrite:
self.commit_user = commit_user
self.static_partition = static_partition
self.file_store_write = self._create_file_store_write(commit_user)
+ if static_partition is not None:
+ # An overwrite replaces state, not an input changelog. Java's
+ # overwrite commit does not publish changelog manifests; avoid
+ # writing unreferenced changelog files in the first place.
+ self.file_store_write.changelog_producer = ChangelogProducer.NONE
self.row_key_extractor =
self._create_row_key_extractor(static_partition)
def _create_file_store_write(self, commit_user):