From 8bc274be9b34d3d9f240219844183b4a1b99c921 Mon Sep 17 00:00:00 2001 From: JingsongLi Date: Sun, 20 Sep 2026 22:09:06 +0800 Subject: [PATCH 1/5] [python] Plan chunk-shuffled reads natively --- paimon-python/pypaimon/read/native_plan.py | 56 ++++++++++++++++++- paimon-python/pypaimon/read/table_scan.py | 41 +++----------- .../tests/native_plan_chunk_shuffle_test.py | 25 +++++++-- .../pypaimon/tests/native_plan_test.py | 25 +++++++++ 4 files changed, 108 insertions(+), 39 deletions(-) diff --git a/paimon-python/pypaimon/read/native_plan.py b/paimon-python/pypaimon/read/native_plan.py index f1e4d0a87068..f174e3530eca 100644 --- a/paimon-python/pypaimon/read/native_plan.py +++ b/paimon-python/pypaimon/read/native_plan.py @@ -371,7 +371,9 @@ def native_plan( row_ranges: Optional[List[Tuple[int, int]]] = None, incremental_range: Optional[Tuple[int, int]] = None, row_position_slice: Optional[Tuple[int, int]] = None, - row_position_shard: Optional[Tuple[int, int]] = None) -> Plan: + row_position_shard: Optional[Tuple[int, int]] = None, + chunk_shuffle: Optional[Tuple[int, int]] = None, + chunk_shuffle_shard: Optional[Tuple[int, int]] = None) -> Plan: """Plan with pypaimon_rust, preserving snapshot metadata. Native conversion or planning failures are handled by TableScan, which @@ -390,12 +392,23 @@ def native_plan( scan = scan.with_row_position_slice(*row_position_slice) if row_position_shard is not None: scan = scan.with_row_position_shard(*row_position_shard) + if chunk_shuffle is not None: + seed, chunk_size = chunk_shuffle + shard_index, shard_count = ( + chunk_shuffle_shard if chunk_shuffle_shard is not None + else (None, None) + ) + scan = scan.with_chunk_shuffle( + str(seed), chunk_size, shard_index, shard_count) rust_plan = scan.plan() rust_splits = rust_plan.splits() pfields = _partition_fields(table) # Trimmed primary keys decode per-file min/max keys (PK merge-on-read). kfields = table.trimmed_primary_keys_fields - splits = [deserialize_split_v1(split.serialize(), pfields, kfields) for split in rust_splits] + splits = [ + _native_split_metadata_view(split, pfields, kfields) + for split in rust_splits + ] if table.options.native_read_enabled(): # Retain the opaque Rust split next to the Python metadata view. The # normal planner/reader contract remains a Python Split list, while @@ -415,3 +428,42 @@ def native_plan( # table without snapshots. Let the Python scanner recover the metadata. raise RuntimeError("Native runtime cannot report an empty plan's snapshot") return Plan(splits, snapshot_id=snapshot_id) + + +def _native_split_metadata_view(rust_split, partition_fields, key_fields): + """Build the Python metadata facade while retaining native-only ranges. + + File-local ranges cannot be represented by Java SplitSerializer v1. They + remain authoritative on ``rust_split``; SlicedSplit is only the public + Python view used for row counts, paths and diagnostics. + """ + serialize_metadata = (getattr(rust_split, 'serialize_metadata', None) + if hasattr(type(rust_split), 'serialize_metadata') + else None) + payload = (serialize_metadata() if callable(serialize_metadata) + else rust_split.serialize()) + split = deserialize_split_v1(payload, partition_fields, key_fields) + + exact_count_fn = (getattr(rust_split, 'exact_merged_row_count', None) + if hasattr(type(rust_split), 'exact_merged_row_count') + else None) + exact_count = exact_count_fn() if callable(exact_count_fn) else None + file_ranges_fn = (getattr(rust_split, 'file_row_ranges', None) + if hasattr(type(rust_split), 'file_row_ranges') + else None) + file_ranges = file_ranges_fn() if callable(file_ranges_fn) else None + if file_ranges is not None: + from pypaimon.read.sliced_split import SlicedSplit + split = SlicedSplit( + split, file_ranges, exact_merged_row_count=exact_count) + else: + from pypaimon.globalindex.indexed_split import IndexedSplit + if isinstance(split, IndexedSplit) and exact_count is not None: + split = IndexedSplit( + split.data_split(), split.row_ranges(), split.scores(), + exact_merged_row_count=exact_count) + elif exact_count is not None and split.merged_row_count() != exact_count: + from pypaimon.read.sliced_split import SlicedSplit + split = SlicedSplit( + split, {}, exact_merged_row_count=exact_count) + return split diff --git a/paimon-python/pypaimon/read/table_scan.py b/paimon-python/pypaimon/read/table_scan.py index 8a6536bc0a7a..7d7e6d074fd4 100755 --- a/paimon-python/pypaimon/read/table_scan.py +++ b/paimon-python/pypaimon/read/table_scan.py @@ -144,6 +144,8 @@ def _native_plan_supported_impl(self) -> bool: return False if getattr(fs, 'chunk_shuffle', None) is not None: fs._validate_chunk_shuffle_compat() + if not native_method_available('TableScan', 'with_chunk_shuffle'): + return False # Positional append distribution needs the stable partition/file order # introduced in 0.4. Older bindings can assign different rows per call. if (not self.table.is_primary_key_table and not fs.data_evolution @@ -272,6 +274,11 @@ def _try_native_plan(self) -> Optional[Plan]: return Plan([]) extra_options['incremental_range'] = self._incremental_snapshot_range chunk_shuffle = fs.chunk_shuffle + if chunk_shuffle is not None: + extra_options['chunk_shuffle'] = chunk_shuffle + if fs.idx_of_this_subtask is not None: + extra_options['chunk_shuffle_shard'] = ( + fs.idx_of_this_subtask, fs.number_of_para_subtasks) if has_distribution and fs.data_evolution and chunk_shuffle is None: if fs.idx_of_this_subtask is not None: extra_options['row_position_shard'] = ( @@ -318,9 +325,7 @@ def _try_native_plan(self) -> Optional[Plan]: and plan.snapshot_id is not None): snapshot = self.table.snapshot_manager().get_snapshot_by_id(plan.snapshot_id) splits = fs._apply_primary_key_sorted_indexes(splits, snapshot) - if chunk_shuffle is not None: - splits = self._chunk_shuffle_splits(splits, plan.snapshot_id) - elif has_distribution: + if chunk_shuffle is None and has_distribution: if self.table.is_primary_key_table: splits = [s for s in splits if s.bucket % fs.number_of_para_subtasks == fs.idx_of_this_subtask] @@ -350,36 +355,6 @@ def _try_native_plan(self) -> Optional[Plan]: "Native plan failed, falling back to the Python scanner: %s", e) return None - def _chunk_shuffle_splits(self, splits, snapshot_id): - """Reuse native file/DV planning with Python's stable chunk assignment.""" - from pypaimon.manifest.schema.manifest_entry import ManifestEntry - from pypaimon.read.scanner.chunk_shuffle_split_generator import ( - AppendChunkShuffleSplitGenerator, DataEvolutionChunkShuffleSplitGenerator, - ) - fs = self.file_scanner - entries, deletions = [], {} - for split in splits: - key = (tuple(split.partition.values), split.bucket) - for index, file in enumerate(split.files): - entries.append(ManifestEntry( - 0, split.partition, split.bucket, self.table.total_buckets, file)) - if split.data_deletion_files and split.data_deletion_files[index] is not None: - deletions.setdefault(key, {})[file.file_name] = split.data_deletion_files[index] - generator_type = (DataEvolutionChunkShuffleSplitGenerator if fs.data_evolution - else AppendChunkShuffleSplitGenerator) - seed, chunk_size = fs.chunk_shuffle - generator = generator_type(self.table, fs.target_split_size, fs.open_file_cost, - deletions, seed=seed, chunk_size=chunk_size, - snapshot_id=snapshot_id) - if fs.idx_of_this_subtask is not None: - generator.with_shard(fs.idx_of_this_subtask, fs.number_of_para_subtasks) - chunks = generator.create_splits(entries) - for split in chunks: - while callable(getattr(split, 'data_split', None)): - split = split.data_split() - split.is_streaming = fs.is_streaming - return chunks - def plan_for_write(self) -> Plan: if self.__auth_query() is not None: raise TableNoPermissionException(self.table.identifier) diff --git a/paimon-python/pypaimon/tests/native_plan_chunk_shuffle_test.py b/paimon-python/pypaimon/tests/native_plan_chunk_shuffle_test.py index ea156d8dccb0..2725458aef88 100644 --- a/paimon-python/pypaimon/tests/native_plan_chunk_shuffle_test.py +++ b/paimon-python/pypaimon/tests/native_plan_chunk_shuffle_test.py @@ -24,13 +24,19 @@ from pypaimon import CatalogFactory, Schema from pypaimon.deletionvectors.bitmap_deletion_vector import BitmapDeletionVector -from pypaimon.read.native_plan import native_runtime_available +from pypaimon.read.native_plan import ( + native_method_available, + native_reader_available, + native_runtime_available, +) from pypaimon.write.commit_message import CommitMessage from pypaimon.write.table_delete import TableDeleteByRowId pytestmark = [pytest.mark.native_plan, pytest.mark.skipif( - not native_runtime_available(), reason='Rust planner required')] + not native_runtime_available() + or not native_method_available('TableScan', 'with_chunk_shuffle'), + reason='Rust native chunk planner required')] @pytest.fixture(params=['append', 'append-dv', 'de', 'de-dv']) @@ -105,7 +111,10 @@ def chunk_table(request, tmp_path): def _chunks(table, seed, chunk_size=3, shard=None, predicate=None, projection=None): results = [] for native in (False, True): - builder = table.copy({'scan.native-plan.enabled': str(native).lower()}).new_read_builder() + builder = table.copy({ + 'scan.native-plan.enabled': str(native).lower(), + 'read.native.enabled': str(native).lower(), + }).new_read_builder() if predicate is not None: builder.with_filter(predicate) if projection is not None: @@ -125,7 +134,15 @@ def _chunks(table, seed, chunk_size=3, shard=None, predicate=None, projection=No for split in plan.splits(): if table.options.options.contains_key('incremental-between-timestamp'): assert split.is_streaming - rows = builder.new_read().to_arrow([split]).to_pylist() + if native: + assert native_reader_available() + assert getattr(split, '_native_split', None) is not None + with patch( + 'pypaimon.read.table_read.TableRead._create_split_read', + side_effect=AssertionError('Python reader was used')): + rows = builder.new_read().to_arrow([split]).to_pylist() + else: + rows = builder.new_read().to_arrow([split]).to_pylist() assert 0 < len(rows) <= chunk_size assert split.merged_row_count() == len(rows) if rows and 'p' in rows[0]: diff --git a/paimon-python/pypaimon/tests/native_plan_test.py b/paimon-python/pypaimon/tests/native_plan_test.py index 8ea40f6b2bb6..a61023e67f19 100644 --- a/paimon-python/pypaimon/tests/native_plan_test.py +++ b/paimon-python/pypaimon/tests/native_plan_test.py @@ -33,6 +33,7 @@ from pypaimon.globalindex.vector_search_result import ScoredGlobalIndexResult from pypaimon.read.native_plan import ( _catalog_options, + _native_split_metadata_view, _predicate_to_native, _read_options, _resolved_schema_json, @@ -141,6 +142,30 @@ def test_resolved_schema_json_uses_canonical_nested_collection_types(self): 'key': 'STRING NOT NULL', 'value': 'INT'}}}]}}}], }) + def test_native_split_metadata_view_preserves_file_ranges_and_exact_count(self): + class NativeSplit: + def serialize_metadata(self): + return b'metadata' + + def file_row_ranges(self): + return {'f.parquet': (2, 5)} + + def exact_merged_row_count(self): + return 3 + + base = Mock() + with patch( + 'pypaimon.read.native_plan.deserialize_split_v1', + return_value=base) as deserialize: + view = _native_split_metadata_view(NativeSplit(), [], []) + + from pypaimon.read.sliced_split import SlicedSplit + self.assertIsInstance(view, SlicedSplit) + self.assertIs(view.data_split(), base) + self.assertEqual(view.shard_file_idx_map(), {'f.parquet': (2, 5)}) + self.assertEqual(view.merged_row_count(), 3) + deserialize.assert_called_once_with(b'metadata', [], []) + def test_resolved_schema_keeps_custom_io_and_rest_on_catalog_path(self): from pypaimon.catalog.catalog_environment import CatalogEnvironment from pypaimon.catalog.jdbc_catalog_loader import JdbcCatalogLoader From dd0af14208068b546202174c7d3e2cd0e263d0e2 Mon Sep 17 00:00:00 2001 From: JingsongLi Date: Sun, 20 Sep 2026 22:48:53 +0800 Subject: [PATCH 2/5] [python] Simplify native chunk shuffle API --- paimon-python/pypaimon/read/native_plan.py | 5 +++-- 1 file changed, 3 insertions(+), 2 deletions(-) diff --git a/paimon-python/pypaimon/read/native_plan.py b/paimon-python/pypaimon/read/native_plan.py index f174e3530eca..b6ffdda29b89 100644 --- a/paimon-python/pypaimon/read/native_plan.py +++ b/paimon-python/pypaimon/read/native_plan.py @@ -398,8 +398,9 @@ def native_plan( chunk_shuffle_shard if chunk_shuffle_shard is not None else (None, None) ) - scan = scan.with_chunk_shuffle( - str(seed), chunk_size, shard_index, shard_count) + scan = scan.with_chunk_shuffle(str(seed), chunk_size) + if shard_index is not None: + scan = scan.with_chunk_shuffle_shard(shard_index, shard_count) rust_plan = scan.plan() rust_splits = rust_plan.splits() pfields = _partition_fields(table) From 3e1a3912cde6dbebb6d9ee94fc9568b5bbc6793d Mon Sep 17 00:00:00 2001 From: JingsongLi Date: Mon, 21 Sep 2026 00:00:21 +0800 Subject: [PATCH 3/5] [python] Reuse IndexedSplit ranges for chunk shuffle --- .../pypaimon/globalindex/indexed_split.py | 17 +-- paimon-python/pypaimon/multimodal/temporal.py | 6 +- .../pypaimon/read/datasource/torch_dataset.py | 6 +- paimon-python/pypaimon/read/native_plan.py | 51 +------ .../scanner/chunk_shuffle_split_generator.py | 134 ++++++++---------- paimon-python/pypaimon/read/sliced_split.py | 11 +- paimon-python/pypaimon/read/split_read.py | 96 ++++++++++--- paimon-python/pypaimon/read/table_scan.py | 2 +- .../data_evolution_deletion_vector_test.py | 53 ++++++- .../tests/native_plan_distribution_test.py | 3 - .../pypaimon/tests/native_plan_test.py | 25 ---- .../chunk_shuffle_split_generator_test.py | 89 +++++++----- 12 files changed, 258 insertions(+), 235 deletions(-) diff --git a/paimon-python/pypaimon/globalindex/indexed_split.py b/paimon-python/pypaimon/globalindex/indexed_split.py index e2f203578d2e..fa7aecf84ce6 100644 --- a/paimon-python/pypaimon/globalindex/indexed_split.py +++ b/paimon-python/pypaimon/globalindex/indexed_split.py @@ -15,8 +15,10 @@ # specific language governing permissions and limitations # under the License. -""" -IndexedSplit wraps a Split with row ranges and optional scores. +"""IndexedSplit wraps a Split with row ranges and optional scores. + +Ranges use the coordinate system of the read path: stable row IDs for data +evolution and physical positions for primary-key/raw append reads. """ from typing import List, Optional @@ -43,19 +45,17 @@ def __init__( data_split: 'Split', row_ranges: List['Range'], scores: Optional[List[float]] = None, - exact_merged_row_count: Optional[int] = None, ): self._data_split = data_split self._row_ranges = row_ranges self._scores = scores - self._exact_merged_row_count = exact_merged_row_count def data_split(self) -> 'Split': """Return the underlying data split.""" return self._data_split def row_ranges(self) -> List['Range']: - """Return the row ranges from global index.""" + """Return ranges in the coordinate system of the read path.""" return self._row_ranges def scores(self) -> Optional[List[float]]: @@ -90,8 +90,6 @@ def row_count(self) -> int: return sum(r.count() for r in self._row_ranges) def merged_row_count(self): - if self._exact_merged_row_count is not None: - return self._exact_merged_row_count return self.row_count # Delegate other properties to data_split @@ -157,7 +155,6 @@ def __eq__(self, other): self._data_split == other._data_split and self._row_ranges == other._row_ranges and self._scores == other._scores - and self._exact_merged_row_count == other._exact_merged_row_count ) def __hash__(self): @@ -166,10 +163,8 @@ def __hash__(self): id(self._data_split), tuple(self._row_ranges), scores_hash, - self._exact_merged_row_count, )) def __repr__(self): return (f"IndexedSplit(data_split={self._data_split}, " - f"row_ranges={self._row_ranges}, scores={self._scores}, " - f"exact_merged_row_count={self._exact_merged_row_count})") + f"row_ranges={self._row_ranges}, scores={self._scores})") diff --git a/paimon-python/pypaimon/multimodal/temporal.py b/paimon-python/pypaimon/multimodal/temporal.py index 50af264bfb5b..9f8a2d514195 100644 --- a/paimon-python/pypaimon/multimodal/temporal.py +++ b/paimon-python/pypaimon/multimodal/temporal.py @@ -1228,11 +1228,7 @@ def fetch(self, row_ids): allowed = Range.and_(wanted, self._split_ranges[split_index]) if not allowed: continue - indexed = IndexedSplit( - split, - allowed, - exact_merged_row_count=sum(r.count() for r in allowed), - ) + indexed = IndexedSplit(split, allowed) if auth_result is not None: indexed = QueryAuthSplit(indexed, auth_result) selected_splits.append(indexed) diff --git a/paimon-python/pypaimon/read/datasource/torch_dataset.py b/paimon-python/pypaimon/read/datasource/torch_dataset.py index 45c1286fe48c..541dd1d987f0 100644 --- a/paimon-python/pypaimon/read/datasource/torch_dataset.py +++ b/paimon-python/pypaimon/read/datasource/torch_dataset.py @@ -218,11 +218,7 @@ def select_indexed_splits( if not allowed: continue - indexed = IndexedSplit( - split, - allowed, - exact_merged_row_count=sum(r.count() for r in allowed), - ) + indexed = IndexedSplit(split, allowed) selected.append( QueryAuthSplit(indexed, auth_result) if auth_result is not None else indexed diff --git a/paimon-python/pypaimon/read/native_plan.py b/paimon-python/pypaimon/read/native_plan.py index b6ffdda29b89..5520f9c165e0 100644 --- a/paimon-python/pypaimon/read/native_plan.py +++ b/paimon-python/pypaimon/read/native_plan.py @@ -373,7 +373,7 @@ def native_plan( row_position_slice: Optional[Tuple[int, int]] = None, row_position_shard: Optional[Tuple[int, int]] = None, chunk_shuffle: Optional[Tuple[int, int]] = None, - chunk_shuffle_shard: Optional[Tuple[int, int]] = None) -> Plan: + shard: Optional[Tuple[int, int]] = None) -> Plan: """Plan with pypaimon_rust, preserving snapshot metadata. Native conversion or planning failures are handled by TableScan, which @@ -394,20 +394,16 @@ def native_plan( scan = scan.with_row_position_shard(*row_position_shard) if chunk_shuffle is not None: seed, chunk_size = chunk_shuffle - shard_index, shard_count = ( - chunk_shuffle_shard if chunk_shuffle_shard is not None - else (None, None) - ) scan = scan.with_chunk_shuffle(str(seed), chunk_size) - if shard_index is not None: - scan = scan.with_chunk_shuffle_shard(shard_index, shard_count) + if shard is not None: + scan = scan.with_shard(*shard) rust_plan = scan.plan() rust_splits = rust_plan.splits() pfields = _partition_fields(table) # Trimmed primary keys decode per-file min/max keys (PK merge-on-read). kfields = table.trimmed_primary_keys_fields splits = [ - _native_split_metadata_view(split, pfields, kfields) + deserialize_split_v1(split.serialize(), pfields, kfields) for split in rust_splits ] if table.options.native_read_enabled(): @@ -429,42 +425,3 @@ def native_plan( # table without snapshots. Let the Python scanner recover the metadata. raise RuntimeError("Native runtime cannot report an empty plan's snapshot") return Plan(splits, snapshot_id=snapshot_id) - - -def _native_split_metadata_view(rust_split, partition_fields, key_fields): - """Build the Python metadata facade while retaining native-only ranges. - - File-local ranges cannot be represented by Java SplitSerializer v1. They - remain authoritative on ``rust_split``; SlicedSplit is only the public - Python view used for row counts, paths and diagnostics. - """ - serialize_metadata = (getattr(rust_split, 'serialize_metadata', None) - if hasattr(type(rust_split), 'serialize_metadata') - else None) - payload = (serialize_metadata() if callable(serialize_metadata) - else rust_split.serialize()) - split = deserialize_split_v1(payload, partition_fields, key_fields) - - exact_count_fn = (getattr(rust_split, 'exact_merged_row_count', None) - if hasattr(type(rust_split), 'exact_merged_row_count') - else None) - exact_count = exact_count_fn() if callable(exact_count_fn) else None - file_ranges_fn = (getattr(rust_split, 'file_row_ranges', None) - if hasattr(type(rust_split), 'file_row_ranges') - else None) - file_ranges = file_ranges_fn() if callable(file_ranges_fn) else None - if file_ranges is not None: - from pypaimon.read.sliced_split import SlicedSplit - split = SlicedSplit( - split, file_ranges, exact_merged_row_count=exact_count) - else: - from pypaimon.globalindex.indexed_split import IndexedSplit - if isinstance(split, IndexedSplit) and exact_count is not None: - split = IndexedSplit( - split.data_split(), split.row_ranges(), split.scores(), - exact_merged_row_count=exact_count) - elif exact_count is not None and split.merged_row_count() != exact_count: - from pypaimon.read.sliced_split import SlicedSplit - split = SlicedSplit( - split, {}, exact_merged_row_count=exact_count) - return split diff --git a/paimon-python/pypaimon/read/scanner/chunk_shuffle_split_generator.py b/paimon-python/pypaimon/read/scanner/chunk_shuffle_split_generator.py index 11bf86c1f72b..ca1906899405 100644 --- a/paimon-python/pypaimon/read/scanner/chunk_shuffle_split_generator.py +++ b/paimon-python/pypaimon/read/scanner/chunk_shuffle_split_generator.py @@ -26,7 +26,6 @@ from pypaimon.manifest.schema.data_file_meta import DataFileMeta from pypaimon.manifest.schema.manifest_entry import ManifestEntry from pypaimon.read.scanner.split_generator import AbstractSplitGenerator -from pypaimon.read.sliced_split import SlicedSplit from pypaimon.read.split import DataSplit, Split from pypaimon.table.row.generic_row import GenericRow from pypaimon.table.source.deletion_file import DeletionFile @@ -47,17 +46,16 @@ def _null_safe_partition_key(partition_values) -> tuple: @dataclass class _PhysicalRowSlice: - """A half-open physical row slice containing visible rows.""" + """File-local inclusive ranges containing only visible rows.""" - start_inclusive: int - end_exclusive: int + row_ranges: List[Range] live_row_count: int - def to_closed_row_id_range(self, first_row_id: int) -> Range: - return Range( - first_row_id + self.start_inclusive, - first_row_id + self.end_exclusive - 1, - ) + def to_closed_row_id_ranges(self, first_row_id: int) -> List[Range]: + return [ + Range(first_row_id + row_range.from_, first_row_id + row_range.to) + for row_range in self.row_ranges + ] class _LiveRowRangeSlicer: @@ -91,8 +89,8 @@ def take(self, expected_live_rows: int) -> Optional[_PhysicalRowSlice]: if self._physical_position >= self._physical_row_count: return None - start = self._physical_position live_rows = 0 + row_ranges = [] while self._physical_position < self._physical_row_count: if self._next_deleted_position is None: @@ -100,17 +98,32 @@ def take(self, expected_live_rows: int) -> Optional[_PhysicalRowSlice]: expected_live_rows - live_rows, self._physical_row_count - self._physical_position, ) - self._physical_position += take - live_rows += take + if take > 0: + row_ranges.append(Range( + self._physical_position, + self._physical_position + take - 1, + )) + self._physical_position += take + live_rows += take else: live_run = self._next_deleted_position - self._physical_position needed = expected_live_rows - live_rows if needed <= live_run: - self._physical_position += needed - live_rows += needed + if needed > 0: + row_ranges.append(Range( + self._physical_position, + self._physical_position + needed - 1, + )) + self._physical_position += needed + live_rows += needed else: - self._physical_position += live_run - live_rows += live_run + if live_run > 0: + row_ranges.append(Range( + self._physical_position, + self._physical_position + live_run - 1, + )) + self._physical_position += live_run + live_rows += live_run if live_rows == expected_live_rows: # Deleted rows have zero live-row weight. Attach a deletion run @@ -118,8 +131,7 @@ def take(self, expected_live_rows: int) -> Optional[_PhysicalRowSlice]: # range starts at a live row (or EOF). self._skip_deleted_positions_at_cursor() return _PhysicalRowSlice( - start, - self._physical_position, + row_ranges, live_rows, ) @@ -129,7 +141,7 @@ def take(self, expected_live_rows: int) -> Optional[_PhysicalRowSlice]: if live_rows == 0: return None - return _PhysicalRowSlice(start, self._physical_position, live_rows) + return _PhysicalRowSlice(row_ranges, live_rows) def _skip_deleted_positions_at_cursor(self) -> None: while ( @@ -187,7 +199,7 @@ class ChunkShuffleSplitGeneratorBase(AbstractSplitGenerator): 7. Map each chunk through :meth:`_chunk_to_split`. Subclasses implement the three abstract hooks. Chunks ride on existing - reader wrappers (``SlicedSplit`` / ``IndexedSplit``). + reader wrapper (``IndexedSplit``). """ def __init__( @@ -323,16 +335,9 @@ def _live_row_slicer( @dataclass class _FileSegment: - """A contiguous slice of a data file inside one chunk. - - start/end are half-open row offsets within the file when the chunk - boundary falls inside the file; both are None when the chunk owns - the full file (so SlicedSplit's shard_file_idx_map can skip it and - treat the file as full — see sliced_split.py:73-78). - """ + """Visible file-local row ranges from a data file inside one chunk.""" file: DataFileMeta - start: Optional[int] - end: Optional[int] + row_ranges: List[Range] live_row_count: int @@ -381,27 +386,13 @@ def _slice_group_into_chunks( if physical_slice is None: break - if ( - physical_slice.start_inclusive == 0 - and physical_slice.end_exclusive == file.row_count - ): - current.append( - _FileSegment( - file, - None, - None, - physical_slice.live_row_count, - ) - ) - else: - current.append( - _FileSegment( - file, - physical_slice.start_inclusive, - physical_slice.end_exclusive, - physical_slice.live_row_count, - ) + current.append( + _FileSegment( + file, + physical_slice.row_ranges, + physical_slice.live_row_count, ) + ) current_rows += physical_slice.live_row_count @@ -412,11 +403,16 @@ def _slice_group_into_chunks( def _chunk_to_split(self, chunk: _Chunk) -> Split: files: List[DataFileMeta] = [] - shard_file_idx_map = {} + row_ranges = [] + split_offset = 0 for seg in chunk.segments: files.append(seg.file) - if seg.start is not None and seg.end is not None: - shard_file_idx_map[seg.file.file_name] = (seg.start, seg.end) + row_ranges.extend( + Range(split_offset + row_range.from_, + split_offset + row_range.to) + for row_range in seg.row_ranges + ) + split_offset += seg.file.row_count # set_file_path is already done once per unique file in # ChunkShuffleSplitGeneratorBase.create_splits. @@ -435,19 +431,11 @@ def _chunk_to_split(self, chunk: _Chunk) -> Split: snapshot_id=self.snapshot_id, ) - exact_merged_row_count = sum( - seg.live_row_count for seg in chunk.segments + return IndexedSplit( + data_split, + Range.sort_and_merge_overlap(row_ranges, True), + scores=None, ) - if ( - shard_file_idx_map - or data_split.merged_row_count() != exact_merged_row_count - ): - return SlicedSplit( - data_split, - shard_file_idx_map, - exact_merged_row_count=exact_merged_row_count, - ) - return data_split # --------------------------------------------------------------------------- @@ -461,11 +449,11 @@ class _AlignedGroupSegment: ``files`` is the entire group (may include blob/vector siblings), so the reader sees every column file even when only a slice of the - group's row_id range lands in this chunk. ``row_range`` is the - inclusive global row_id range this segment owns. + group's row_id range lands in this chunk. ``row_ranges`` are the + inclusive global row-id ranges of visible rows this segment owns. """ files: List[DataFileMeta] - row_range: Range + row_ranges: List[Range] live_row_count: int @@ -554,11 +542,10 @@ def _slice_group_into_chunks( physical_slice = slicer.take(avail) if physical_slice is None: break - seg_range = physical_slice.to_closed_row_id_range(first_row_id) current.append( _AlignedGroupSegment( group_files, - seg_range, + physical_slice.to_closed_row_id_ranges(first_row_id), physical_slice.live_row_count, ) ) @@ -573,14 +560,14 @@ def _chunk_to_split(self, chunk: _Chunk) -> Split: segments = chunk.segments if len(segments) == 1: all_files = segments[0].files - row_ranges = [segments[0].row_range] + row_ranges = list(segments[0].row_ranges) else: all_files = [] row_ranges = [] for seg in segments: all_files.extend(seg.files) - row_ranges.append(seg.row_range) - row_ranges.sort(key=lambda r: r.from_) + row_ranges.extend(seg.row_ranges) + row_ranges = Range.sort_and_merge_overlap(row_ranges, True) data_deletion_files = self._get_deletion_files_for_split( all_files, @@ -599,9 +586,6 @@ def _chunk_to_split(self, chunk: _Chunk) -> Split: data_split, row_ranges, scores=None, - exact_merged_row_count=sum( - seg.live_row_count for seg in segments - ), ) @staticmethod diff --git a/paimon-python/pypaimon/read/sliced_split.py b/paimon-python/pypaimon/read/sliced_split.py index 8d69f644c6a9..c3a9c5b030a7 100644 --- a/paimon-python/pypaimon/read/sliced_split.py +++ b/paimon-python/pypaimon/read/sliced_split.py @@ -19,7 +19,7 @@ SlicedSplit wraps a Split with file index ranges for shard/slice processing. """ -from typing import Dict, List, Optional, Tuple +from typing import Dict, List, Tuple from pypaimon.read.split import Split @@ -41,11 +41,9 @@ def __init__( self, data_split: 'Split', shard_file_idx_map: Dict[str, Tuple[int, int]], - exact_merged_row_count: Optional[int] = None, ): self._data_split = data_split self._shard_file_idx_map = shard_file_idx_map - self._exact_merged_row_count = exact_merged_row_count def data_split(self) -> 'Split': return self._data_split @@ -112,8 +110,6 @@ def _get_sliced_file_row_count(self, file: 'DataFileMeta') -> int: return file.row_count def merged_row_count(self): - if self._exact_merged_row_count is not None: - return self._exact_merged_row_count if not self._shard_file_idx_map: return self._data_split.merged_row_count() @@ -197,17 +193,14 @@ def __eq__(self, other): return ( self._data_split == other._data_split and self._shard_file_idx_map == other._shard_file_idx_map - and self._exact_merged_row_count == other._exact_merged_row_count ) def __hash__(self): return hash(( id(self._data_split), tuple(sorted(self._shard_file_idx_map.items())), - self._exact_merged_row_count, )) def __repr__(self): return (f"SlicedSplit(data_split={self._data_split}, " - f"shard_file_idx_map={self._shard_file_idx_map}, " - f"exact_merged_row_count={self._exact_merged_row_count})") + f"shard_file_idx_map={self._shard_file_idx_map})") diff --git a/paimon-python/pypaimon/read/split_read.py b/paimon-python/pypaimon/read/split_read.py index 4f1fbd64bfd3..cef05f8146aa 100644 --- a/paimon-python/pypaimon/read/split_read.py +++ b/paimon-python/pypaimon/read/split_read.py @@ -229,19 +229,25 @@ def _push_down_predicate(self) -> Optional[Predicate]: def create_reader(self) -> RecordReader: """Create a record reader for the given split.""" - # row_ranges: from IndexedSplit (ANN vector search), a list of discrete global row ID ranges. + # row_ranges: stable global row IDs for data-evolution reads. + # physical_row_ranges: file-local physical positions for IndexedSplit raw reads. # shard_range: from SlicedSplit (parallel shard scan), a contiguous [start, end) row range within the file. def file_reader_supplier(self, file: DataFileMeta, for_merge_read: bool, read_fields: List[str], row_tracking_enabled: bool, row_ranges: Optional[List[Range]] = None, + physical_row_ranges: Optional[List[Range]] = None, shard_range: Optional[Tuple[int, int]] = None) -> RecordBatchReader: + if row_ranges is not None and physical_row_ranges is not None: + raise ValueError( + "row_ranges and physical_row_ranges cannot be used together") ( read_file_fields, read_arrow_predicate, read_paimon_predicate, ) = self._get_fields_and_predicate(file.schema_id, read_fields) if (file.file_name in self.deletion_file_readers - or (for_merge_read and self.row_ranges is not None)): + or (for_merge_read and self.row_ranges is not None) + or physical_row_ranges is not None): # DVs and indexed PK ranges refer to physical file positions. # Filtering or skipping row groups here would renumber those rows. # Apply the residual predicate after position selection and merging. @@ -255,18 +261,29 @@ def file_reader_supplier(self, file: DataFileMeta, for_merge_read: bool, batch_size = self.table.options.read_batch_size() effective_row_ranges = None + row_range_base = file.first_row_id if row_ranges is not None: effective_row_ranges = Range.and_(row_ranges, [file.row_id_range()]) if len(effective_row_ranges) == 0: return EmptyRecordBatchReader() + elif physical_row_ranges is not None: + effective_row_ranges = Range.and_( + physical_row_ranges, + [Range(0, file.row_count - 1)], + ) + row_range_base = 0 + if len(effective_row_ranges) == 0: + return EmptyRecordBatchReader() row_sidecar_file = self._row_sidecar_file_name(file) - if row_sidecar_file is not None and self._should_read_row_sidecar( - file, - effective_row_ranges, - row_sidecar_file, - self.table.options.data_evolution_row_sidecar_max_selected_rows(), - self.table.options.data_evolution_row_sidecar_max_selection_ratio()): + if (physical_row_ranges is None + and row_sidecar_file is not None + and self._should_read_row_sidecar( + file, + effective_row_ranges, + row_sidecar_file, + self.table.options.data_evolution_row_sidecar_max_selected_rows(), + self.table.options.data_evolution_row_sidecar_max_selection_ratio())): file_path = self._aligned_extra_file_path(file, row_sidecar_file) file_format = ROW_SIDECAR_FORMAT @@ -283,7 +300,7 @@ def file_reader_supplier(self, file: DataFileMeta, for_merge_read: bool, CoreOptions.FILE_FORMAT_ROW) if file_format in row_index_formats: row_indices = [ - row_id - file.first_row_id + row_id - row_range_base for row_range in effective_row_ranges for row_id in range(row_range.from_, row_range.to + 1) ] @@ -293,10 +310,10 @@ def file_reader_supplier(self, file: DataFileMeta, for_merge_read: bool, merged_ranges = Range.sort_and_merge_overlap( effective_row_ranges, True) for r in merged_ranges: - start = max(0, r.from_ - file.first_row_id) + start = max(0, r.from_ - row_range_base) end = min( file.row_count - 1, - r.to - file.first_row_id, + r.to - row_range_base, ) if end >= start: parquet_row_ranges.append((start, end)) @@ -494,10 +511,11 @@ def file_reader_supplier(self, file: DataFileMeta, for_merge_read: bool, target_data_fields=target_fields) # For non-Vortex formats, wrap with RowIdFilterRecordBatchReader - if (row_ranges is not None + if (effective_row_ranges is not None and row_indices is None and parquet_row_ranges is None): - reader = RowIdFilterRecordBatchReader(reader, file.first_row_id, effective_row_ranges) + reader = RowIdFilterRecordBatchReader( + reader, row_range_base, effective_row_ranges) # For formats without native shard support, wrap with ShardBatchReader if shard_range is not None and file_format not in ( @@ -828,6 +846,30 @@ def _genarate_deletion_file_readers(self): self.table.file_io, df) +def _split_local_row_ranges_by_file( + files: List[DataFileMeta], + row_ranges: List[Range]) -> Dict[str, List[Range]]: + """Map split-local physical positions to file-local ranges. + + The split coordinate is the concatenation of ``files`` in list order, + matching Java ``IndexedSplit`` raw-read semantics and the native reader. + """ + ranges_by_file = {} + split_offset = 0 + for file in files: + selected = Range.and_( + row_ranges, + [Range(split_offset, split_offset + file.row_count - 1)], + ) + ranges_by_file[file.file_name] = [ + Range(row_range.from_ - split_offset, + row_range.to - split_offset) + for row_range in selected + ] + split_offset += file.row_count + return ranges_by_file + + class RawFileSplitRead(SplitRead): def __init__( self, @@ -839,6 +881,12 @@ def __init__( outer_extract_name_paths: Optional[List[List[str]]] = None, outer_flat_read_type: Optional[List[DataField]] = None, limit: Optional[int] = None): + self._physical_row_ranges = {} + actual_split = split + if isinstance(split, IndexedSplit): + self._physical_row_ranges = _split_local_row_ranges_by_file( + split.files, split.row_ranges()) + actual_split = split.data_split() # Nested-leaf projection is NOT pushed down by name: a leaf path is # only valid against the latest schema, while each data file stores # its own (possibly renamed / retyped) sub-fields. Instead the read @@ -849,7 +897,7 @@ def __init__( table=table, predicate=predicate, read_type=read_type, - split=split, + split=actual_split, row_tracking_enabled=row_tracking_enabled, nested_name_paths=None, limit=limit) @@ -858,6 +906,10 @@ def __init__( def raw_reader_supplier(self, file: DataFileMeta, dv_factory: Optional[Callable] = None) -> Optional[RecordReader]: read_fields = self._get_final_read_data_fields() + physical_row_ranges = getattr( + self, '_physical_row_ranges', {}).get(file.file_name) + if physical_row_ranges == []: + return None # Check if this is a SlicedSplit to get shard_file_idx_map shard_file_idx_map = ( self.split.shard_file_idx_map() if isinstance(self.split, SlicedSplit) else {} @@ -871,16 +923,27 @@ def raw_reader_supplier(self, file: DataFileMeta, dv_factory: Optional[Callable] for_merge_read=False, read_fields=read_fields, row_tracking_enabled=True, + physical_row_ranges=physical_row_ranges, shard_range=(start_pos, end_pos)) else: file_batch_reader = self.file_reader_supplier( file=file, for_merge_read=False, read_fields=read_fields, - row_tracking_enabled=True) + row_tracking_enabled=True, + physical_row_ranges=physical_row_ranges) dv = dv_factory() if dv_factory else None if dv: - if file.file_name in shard_file_idx_map: + if physical_row_ranges is not None: + dv = PositionMappedDeletionVector( + dv, + row_positions=[ + position + for row_range in physical_row_ranges + for position in range(row_range.from_, row_range.to + 1) + ], + ) + elif file.file_name in shard_file_idx_map: dv = PositionMappedDeletionVector( dv, file_offset=start_pos, @@ -912,6 +975,7 @@ def create_reader(self) -> RecordReader: and (self.table.is_primary_key_table or not self._arrow_filter_pushdown_enabled or self.deletion_file_readers + or self._physical_row_ranges or any(file.schema_id != self.table.table_schema.id for file in self.split.files))): reader = FilterRecordBatchReader( diff --git a/paimon-python/pypaimon/read/table_scan.py b/paimon-python/pypaimon/read/table_scan.py index 7d7e6d074fd4..fe1743750c1b 100755 --- a/paimon-python/pypaimon/read/table_scan.py +++ b/paimon-python/pypaimon/read/table_scan.py @@ -277,7 +277,7 @@ def _try_native_plan(self) -> Optional[Plan]: if chunk_shuffle is not None: extra_options['chunk_shuffle'] = chunk_shuffle if fs.idx_of_this_subtask is not None: - extra_options['chunk_shuffle_shard'] = ( + extra_options['shard'] = ( fs.idx_of_this_subtask, fs.number_of_para_subtasks) if has_distribution and fs.data_evolution and chunk_shuffle is None: if fs.idx_of_this_subtask is not None: diff --git a/paimon-python/pypaimon/tests/data_evolution_deletion_vector_test.py b/paimon-python/pypaimon/tests/data_evolution_deletion_vector_test.py index 4ef3748f78c1..fc272fdaab81 100644 --- a/paimon-python/pypaimon/tests/data_evolution_deletion_vector_test.py +++ b/paimon-python/pypaimon/tests/data_evolution_deletion_vector_test.py @@ -27,6 +27,7 @@ from pypaimon.deletionvectors.bitmap_deletion_vector import BitmapDeletionVector from pypaimon.manifest.schema.data_file_meta import DataFileMeta from pypaimon.manifest.schema.simple_stats import SimpleStats +from pypaimon.globalindex.indexed_split import IndexedSplit from pypaimon.read.reader.concat_batch_reader import ( BlobFallbackBatchReader, DataEvolutionMergeReader, @@ -35,7 +36,10 @@ from pypaimon.read.reader.iface.record_batch_reader import RecordBatchReader from pypaimon.read.sliced_split import SlicedSplit from pypaimon.read.split import DataSplit -from pypaimon.read.split_read import RawFileSplitRead +from pypaimon.read.split_read import ( + RawFileSplitRead, + _split_local_row_ranges_by_file, +) from pypaimon.table.row.blob import Blob, BlobData from pypaimon.table.row.generic_row import GenericRow from pypaimon.table.source.deletion_file import DeletionFile @@ -196,6 +200,53 @@ def test_append_sliced_reader_maps_positions_to_original_file_offsets(self): ) self.assertIsNone(reader.read_arrow_batch()) + def test_append_indexed_ranges_cross_file_boundary_and_map_dv_positions(self): + first = _file("first.parquet", None, 5, 1) + second = _file("second.parquet", None, 4, 1) + self.assertEqual( + _split_local_row_ranges_by_file( + [first, second], + [Range(3, 4), Range(5, 7)], + ), + { + "first.parquet": [Range(3, 4)], + "second.parquet": [Range(0, 2)], + }, + ) + + data_split = DataSplit( + files=[first, second], + partition=GenericRow([], []), + bucket=0, + raw_convertible=True, + data_deletion_files=None, + ) + indexed = IndexedSplit(data_split, [Range(3, 4), Range(5, 7)]) + split_read = RawFileSplitRead.__new__(RawFileSplitRead) + split_read.split = indexed.data_split() + split_read._physical_row_ranges = _split_local_row_ranges_by_file( + indexed.files, indexed.row_ranges()) + split_read._get_final_read_data_fields = Mock(return_value=[]) + split_read.file_reader_supplier = Mock( + return_value=_OneBatchReader([0, 1, 2]) + ) + deletion_vector = BitmapDeletionVector() + deletion_vector.delete(1) + + reader = split_read.raw_reader_supplier( + second, + dv_factory=lambda: deletion_vector, + ) + + self.assertEqual([0, 2], reader.read_arrow_batch().column(0).to_pylist()) + split_read.file_reader_supplier.assert_called_once_with( + file=second, + for_merge_read=False, + read_fields=[], + row_tracking_enabled=True, + physical_row_ranges=[Range(0, 2)], + ) + def test_data_evolution_merge_reader_handles_fully_deleted_file(self): deletion_vector = BitmapDeletionVector() deletion_vector.delete(0) diff --git a/paimon-python/pypaimon/tests/native_plan_distribution_test.py b/paimon-python/pypaimon/tests/native_plan_distribution_test.py index efb0f6cebeff..f6b8d029e8d9 100644 --- a/paimon-python/pypaimon/tests/native_plan_distribution_test.py +++ b/paimon-python/pypaimon/tests/native_plan_distribution_test.py @@ -353,9 +353,6 @@ def test_partial_dv_split_cannot_be_estimated_or_dropped_by_limit(self): plan, _ = self._read(table, False, slice_=(2, 9)) sliced = plan.splits()[0] self.assertIsNone(sliced.merged_row_count()) - exact = SlicedSplit(sliced.data_split(), sliced.shard_file_idx_map(), - exact_merged_row_count=int(expected_key == 2)) - self.assertEqual(exact.merged_row_count(), int(expected_key == 2)) _, actual = self._read(table, False, slice_=(2, 9), limit=1) self.assertEqual(actual, [rows[expected_key]]) diff --git a/paimon-python/pypaimon/tests/native_plan_test.py b/paimon-python/pypaimon/tests/native_plan_test.py index a61023e67f19..8ea40f6b2bb6 100644 --- a/paimon-python/pypaimon/tests/native_plan_test.py +++ b/paimon-python/pypaimon/tests/native_plan_test.py @@ -33,7 +33,6 @@ from pypaimon.globalindex.vector_search_result import ScoredGlobalIndexResult from pypaimon.read.native_plan import ( _catalog_options, - _native_split_metadata_view, _predicate_to_native, _read_options, _resolved_schema_json, @@ -142,30 +141,6 @@ def test_resolved_schema_json_uses_canonical_nested_collection_types(self): 'key': 'STRING NOT NULL', 'value': 'INT'}}}]}}}], }) - def test_native_split_metadata_view_preserves_file_ranges_and_exact_count(self): - class NativeSplit: - def serialize_metadata(self): - return b'metadata' - - def file_row_ranges(self): - return {'f.parquet': (2, 5)} - - def exact_merged_row_count(self): - return 3 - - base = Mock() - with patch( - 'pypaimon.read.native_plan.deserialize_split_v1', - return_value=base) as deserialize: - view = _native_split_metadata_view(NativeSplit(), [], []) - - from pypaimon.read.sliced_split import SlicedSplit - self.assertIsInstance(view, SlicedSplit) - self.assertIs(view.data_split(), base) - self.assertEqual(view.shard_file_idx_map(), {'f.parquet': (2, 5)}) - self.assertEqual(view.merged_row_count(), 3) - deserialize.assert_called_once_with(b'metadata', [], []) - def test_resolved_schema_keeps_custom_io_and_rest_on_catalog_path(self): from pypaimon.catalog.catalog_environment import CatalogEnvironment from pypaimon.catalog.jdbc_catalog_loader import JdbcCatalogLoader diff --git a/paimon-python/pypaimon/tests/scanner/chunk_shuffle_split_generator_test.py b/paimon-python/pypaimon/tests/scanner/chunk_shuffle_split_generator_test.py index b35f120f7529..cfb6335ba5cc 100644 --- a/paimon-python/pypaimon/tests/scanner/chunk_shuffle_split_generator_test.py +++ b/paimon-python/pypaimon/tests/scanner/chunk_shuffle_split_generator_test.py @@ -170,34 +170,32 @@ def test_slices_by_live_rows_and_absorbs_deletion_runs(self): second = slicer.take(3) self.assertEqual( - ( - first.start_inclusive, - first.end_exclusive, - first.live_row_count, - ), - (0, 6, 3), + (first.row_ranges, first.live_row_count), + ([Range(1, 2), Range(5, 5)], 3), ) self.assertEqual( - ( - second.start_inclusive, - second.end_exclusive, - second.live_row_count, - ), - (6, 10, 3), + (second.row_ranges, second.live_row_count), + ([Range(6, 8)], 3), + ) + self.assertEqual( + first.to_closed_row_id_ranges(100), + [Range(101, 102), Range(105, 105)], ) - self.assertEqual(first.to_closed_row_id_range(100), Range(100, 105)) self.assertIsNone(slicer.take(3)) + alternating = _LiveRowRangeSlicer(8, iter([1, 3, 5, 7])) + self.assertEqual( + alternating.take(3).row_ranges, + [Range(0, 0), Range(2, 2), Range(4, 4)], + ) + self.assertEqual(alternating.take(3).row_ranges, [Range(6, 6)]) + def test_returns_smaller_tail_and_skips_fully_deleted_source(self): tail_slicer = _LiveRowRangeSlicer(6, iter([1, 4])) tail = tail_slicer.take(10) self.assertEqual( - ( - tail.start_inclusive, - tail.end_exclusive, - tail.live_row_count, - ), - (0, 6, 4), + (tail.row_ranges, tail.live_row_count), + ([Range(0, 0), Range(2, 3), Range(5, 5)], 4), ) self.assertIsNone(tail_slicer.take(1)) @@ -227,10 +225,12 @@ def test_full_files_no_truncation(self): ] gen = _make_generator(seed=1, chunk_size=100) splits = gen.create_splits(entries) - # 3 chunks, each holding exactly one whole file → all DataSplit, no SlicedSplit + # Chunk output always uses IndexedSplit so row ranges carry the + # coordinate contract through stable split serialization. self.assertEqual(len(splits), 3) for s in splits: - self.assertIsInstance(s, DataSplit) + self.assertIsInstance(s, IndexedSplit) + self.assertEqual(s.row_ranges(), [Range(0, 99)]) self.assertEqual(s.row_count, 100) def test_chunk_truncates_inside_file(self): @@ -239,14 +239,17 @@ def test_chunk_truncates_inside_file(self): gen = _make_generator(seed=1, chunk_size=100, snapshot_id=7) splits = gen.create_splits(entries) self.assertEqual(len(splits), 3) - # All three chunks slice the same file → all SlicedSplit + # All three chunks carry split-local physical positions. for s in splits: - self.assertIsInstance(s, SlicedSplit) + self.assertIsInstance(s, IndexedSplit) self.assertEqual(s.snapshot_id, 7) - # union of (start, end) intervals must cover [0, 250) - intervals = sorted(s.shard_file_idx_map()['f1'] for s in splits) - self.assertEqual(intervals, [(0, 100), (100, 200), (200, 250)]) - total = sum(end - start for start, end in intervals) + intervals = sorted( + (row_range.from_, row_range.to) + for split in splits + for row_range in split.row_ranges() + ) + self.assertEqual(intervals, [(0, 99), (100, 199), (200, 249)]) + total = sum(end - start + 1 for start, end in intervals) self.assertEqual(total, 250) def test_chunk_spans_multiple_files(self): @@ -271,8 +274,8 @@ def test_chunk_size_larger_than_total(self): gen = _make_generator(seed=1, chunk_size=1000) splits = gen.create_splits(entries) self.assertEqual(len(splits), 1) - # No truncation — full files inside one chunk → DataSplit not SlicedSplit - self.assertIsInstance(splits[0], DataSplit) + self.assertIsInstance(splits[0], IndexedSplit) + self.assertEqual(splits[0].row_ranges(), [Range(0, 59)]) self.assertEqual(_split_rows(splits[0]), 60) def test_deletion_vector_slices_by_live_rows_and_is_attached(self): @@ -293,8 +296,12 @@ def test_deletion_vector_slices_by_live_rows_and_is_attached(self): splits = gen.create_splits([entry]) self.assertEqual( - sorted(s.shard_file_idx_map()['f1'] for s in splits), - [(0, 6), (6, 10)], + sorted( + (row_range.from_, row_range.to) + for split in splits + for row_range in split.row_ranges() + ), + [(1, 2), (5, 5), (6, 8)], ) self.assertEqual( sorted(s.merged_row_count() for s in splits), @@ -323,9 +330,9 @@ def test_whole_file_unknown_dv_cardinality_preserves_live_row_count(self): splits = gen.create_splits([entry]) self.assertEqual(len(splits), 1) - self.assertIsInstance(splits[0], SlicedSplit) - self.assertEqual(splits[0].shard_file_idx_map(), {}) - self.assertEqual(splits[0].row_count, 10) + self.assertIsInstance(splits[0], IndexedSplit) + self.assertEqual(splits[0].row_ranges(), [Range(1, 2), Range(5, 8)]) + self.assertEqual(splits[0].row_count, 6) self.assertEqual(splits[0].merged_row_count(), 6) read.assert_called_once_with(gen.table.file_io, deletion_file) @@ -371,11 +378,15 @@ def test_two_deletion_vector_files_share_one_live_row_chunk(self): splits = gen.create_splits(entries) self.assertEqual(len(splits), 1) - self.assertIsInstance(splits[0], DataSplit) + self.assertIsInstance(splits[0], IndexedSplit) self.assertEqual( [file.file_name for file in splits[0].files], ['f1', 'f2'], ) + self.assertEqual( + splits[0].row_ranges(), + [Range(0, 0), Range(2, 2), Range(4, 4), Range(6, 8)], + ) self.assertEqual( splits[0].data_deletion_files, [first_deletion_file, second_deletion_file], @@ -853,7 +864,10 @@ def test_deletion_vector_uses_anchor_and_slices_by_live_rows(self): for split in splits for row_range in split.row_ranges() ) - self.assertEqual(ranges, [(100, 104), (105, 109)]) + self.assertEqual( + ranges, + [(100, 100), (103, 104), (105, 106), (108, 108)], + ) self.assertEqual( sorted(split.merged_row_count() for split in splits), [3, 3], @@ -919,7 +933,8 @@ def test_multiple_deletion_vector_groups_share_one_live_row_chunk(self): (row_range.from_, row_range.to) for row_range in splits[0].row_ranges() ], - [(100, 104), (200, 205)], + [(100, 100), (102, 102), (104, 104), + (201, 201), (203, 204)], ) self.assertEqual( [file.file_name for file in splits[0].files], From fb79b662af27a5c55f1982cf7e562b812eb9e74c Mon Sep 17 00:00:00 2001 From: JingsongLi Date: Mon, 21 Sep 2026 08:46:28 +0800 Subject: [PATCH 4/5] [python] Remove unused native plan test import --- paimon-python/pypaimon/tests/native_plan_distribution_test.py | 1 - 1 file changed, 1 deletion(-) diff --git a/paimon-python/pypaimon/tests/native_plan_distribution_test.py b/paimon-python/pypaimon/tests/native_plan_distribution_test.py index f6b8d029e8d9..00d3dfcd5b63 100644 --- a/paimon-python/pypaimon/tests/native_plan_distribution_test.py +++ b/paimon-python/pypaimon/tests/native_plan_distribution_test.py @@ -32,7 +32,6 @@ native_version_at_least, native_runtime_available, ) -from pypaimon.read.sliced_split import SlicedSplit from pypaimon.table.row.generic_row import GenericRow from pypaimon.utils.range import Range from pypaimon.write.commit_message import CommitMessage From 154af1a638b71efe1f96b0be79e4c03514ca43b7 Mon Sep 17 00:00:00 2001 From: JingsongLi Date: Mon, 21 Sep 2026 09:19:10 +0800 Subject: [PATCH 5/5] [python] Align row-tracked chunk range coordinates --- .../pypaimon/globalindex/indexed_split.py | 5 +- .../scanner/chunk_shuffle_split_generator.py | 14 +++- paimon-python/pypaimon/read/split_read.py | 47 +++++++++-- .../data_evolution_deletion_vector_test.py | 25 ++++++ .../tests/native_plan_chunk_shuffle_test.py | 79 ++++++++++++++----- .../chunk_shuffle_split_generator_test.py | 28 ++++++- 6 files changed, 168 insertions(+), 30 deletions(-) diff --git a/paimon-python/pypaimon/globalindex/indexed_split.py b/paimon-python/pypaimon/globalindex/indexed_split.py index fa7aecf84ce6..cd6b216d7fad 100644 --- a/paimon-python/pypaimon/globalindex/indexed_split.py +++ b/paimon-python/pypaimon/globalindex/indexed_split.py @@ -17,8 +17,9 @@ """IndexedSplit wraps a Split with row ranges and optional scores. -Ranges use the coordinate system of the read path: stable row IDs for data -evolution and physical positions for primary-key/raw append reads. +Ranges use the coordinate system of the table read path: stable row IDs for +row-tracked tables and split-local physical positions for tables without row +tracking. """ from typing import List, Optional diff --git a/paimon-python/pypaimon/read/scanner/chunk_shuffle_split_generator.py b/paimon-python/pypaimon/read/scanner/chunk_shuffle_split_generator.py index ca1906899405..6ed8846afe06 100644 --- a/paimon-python/pypaimon/read/scanner/chunk_shuffle_split_generator.py +++ b/paimon-python/pypaimon/read/scanner/chunk_shuffle_split_generator.py @@ -404,12 +404,22 @@ def _slice_group_into_chunks( def _chunk_to_split(self, chunk: _Chunk) -> Split: files: List[DataFileMeta] = [] row_ranges = [] + ranges_use_row_ids = self.table.options.row_tracking_enabled() split_offset = 0 for seg in chunk.segments: files.append(seg.file) + if ranges_use_row_ids: + if seg.file.first_row_id is None: + raise ValueError( + "Row-tracked file '%s' is missing first_row_id" + % seg.file.file_name + ) + range_base = seg.file.first_row_id + else: + range_base = split_offset row_ranges.extend( - Range(split_offset + row_range.from_, - split_offset + row_range.to) + Range(range_base + row_range.from_, + range_base + row_range.to) for row_range in seg.row_ranges ) split_offset += seg.file.row_count diff --git a/paimon-python/pypaimon/read/split_read.py b/paimon-python/pypaimon/read/split_read.py index cef05f8146aa..71b34b285346 100644 --- a/paimon-python/pypaimon/read/split_read.py +++ b/paimon-python/pypaimon/read/split_read.py @@ -230,7 +230,8 @@ def create_reader(self) -> RecordReader: """Create a record reader for the given split.""" # row_ranges: stable global row IDs for data-evolution reads. - # physical_row_ranges: file-local physical positions for IndexedSplit raw reads. + # physical_row_ranges: file-local positions after resolving the table's + # IndexedSplit coordinate system. # shard_range: from SlicedSplit (parallel shard scan), a contiguous [start, end) row range within the file. def file_reader_supplier(self, file: DataFileMeta, for_merge_read: bool, read_fields: List[str], row_tracking_enabled: bool, @@ -851,8 +852,8 @@ def _split_local_row_ranges_by_file( row_ranges: List[Range]) -> Dict[str, List[Range]]: """Map split-local physical positions to file-local ranges. - The split coordinate is the concatenation of ``files`` in list order, - matching Java ``IndexedSplit`` raw-read semantics and the native reader. + The split coordinate is the concatenation of ``files`` in list order for + tables without row tracking, matching the native reader. """ ranges_by_file = {} split_offset = 0 @@ -870,6 +871,39 @@ def _split_local_row_ranges_by_file( return ranges_by_file +def _row_ranges_by_file( + files: List[DataFileMeta], + row_ranges: List[Range], + ranges_use_row_ids: bool) -> Dict[str, List[Range]]: + """Map table-path row ranges to file-local physical positions. + + Row-tracked append tables use stable global row IDs. Tables without row + tracking use the split-local concatenation handled by + :func:`_split_local_row_ranges_by_file`. + """ + if not ranges_use_row_ids: + return _split_local_row_ranges_by_file(files, row_ranges) + + ranges_by_file = {} + for file in files: + if file.first_row_id is None: + raise ValueError( + "Row-tracked file '%s' is missing first_row_id" + % file.file_name + ) + first_row_id = file.first_row_id + selected = Range.and_( + row_ranges, + [Range(first_row_id, first_row_id + file.row_count - 1)], + ) + ranges_by_file[file.file_name] = [ + Range(row_range.from_ - first_row_id, + row_range.to - first_row_id) + for row_range in selected + ] + return ranges_by_file + + class RawFileSplitRead(SplitRead): def __init__( self, @@ -884,8 +918,11 @@ def __init__( self._physical_row_ranges = {} actual_split = split if isinstance(split, IndexedSplit): - self._physical_row_ranges = _split_local_row_ranges_by_file( - split.files, split.row_ranges()) + self._physical_row_ranges = _row_ranges_by_file( + split.files, + split.row_ranges(), + row_tracking_enabled, + ) actual_split = split.data_split() # Nested-leaf projection is NOT pushed down by name: a leaf path is # only valid against the latest schema, while each data file stores diff --git a/paimon-python/pypaimon/tests/data_evolution_deletion_vector_test.py b/paimon-python/pypaimon/tests/data_evolution_deletion_vector_test.py index fc272fdaab81..59d17c797e1a 100644 --- a/paimon-python/pypaimon/tests/data_evolution_deletion_vector_test.py +++ b/paimon-python/pypaimon/tests/data_evolution_deletion_vector_test.py @@ -38,6 +38,7 @@ from pypaimon.read.split import DataSplit from pypaimon.read.split_read import ( RawFileSplitRead, + _row_ranges_by_file, _split_local_row_ranges_by_file, ) from pypaimon.table.row.blob import Blob, BlobData @@ -247,6 +248,30 @@ def test_append_indexed_ranges_cross_file_boundary_and_map_dv_positions(self): physical_row_ranges=[Range(0, 2)], ) + def test_row_tracked_append_ranges_map_global_row_ids(self): + first = _file("first.parquet", 100, 5, 1) + second = _file("second.parquet", 200, 4, 1) + + self.assertEqual( + _row_ranges_by_file( + [first, second], + [Range(103, 104), Range(200, 202)], + ranges_use_row_ids=True, + ), + { + "first.parquet": [Range(3, 4)], + "second.parquet": [Range(0, 2)], + }, + ) + + second.first_row_id = None + with self.assertRaisesRegex(ValueError, "missing first_row_id"): + _row_ranges_by_file( + [first, second], + [Range(103, 104)], + ranges_use_row_ids=True, + ) + def test_data_evolution_merge_reader_handles_fully_deleted_file(self): deletion_vector = BitmapDeletionVector() deletion_vector.delete(0) diff --git a/paimon-python/pypaimon/tests/native_plan_chunk_shuffle_test.py b/paimon-python/pypaimon/tests/native_plan_chunk_shuffle_test.py index 2725458aef88..b13087ee45ef 100644 --- a/paimon-python/pypaimon/tests/native_plan_chunk_shuffle_test.py +++ b/paimon-python/pypaimon/tests/native_plan_chunk_shuffle_test.py @@ -39,16 +39,18 @@ reason='Rust native chunk planner required')] -@pytest.fixture(params=['append', 'append-dv', 'de', 'de-dv']) -def chunk_table(request, tmp_path): - de = request.param.startswith('de') - dv = request.param.endswith('dv') +def _create_chunk_table(kind, tmp_path): + de = kind.startswith('de') + dv = kind.endswith('dv') + row_tracking = de or kind.startswith('append-row-tracking') fields = [('id', pa.int64()), ('p', pa.string())] options = {'file.format': 'parquet', 'source.split.target-size': '1b', 'source.split.open-file-cost': '1b'} + if row_tracking: + options['row-tracking.enabled'] = 'true' if de: fields.append(('payload', pa.large_binary())) - options.update({'data-evolution.enabled': 'true', 'row-tracking.enabled': 'true', + options.update({'data-evolution.enabled': 'true', 'blob.target-file-size': '1b'}) if dv: options['deletion-vectors.enabled'] = 'true' @@ -108,49 +110,86 @@ def chunk_table(request, tmp_path): return table, expected, before_deletes -def _chunks(table, seed, chunk_size=3, shard=None, predicate=None, projection=None): +@pytest.fixture(params=['append', 'append-dv', 'de', 'de-dv']) +def chunk_table(request, tmp_path): + return _create_chunk_table(request.param, tmp_path) + + +@pytest.fixture +def row_tracking_append_table(tmp_path): + return _create_chunk_table('append-row-tracking', tmp_path) + + +def _chunks(table, seed, chunk_size=3, shard=None, predicate=None, projection=None, + runtimes=((False, False), (True, True))): results = [] - for native in (False, True): - builder = table.copy({ - 'scan.native-plan.enabled': str(native).lower(), - 'read.native.enabled': str(native).lower(), + for native_plan_enabled, native_read_enabled in runtimes: + plan_builder = table.copy({ + 'scan.native-plan.enabled': str(native_plan_enabled).lower(), + # Retain the opaque Rust split for every Rust-planned case. The + # separately configured read builder still decides which reader + # consumes it. + 'read.native.enabled': str(native_plan_enabled).lower(), }).new_read_builder() if predicate is not None: - builder.with_filter(predicate) + plan_builder.with_filter(predicate) if projection is not None: - builder.with_projection(projection) - scan = builder.new_scan().with_chunk_shuffle(seed, chunk_size) + plan_builder.with_projection(projection) + scan = plan_builder.new_scan().with_chunk_shuffle(seed, chunk_size) if shard is not None: scan.with_shard(*shard) - if native: + if native_plan_enabled: with patch.object(scan.file_scanner, 'scan', side_effect=AssertionError('native fallback')), \ patch.object(scan.file_scanner, 'plan_files', side_effect=AssertionError('Python manifest scan')): plan = scan.plan() else: plan = scan.plan() + + read_builder = table.copy({ + 'scan.native-plan.enabled': 'false', + 'read.native.enabled': str(native_read_enabled).lower(), + }).new_read_builder() + if predicate is not None: + read_builder.with_filter(predicate) + if projection is not None: + read_builder.with_projection(projection) assert all( split.snapshot_id == plan.snapshot_id for split in plan.splits()) chunks = [] for split in plan.splits(): if table.options.options.contains_key('incremental-between-timestamp'): assert split.is_streaming - if native: - assert native_reader_available() + if native_plan_enabled: assert getattr(split, '_native_split', None) is not None + if native_read_enabled: + assert native_reader_available() with patch( 'pypaimon.read.table_read.TableRead._create_split_read', side_effect=AssertionError('Python reader was used')): - rows = builder.new_read().to_arrow([split]).to_pylist() + rows = read_builder.new_read().to_arrow([split]).to_pylist() else: - rows = builder.new_read().to_arrow([split]).to_pylist() + rows = read_builder.new_read().to_arrow([split]).to_pylist() assert 0 < len(rows) <= chunk_size assert split.merged_row_count() == len(rows) if rows and 'p' in rows[0]: assert len({row['p'] for row in rows}) == 1 chunks.append(sorted(rows, key=lambda row: row['id'])) results.append((plan.snapshot_id, chunks)) - assert results[0] == results[1] - return results[1] + assert all(result == results[0] for result in results[1:]) + return results[-1] + + +def test_row_tracking_append_chunks_cross_runtime(row_tracking_append_table): + table, expected, _ = row_tracking_append_table + _, chunks = _chunks( + table, + seed=42, + runtimes=((False, False), (True, False), (False, True), (True, True)), + ) + assert sorted( + (row for chunk in chunks for row in chunk), + key=lambda row: row['id'], + ) == expected @pytest.mark.parametrize('seed', [-11, 42, 2 ** 70]) diff --git a/paimon-python/pypaimon/tests/scanner/chunk_shuffle_split_generator_test.py b/paimon-python/pypaimon/tests/scanner/chunk_shuffle_split_generator_test.py index cfb6335ba5cc..4feaa940a62f 100644 --- a/paimon-python/pypaimon/tests/scanner/chunk_shuffle_split_generator_test.py +++ b/paimon-python/pypaimon/tests/scanner/chunk_shuffle_split_generator_test.py @@ -49,10 +49,13 @@ def _mock_table(table_path='/tmp/_chunk_shuffle_test_path'): table = Mock() table.table_path = table_path table.options = Mock() + table.options.row_tracking_enabled.return_value = False return table -def _mock_entry(partition_values, bucket, file_name, row_count, file_size=1024): +def _mock_entry( + partition_values, bucket, file_name, row_count, file_size=1024, + first_row_id=None): entry = Mock() entry.partition = Mock() entry.partition.values = partition_values @@ -61,6 +64,7 @@ def _mock_entry(partition_values, bucket, file_name, row_count, file_size=1024): entry.file.file_name = file_name entry.file.file_size = file_size entry.file.row_count = row_count + entry.file.first_row_id = first_row_id # Swallow set_file_path so we don't need to mock partition path encoding. entry.file.set_file_path = Mock() return entry @@ -266,6 +270,28 @@ def test_chunk_spans_multiple_files(self): total_rows = sum(_split_rows(s) for s in splits) self.assertEqual(total_rows, 90) + def test_row_tracking_chunks_use_global_row_ids(self): + table = _mock_table() + table.options.row_tracking_enabled.return_value = True + entries = [ + _mock_entry([], 0, 'f1', 2, first_row_id=100), + _mock_entry([], 0, 'f2', 2, first_row_id=200), + ] + splits = _make_generator( + seed=1, chunk_size=3, table=table).create_splits(entries) + + ranges = sorted( + (row_range.from_, row_range.to) + for split in splits + for row_range in split.row_ranges() + ) + self.assertEqual(ranges, [(100, 101), (200, 200), (201, 201)]) + + entries[0].file.first_row_id = None + with self.assertRaisesRegex(ValueError, 'missing first_row_id'): + _make_generator( + seed=1, chunk_size=3, table=table).create_splits(entries) + def test_chunk_size_larger_than_total(self): entries = [ _mock_entry([], 0, 'f1', 30),