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):

Reply via email to