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 60ae8eb83a [python] Reuse manifest merge results on commit retry 
(#9040)
60ae8eb83a is described below

commit 60ae8eb83abfe3d22683f65b7eedf74f65fdcc48
Author: XiaoHongbo <[email protected]>
AuthorDate: Wed Aug 5 14:55:07 2026 +0800

    [python] Reuse manifest merge results on commit retry (#9040)
---
 .../pypaimon/tests/file_store_commit_test.py       | 290 ++++++++++++++++++++-
 paimon-python/pypaimon/write/file_store_commit.py  | 110 +++++++-
 2 files changed, 391 insertions(+), 9 deletions(-)

diff --git a/paimon-python/pypaimon/tests/file_store_commit_test.py 
b/paimon-python/pypaimon/tests/file_store_commit_test.py
index f839704114..50f39d3ed3 100644
--- a/paimon-python/pypaimon/tests/file_store_commit_test.py
+++ b/paimon-python/pypaimon/tests/file_store_commit_test.py
@@ -22,13 +22,19 @@ from unittest.mock import MagicMock, Mock, patch
 
 from pypaimon.manifest.schema.data_file_meta import DataFileMeta
 from pypaimon.manifest.schema.manifest_entry import ManifestEntry
+from pypaimon.manifest.schema.manifest_file_meta import ManifestFileMeta
+from pypaimon.manifest.schema.simple_stats import SimpleStats
 from pypaimon.snapshot.snapshot_commit import PartitionStatistics
-from pypaimon.table.row.generic_row import GenericRow
+from pypaimon.table.row.binary_row import BinaryRow
+from pypaimon.table.row.generic_row import GenericRow, GenericRowSerializer
 from pypaimon.write.commit.row_id_conflict_rewriter import RowIdRewriteResult
 from pypaimon.write.commit_message import CommitMessage
 from pypaimon.write.file_store_commit import (
     FileStoreCommit,
+    ManifestMergeResult,
+    RetryResult,
     RewriteResult,
+    _try_replace_manifest_files,
 )
 
 
@@ -59,6 +65,288 @@ class TestFileStoreCommit(unittest.TestCase):
             commit_user='test_user'
         )
 
+    @staticmethod
+    def _manifest_meta(name):
+        row = GenericRowSerializer.to_bytes(GenericRow([], []))
+        return ManifestFileMeta(
+            file_name=name,
+            file_size=10,
+            num_added_files=1,
+            num_deleted_files=0,
+            partition_stats=SimpleStats(
+                BinaryRow(row, []), BinaryRow(row, []), []),
+            schema_id=0,
+        )
+
+    def test_replace_manifest_files_uses_stable_value_equality(
+            self, mock_manifest_list_manager, mock_manifest_file_manager):
+        previous = [self._manifest_meta('a'), self._manifest_meta('b')]
+        current = [
+            self._manifest_meta('prefix'),
+            self._manifest_meta('a'),
+            self._manifest_meta('b'),
+            self._manifest_meta('suffix'),
+        ]
+        replacement = [self._manifest_meta('merged')]
+
+        result = _try_replace_manifest_files(
+            current, previous, replacement)
+
+        self.assertEqual(
+            ['prefix', 'merged', 'suffix'],
+            [manifest.file_name for manifest in result],
+        )
+        self.assertIsNot(current[1], previous[0])
+
+    def test_replace_manifest_files_preserves_order_and_empty_semantics(
+            self, mock_manifest_list_manager, mock_manifest_file_manager):
+        a = self._manifest_meta('a')
+        b = self._manifest_meta('b')
+        merged = self._manifest_meta('merged')
+
+        self.assertIsNone(_try_replace_manifest_files(
+            [a, self._manifest_meta('x'), b], [a, b], [merged]))
+        self.assertEqual(
+            ['a', 'merged'],
+            [manifest.file_name for manifest in _try_replace_manifest_files(
+                [a, self._manifest_meta('a'), b], [a, b], [merged])],
+        )
+        self.assertEqual(
+            [merged], _try_replace_manifest_files([], [], [merged]))
+        self.assertIsNone(_try_replace_manifest_files([a], [], [merged]))
+
+    def test_manifest_merge_result_copies_and_freezes_lists(
+            self, mock_manifest_list_manager, mock_manifest_file_manager):
+        before = [self._manifest_meta('before')]
+        after = [self._manifest_meta('after')]
+
+        result = ManifestMergeResult(before, after)
+        before.clear()
+        after.clear()
+
+        self.assertIsInstance(result.merge_before_manifests, tuple)
+        self.assertIsInstance(result.merge_after_manifests, tuple)
+        self.assertEqual(
+            ['before'],
+            [manifest.file_name
+             for manifest in result.merge_before_manifests],
+        )
+        self.assertEqual(
+            ['after'],
+            [manifest.file_name
+             for manifest in result.merge_after_manifests],
+        )
+
+    def _run_manifest_commit_attempt(self, commit_side_effect=None,
+                                     commit_result=None, retry_result=None,
+                                     existing_manifests=None,
+                                     merged_manifests=None):
+        file_store_commit = self._create_file_store_commit()
+        self.mock_table.identifier = 'default.test_table'
+        self.mock_table.table_schema.id = 7
+        self.mock_table.options.row_tracking_enabled.return_value = False
+
+        snapshot_commit = MagicMock()
+        snapshot_commit.__enter__.return_value = snapshot_commit
+        snapshot_commit.__exit__.return_value = False
+        snapshot_commit.commit.side_effect = commit_side_effect
+        snapshot_commit.commit.return_value = commit_result
+        file_store_commit.snapshot_commit = snapshot_commit
+
+        before = self._manifest_meta('before')
+        after = self._manifest_meta('after')
+        existing_manifests = (
+            [before] if existing_manifests is None
+            else existing_manifests)
+        merged_manifests = (
+            [after] if merged_manifests is None
+            else merged_manifests)
+        delta = self._manifest_meta('delta')
+        file_store_commit._write_manifest_files = Mock(
+            return_value=[delta])
+        file_store_commit._generate_partition_statistics = Mock(
+            return_value=[])
+        file_store_commit.manifest_list_manager.read_all.return_value = (
+            existing_manifests)
+        file_store_commit.manifest_file_merger = Mock()
+        file_store_commit.manifest_file_merger.merge.return_value = (
+            merged_manifests, merged_manifests)
+        file_store_commit._clean_up_reuse_tmp_manifests = Mock()
+        file_store_commit._clean_up_no_reuse_tmp_manifests = Mock()
+
+        latest_snapshot = Mock(
+            id=3,
+            uuid='base-snapshot-uuid',
+            total_record_count=10,
+            index_manifest=None,
+        )
+        commit_entry = Mock(kind=0)
+        commit_entry.file.row_count = 2
+        result = file_store_commit._try_commit_once(
+            retry_result=retry_result,
+            commit_kind='APPEND',
+            commit_entries=[commit_entry],
+            changelog_entries=[],
+            commit_identifier=11,
+            latest_snapshot=latest_snapshot,
+        )
+        return file_store_commit, result
+
+    def test_false_atomic_commit_retains_manifest_merge_result(
+            self, mock_manifest_list_manager, mock_manifest_file_manager):
+        file_store_commit, result = self._run_manifest_commit_attempt(
+            commit_result=False)
+
+        self.assertIsInstance(result, RetryResult)
+        self.assertIsNone(result.exception)
+        self.assertEqual(
+            ['before'],
+            [manifest.file_name for manifest
+             in result.manifest_merge_result.merge_before_manifests],
+        )
+        self.assertEqual(
+            ['after'],
+            [manifest.file_name for manifest
+             in result.manifest_merge_result.merge_after_manifests],
+        )
+        file_store_commit.manifest_file_merger.merge.assert_called_once()
+
+    def test_atomic_commit_exception_does_not_retain_manifest_merge_result(
+            self, mock_manifest_list_manager, mock_manifest_file_manager):
+        failure = TimeoutError('lost commit response')
+        file_store_commit, result = self._run_manifest_commit_attempt(
+            commit_side_effect=failure)
+
+        self.assertIsInstance(result, RetryResult)
+        self.assertIs(failure, result.exception)
+        self.assertTrue(result.commit_result_may_be_uncertain)
+        self.assertIsNone(result.manifest_merge_result)
+        file_store_commit._clean_up_reuse_tmp_manifests.assert_not_called()
+        file_store_commit._clean_up_no_reuse_tmp_manifests.assert_not_called()
+
+    def test_retry_reuses_manifest_merge_and_preserves_surrounding_files(
+            self, mock_manifest_list_manager, mock_manifest_file_manager):
+        previous_before = [
+            self._manifest_meta('before-a'),
+            self._manifest_meta('before-b'),
+        ]
+        previous_after = [self._manifest_meta('merged')]
+        retry_result = RetryResult(
+            Mock(id=3),
+            manifest_merge_result=ManifestMergeResult(
+                previous_before, previous_after),
+        )
+        current = [
+            self._manifest_meta('prefix'),
+            self._manifest_meta('before-a'),
+            self._manifest_meta('before-b'),
+            self._manifest_meta('suffix'),
+        ]
+
+        file_store_commit, result = self._run_manifest_commit_attempt(
+            commit_result=False,
+            retry_result=retry_result,
+            existing_manifests=current,
+        )
+
+        self.assertIsInstance(result, RetryResult)
+        self.assertEqual(
+            ['prefix', 'before-a', 'before-b', 'suffix'],
+            [manifest.file_name for manifest
+             in result.manifest_merge_result.merge_before_manifests],
+        )
+        self.assertEqual(
+            ['prefix', 'merged', 'suffix'],
+            [manifest.file_name for manifest
+             in result.manifest_merge_result.merge_after_manifests],
+        )
+        file_store_commit.manifest_file_merger.merge.assert_not_called()
+        base_manifests = (
+            file_store_commit.manifest_list_manager.write
+            .call_args_list[-1].args[1])
+        self.assertEqual(
+            ['prefix', 'merged', 'suffix'],
+            [manifest.file_name for manifest in base_manifests],
+        )
+
+    def test_retry_skips_manifest_merge_when_previous_input_is_not_contiguous(
+            self, mock_manifest_list_manager, mock_manifest_file_manager):
+        previous_before = [
+            self._manifest_meta('before-a'),
+            self._manifest_meta('before-b'),
+        ]
+        retry_result = RetryResult(
+            Mock(id=3),
+            manifest_merge_result=ManifestMergeResult(
+                previous_before, [self._manifest_meta('merged')]),
+        )
+        current = [
+            self._manifest_meta('before-a'),
+            self._manifest_meta('between'),
+            self._manifest_meta('before-b'),
+        ]
+
+        file_store_commit, result = self._run_manifest_commit_attempt(
+            commit_result=False,
+            retry_result=retry_result,
+            existing_manifests=current,
+        )
+
+        self.assertIsInstance(result, RetryResult)
+        self.assertIsNone(result.manifest_merge_result)
+        file_store_commit.manifest_file_merger.merge.assert_not_called()
+        base_manifests = (
+            file_store_commit.manifest_list_manager.write
+            .call_args_list[-1].args[1])
+        self.assertEqual(current, base_manifests)
+
+    def test_manifest_merge_runs_once_across_multiple_retries(
+            self, mock_manifest_list_manager, mock_manifest_file_manager):
+        first_commit, retry_result = self._run_manifest_commit_attempt(
+            commit_result=False,
+            existing_manifests=[self._manifest_meta('before')],
+            merged_manifests=[self._manifest_meta('merged')],
+        )
+        first_commit.manifest_file_merger.merge.assert_called_once()
+
+        unchanged_retry, retry_result = self._run_manifest_commit_attempt(
+            commit_result=False,
+            retry_result=retry_result,
+            existing_manifests=[self._manifest_meta('before')],
+        )
+        retry_commits = [unchanged_retry]
+        current_names = ['before']
+        for suffix in ['concurrent-1', 'concurrent-2']:
+            current_names.append(suffix)
+            retry_commit, retry_result = self._run_manifest_commit_attempt(
+                commit_result=False,
+                retry_result=retry_result,
+                existing_manifests=[
+                    self._manifest_meta(name) for name in current_names
+                ],
+            )
+            retry_commits.append(retry_commit)
+
+        final_names = current_names + ['concurrent-3']
+        final_commit, result = self._run_manifest_commit_attempt(
+            commit_result=True,
+            retry_result=retry_result,
+            existing_manifests=[
+                self._manifest_meta(name) for name in final_names
+            ],
+        )
+
+        self.assertTrue(result.is_success())
+        for retry_commit in retry_commits + [final_commit]:
+            retry_commit.manifest_file_merger.merge.assert_not_called()
+        base_manifests = (
+            final_commit.manifest_list_manager.write
+            .call_args_list[-1].args[1])
+        self.assertEqual(
+            ['merged', 'concurrent-1', 'concurrent-2', 'concurrent-3'],
+            [manifest.file_name for manifest in base_manifests],
+        )
+
     def test_generate_partition_statistics_single_partition_single_file(
             self, mock_manifest_list_manager, mock_manifest_file_manager):
         """Test partition statistics generation with single partition and 
single file."""
diff --git a/paimon-python/pypaimon/write/file_store_commit.py 
b/paimon-python/pypaimon/write/file_store_commit.py
index 405052b8f5..f9012c5eaf 100644
--- a/paimon-python/pypaimon/write/file_store_commit.py
+++ b/paimon-python/pypaimon/write/file_store_commit.py
@@ -34,7 +34,7 @@ from pypaimon.manifest.schema.manifest_file_meta import 
ManifestFileMeta
 from pypaimon.snapshot.snapshot import Snapshot
 from pypaimon.snapshot.snapshot_commit import (PartitionStatistics,
                                                SnapshotCommit)
-from pypaimon.table.row.generic_row import GenericRow
+from pypaimon.table.row.generic_row import GenericRow, GenericRowSerializer
 from pypaimon.table.row.offset_row import OffsetRow
 from pypaimon.write.commit.commit_rollback import CommitRollback
 from pypaimon.write.commit.commit_scanner import CommitScanner
@@ -70,17 +70,76 @@ class SuccessResult(CommitResult):
         return True
 
 
+def _manifest_file_key(manifest: ManifestFileMeta):
+    stats = manifest.partition_stats
+    return (
+        manifest.file_name,
+        manifest.file_size,
+        manifest.num_added_files,
+        manifest.num_deleted_files,
+        GenericRowSerializer.to_bytes(stats.min_values),
+        GenericRowSerializer.to_bytes(stats.max_values),
+        tuple(stats.null_counts) if stats.null_counts is not None else None,
+        manifest.schema_id,
+        manifest.min_row_id,
+        manifest.max_row_id,
+    )
+
+
+def _try_replace_manifest_files(current, replaced, replacement):
+    """Replace the first contiguous occurrence while preserving list order."""
+    current = list(current)
+    replaced = list(replaced)
+    replacement = list(replacement)
+    if not replaced:
+        return replacement if not current else None
+
+    current_keys = [_manifest_file_key(manifest) for manifest in current]
+    replaced_keys = [_manifest_file_key(manifest) for manifest in replaced]
+    for start in range(len(current) - len(replaced) + 1):
+        if current_keys[start:start + len(replaced)] == replaced_keys:
+            return (
+                current[:start]
+                + replacement
+                + current[start + len(replaced):]
+            )
+    return None
+
+
+class ManifestMergeResult:
+    """Manifest merge input and output retained for a deterministic retry."""
+
+    def __init__(self, merge_before_manifests, merge_after_manifests):
+        self.merge_before_manifests = tuple(merge_before_manifests)
+        self.merge_after_manifests = tuple(merge_after_manifests)
+
+
+def _try_reuse_manifest_merge_result(retry_result, current_manifests):
+    if (retry_result is None
+            or retry_result.commit_result_may_be_uncertain
+            or retry_result.manifest_merge_result is None):
+        return None
+    previous = retry_result.manifest_merge_result
+    return _try_replace_manifest_files(
+        current_manifests,
+        previous.merge_before_manifests,
+        previous.merge_after_manifests,
+    )
+
+
 class RetryResult(CommitResult):
 
     def __init__(self, latest_snapshot, exception: Optional[Exception] = None,
                  base_data_files: Optional[List[ManifestEntry]] = None,
-                 commit_result_may_be_uncertain: bool = False):
+                 commit_result_may_be_uncertain: bool = False,
+                 manifest_merge_result: Optional[ManifestMergeResult] = None):
         self.latest_snapshot = latest_snapshot
         self.exception = exception
         self.commit_result_may_be_uncertain = commit_result_may_be_uncertain
         # Base entries as of latest_snapshot, carried so the next attempt 
reuses
         # them and reads only the incremental changes.
         self.base_data_files = base_data_files
+        self.manifest_merge_result = manifest_merge_result
 
     def is_success(self) -> bool:
         return False
@@ -597,7 +656,10 @@ class FileStoreCommit:
         changelog_manifest_list_name = None
         changelog_manifest_list_size = None
         changelog_record_count = None
+        merge_before_manifests = []
+        merge_after_manifests = []
         merge_new_files = []
+        skip_manifest_merge_on_retry = False
         try:
             new_manifest_file_metas = 
self._write_manifest_files(commit_entries, new_manifest_file)
             self.manifest_list_manager.write(delta_manifest_list, 
new_manifest_file_metas)
@@ -620,15 +682,33 @@ class FileStoreCommit:
             # process existing_manifest
             total_record_count = 0
             if latest_snapshot:
-                existing_manifest_files = 
self.manifest_list_manager.read_all(latest_snapshot)
+                merge_before_manifests = self.manifest_list_manager.read_all(
+                    latest_snapshot)
                 previous_record_count = latest_snapshot.total_record_count
                 if previous_record_count:
                     total_record_count += previous_record_count
+
+            reused_manifests = _try_reuse_manifest_merge_result(
+                retry_result, merge_before_manifests)
+            skip_manifest_merge_on_retry = (
+                reused_manifests is None and retry_result is not None)
+            if reused_manifests is not None:
+                merge_after_manifests = reused_manifests
+                old_names = {
+                    manifest.file_name for manifest in merge_before_manifests
+                }
+                merge_new_files = [
+                    manifest for manifest in merge_after_manifests
+                    if manifest.file_name not in old_names
+                ]
+            elif skip_manifest_merge_on_retry:
+                merge_after_manifests = merge_before_manifests
             else:
-                existing_manifest_files = []
-            merged_manifest_files, merge_new_files = 
self.manifest_file_merger.merge(
-                existing_manifest_files)
-            self.manifest_list_manager.write(base_manifest_list, 
merged_manifest_files)
+                merge_after_manifests, merge_new_files = (
+                    self.manifest_file_merger.merge(
+                        merge_before_manifests))
+            self.manifest_list_manager.write(
+                base_manifest_list, merge_after_manifests)
 
             delta_record_count = 0
             for entry in commit_entries:
@@ -698,7 +778,20 @@ class FileStoreCommit:
                         commit_kind,
                         commit_time_s,
                     )
-                    return RetryResult(latest_snapshot, None, 
base_data_files=base_data_files)
+                    manifest_merge_result = (
+                        None
+                        if skip_manifest_merge_on_retry
+                        else ManifestMergeResult(
+                            merge_before_manifests,
+                            merge_after_manifests,
+                        )
+                    )
+                    return RetryResult(
+                        latest_snapshot,
+                        None,
+                        base_data_files=base_data_files,
+                        manifest_merge_result=manifest_merge_result,
+                    )
         except Exception as e:
             # Commit exception, not sure about the situation and should not 
clean up the files
             logger.warning("Retry commit for exception.", exc_info=True)
@@ -707,6 +800,7 @@ class FileStoreCommit:
                 e,
                 base_data_files=base_data_files,
                 commit_result_may_be_uncertain=True,
+                manifest_merge_result=None,
             )
 
         logger.info(

Reply via email to