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 a9259cf217 [python] Group row tracking commits by partition (#9312)
a9259cf217 is described below

commit a9259cf2174fb36b2aa0171ea875ed1a7082ce91
Author: Yann Byron <[email protected]>
AuthorDate: Thu Aug 20 15:40:30 2026 +0800

    [python] Group row tracking commits by partition (#9312)
---
 .../pypaimon/common/options/core_options.py        |  15 ++
 .../pypaimon/tests/file_store_commit_test.py       | 200 +++++++++++++++++++++
 paimon-python/pypaimon/write/file_store_commit.py  |  21 ++-
 3 files changed, 235 insertions(+), 1 deletion(-)

diff --git a/paimon-python/pypaimon/common/options/core_options.py 
b/paimon-python/pypaimon/common/options/core_options.py
index 7cde748d64..e76db78870 100644
--- a/paimon-python/pypaimon/common/options/core_options.py
+++ b/paimon-python/pypaimon/common/options/core_options.py
@@ -703,6 +703,17 @@ class CoreOptions:
         .with_description("Whether to enable row tracking.")
     )
 
+    ROW_TRACKING_PARTITION_GROUP_ON_COMMIT: ConfigOption[bool] = (
+        ConfigOptions.key("row-tracking.partition-group-on-commit")
+        .boolean_type()
+        .default_value(True)
+        .with_description(
+            "When row-tracking is enabled, whether to group new file metas "
+            "by partition before commit, so that assigned row IDs are "
+            "contiguous within each partition."
+        )
+    )
+
     DATA_EVOLUTION_ENABLED: ConfigOption[bool] = (
         ConfigOptions.key("data-evolution.enabled")
         .boolean_type()
@@ -1355,6 +1366,10 @@ class CoreOptions:
     def row_tracking_enabled(self, default=None):
         return self.options.get(CoreOptions.ROW_TRACKING_ENABLED, default)
 
+    def row_tracking_partition_group_on_commit(self, default=None):
+        return self.options.get(
+            CoreOptions.ROW_TRACKING_PARTITION_GROUP_ON_COMMIT, default)
+
     def data_evolution_enabled(self, default=None):
         return self.options.get(CoreOptions.DATA_EVOLUTION_ENABLED, default)
 
diff --git a/paimon-python/pypaimon/tests/file_store_commit_test.py 
b/paimon-python/pypaimon/tests/file_store_commit_test.py
index 1bb2852b78..3576d659cb 100644
--- a/paimon-python/pypaimon/tests/file_store_commit_test.py
+++ b/paimon-python/pypaimon/tests/file_store_commit_test.py
@@ -20,6 +20,8 @@ import uuid
 from datetime import datetime
 from unittest.mock import MagicMock, Mock, patch
 
+from pypaimon.common.options.core_options import CoreOptions
+from pypaimon.common.options.options import Options
 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
@@ -56,6 +58,204 @@ class TestAbortCommitMessages(unittest.TestCase):
             _abort_commit_messages(table, [message])
 
 
+class TestFileStoreCommitRowTracking(unittest.TestCase):
+
+    def setUp(self):
+        self.mock_table = Mock()
+        self.mock_table.partition_keys = ['dt', 'region']
+        self.mock_table.current_branch.return_value = 'main'
+        self.mock_table.table_path = '/test/table/path'
+        self.mock_table.file_io = Mock()
+        self.mock_table.options.manifest_target_size.return_value = 8 * 1024 * 
1024
+        self.mock_table.options.manifest_merge_min_count.return_value = 30
+        self.mock_snapshot_commit = Mock()
+
+    def _create_file_store_commit(self):
+        with patch('pypaimon.write.file_store_commit.ManifestFileManager'), \
+                patch('pypaimon.write.file_store_commit.ManifestListManager'):
+            return FileStoreCommit(
+                snapshot_commit=self.mock_snapshot_commit,
+                table=self.mock_table,
+                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,
+        )
+
+    @staticmethod
+    def _data_file(name, row_count):
+        return DataFileMeta.create(
+            file_name=name,
+            file_size=10,
+            row_count=row_count,
+            min_key=GenericRow([], []),
+            max_key=GenericRow([], []),
+            key_stats=SimpleStats.empty_stats(),
+            value_stats=SimpleStats.empty_stats(),
+            min_sequence_number=0,
+            max_sequence_number=0,
+            schema_id=0,
+            level=0,
+            extra_files=[],
+            file_source=0,
+        )
+
+    @classmethod
+    def _append_entry(cls, partition, name, row_count):
+        return ManifestEntry(
+            kind=0,
+            partition=GenericRow(list(partition), None),
+            bucket=0,
+            total_buckets=1,
+            file=cls._data_file(name, row_count),
+        )
+
+    def test_groups_files_by_partition_before_assigning_row_ids(self):
+        file_store_commit = self._create_file_store_commit()
+        entries = [
+            self._append_entry(('dt=1',), 'a-1.parquet', 2),
+            self._append_entry(('dt=2',), 'b-1.parquet', 3),
+            self._append_entry(('dt=1',), 'a-2.parquet', 4),
+        ]
+
+        grouped = file_store_commit._group_commit_entries_by_partition(entries)
+        assigned, next_row_id = file_store_commit._assign_row_tracking_meta(
+            10, grouped)
+
+        self.assertEqual(19, next_row_id)
+        self.assertEqual(
+            [
+                ('a-1.parquet', 10),
+                ('a-2.parquet', 12),
+                ('b-1.parquet', 16),
+            ],
+            [(entry.file.file_name, entry.file.first_row_id)
+             for entry in assigned],
+        )
+
+    def test_partition_group_on_commit_option_defaults_to_true(self):
+        options = CoreOptions(Options({}))
+
+        self.assertTrue(options.row_tracking_partition_group_on_commit())
+        self.assertFalse(
+            CoreOptions(Options({
+                "row-tracking.partition-group-on-commit": "false",
+            })).row_tracking_partition_group_on_commit())
+
+    def test_keeps_order_when_partition_grouping_disabled(self):
+        file_store_commit = self._create_file_store_commit()
+        entries = [
+            self._append_entry(('dt=1',), 'a-1.parquet', 2),
+            self._append_entry(('dt=2',), 'b-1.parquet', 3),
+            self._append_entry(('dt=1',), 'a-2.parquet', 4),
+        ]
+
+        assigned, next_row_id = file_store_commit._assign_row_tracking_meta(
+            10, entries)
+
+        self.assertEqual(19, next_row_id)
+        self.assertEqual(
+            [
+                ('a-1.parquet', 10),
+                ('b-1.parquet', 12),
+                ('a-2.parquet', 15),
+            ],
+            [(entry.file.file_name, entry.file.first_row_id)
+             for entry in assigned],
+        )
+
+    def _commit_row_tracking_entries(self, options):
+        file_store_commit = self._create_file_store_commit()
+        self.mock_table.identifier = 'default.test_table'
+        self.mock_table.table_schema = Mock()
+        self.mock_table.table_schema.id = 7
+        self.mock_table.options = CoreOptions(Options(options))
+
+        snapshot_commit = MagicMock()
+        snapshot_commit.__enter__.return_value = snapshot_commit
+        snapshot_commit.__exit__.return_value = False
+        snapshot_commit.commit.return_value = True
+        file_store_commit.snapshot_commit = snapshot_commit
+        file_store_commit.manifest_list_manager.read_all.return_value = []
+        file_store_commit.manifest_file_merger = Mock()
+        file_store_commit.manifest_file_merger.merge.return_value = ([], [])
+        file_store_commit._generate_partition_statistics = Mock(
+            return_value=[])
+
+        written_entries = []
+
+        def capture_entries(entries, manifest_file_name):
+            del manifest_file_name
+            written_entries.extend(entries)
+            return [self._manifest_meta('delta')]
+
+        file_store_commit._write_manifest_files = Mock(
+            side_effect=capture_entries)
+
+        entries = [
+            self._append_entry(('dt=1',), 'a-1.parquet', 2),
+            self._append_entry(('dt=2',), 'b-1.parquet', 3),
+            self._append_entry(('dt=1',), 'a-2.parquet', 4),
+        ]
+
+        file_store_commit._try_commit_once(
+            retry_result=None,
+            commit_kind="APPEND",
+            commit_entries=entries,
+            changelog_entries=[],
+            commit_identifier=11,
+            latest_snapshot=None,
+        )
+
+        committed_snapshot = snapshot_commit.commit.call_args[0][1]
+        return committed_snapshot, written_entries
+
+    def test_commit_applies_default_partition_grouping_before_row_id(self):
+        committed_snapshot, written_entries = (
+            self._commit_row_tracking_entries({
+                "row-tracking.enabled": "true",
+            }))
+
+        self.assertEqual(9, committed_snapshot.next_row_id)
+        self.assertEqual(
+            [
+                ('a-1.parquet', 0),
+                ('a-2.parquet', 2),
+                ('b-1.parquet', 6),
+            ],
+            [(entry.file.file_name, entry.file.first_row_id)
+             for entry in written_entries],
+        )
+
+    def test_commit_respects_disabled_partition_grouping(self):
+        committed_snapshot, written_entries = (
+            self._commit_row_tracking_entries({
+                "row-tracking.enabled": "true",
+                "row-tracking.partition-group-on-commit": "false",
+            }))
+
+        self.assertEqual(9, committed_snapshot.next_row_id)
+        self.assertEqual(
+            [
+                ('a-1.parquet', 0),
+                ('b-1.parquet', 2),
+                ('a-2.parquet', 5),
+            ],
+            [(entry.file.file_name, entry.file.first_row_id)
+             for entry in written_entries],
+        )
+
+
 @patch('pypaimon.write.file_store_commit.ManifestFileManager')
 @patch('pypaimon.write.file_store_commit.ManifestListManager')
 class TestFileStoreCommit(unittest.TestCase):
diff --git a/paimon-python/pypaimon/write/file_store_commit.py 
b/paimon-python/pypaimon/write/file_store_commit.py
index 2f64c59365..7bb6db4dad 100644
--- a/paimon-python/pypaimon/write/file_store_commit.py
+++ b/paimon-python/pypaimon/write/file_store_commit.py
@@ -673,8 +673,14 @@ class FileStoreCommit:
         next_row_id = None
         if row_tracking_enabled:
             commit_entries = self._assign_snapshot_id(new_snapshot_id, 
commit_entries)
+            group_by_partition = (
+                self.table.options.row_tracking_partition_group_on_commit())
+            if group_by_partition:
+                commit_entries = self._group_commit_entries_by_partition(
+                    commit_entries)
             first_row_id_start = self._get_next_row_id_start(latest_snapshot)
-            commit_entries, next_row_id = 
self._assign_row_tracking_meta(first_row_id_start, commit_entries)
+            commit_entries, next_row_id = self._assign_row_tracking_meta(
+                first_row_id_start, commit_entries)
 
         changelog_manifest_list_name = None
         changelog_manifest_list_size = None
@@ -1184,6 +1190,19 @@ class FileStoreCommit:
             return latest_snapshot.next_row_id
         return 0
 
+    @staticmethod
+    def _group_commit_entries_by_partition(
+            commit_entries: List[ManifestEntry]) -> List[ManifestEntry]:
+        grouped = {}
+        for entry in commit_entries:
+            key = tuple(entry.partition.values)
+            grouped.setdefault(key, []).append(entry)
+        return [
+            entry
+            for entries in grouped.values()
+            for entry in entries
+        ]
+
     def _assign_row_tracking_meta(self, first_row_id_start: int, 
commit_entries: List[ManifestEntry]):
         """Assign row tracking metadata (first_row_id) to new files.
 

Reply via email to