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 1bc2588b53 [python] chunk_shuffle should skip deleted records during
counting (#8846)
1bc2588b53 is described below
commit 1bc2588b538af3564972dddad7b782da9fd6b547
Author: Faiz <[email protected]>
AuthorDate: Tue Jul 28 21:15:47 2026 +0800
[python] chunk_shuffle should skip deleted records during counting (#8846)
---
.../pypaimon/globalindex/indexed_split.py | 25 +-
.../read/scanner/chunk_shuffle_split_generator.py | 325 +++++++++++++++--
.../pypaimon/read/scanner/file_scanner.py | 2 -
paimon-python/pypaimon/read/sliced_split.py | 24 +-
paimon-python/pypaimon/read/split_read.py | 5 +
.../tests/data_evolution_deletion_vector_test.py | 36 ++
.../scanner/chunk_shuffle_split_generator_test.py | 392 ++++++++++++++++++++-
7 files changed, 745 insertions(+), 64 deletions(-)
diff --git a/paimon-python/pypaimon/globalindex/indexed_split.py
b/paimon-python/pypaimon/globalindex/indexed_split.py
index 18b75a0af0..bd4d98a1d1 100644
--- a/paimon-python/pypaimon/globalindex/indexed_split.py
+++ b/paimon-python/pypaimon/globalindex/indexed_split.py
@@ -30,11 +30,13 @@ class IndexedSplit(Split):
self,
data_split: 'Split',
row_ranges: List['Range'],
- scores: Optional[List[float]] = None
+ 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."""
@@ -76,6 +78,8 @@ class IndexedSplit(Split):
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
@@ -133,14 +137,23 @@ class IndexedSplit(Split):
def __eq__(self, other):
if not isinstance(other, IndexedSplit):
return False
- return (self._data_split == other._data_split and
- self._row_ranges == other._row_ranges and
- self._scores == other._scores)
+ return (
+ 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):
scores_hash = tuple(self._scores) if self._scores else None
- return hash((id(self._data_split), tuple(self._row_ranges),
scores_hash))
+ return hash((
+ 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"row_ranges={self._row_ranges}, scores={self._scores}, "
+ f"exact_merged_row_count={self._exact_merged_row_count})")
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 6493b2e6fb..89216322cf 100644
--- a/paimon-python/pypaimon/read/scanner/chunk_shuffle_split_generator.py
+++ b/paimon-python/pypaimon/read/scanner/chunk_shuffle_split_generator.py
@@ -19,8 +19,9 @@ import random
from abc import abstractmethod
from collections import defaultdict
from dataclasses import dataclass
-from typing import Any, List, Optional, Tuple
+from typing import Any, Iterator, List, Optional, Tuple
+from pypaimon.deletionvectors.deletion_vector import DeletionVector
from pypaimon.globalindex.indexed_split import IndexedSplit
from pypaimon.manifest.schema.data_file_meta import DataFileMeta
from pypaimon.manifest.schema.manifest_entry import ManifestEntry
@@ -28,6 +29,8 @@ 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
+from pypaimon.utils.data_evolution_utils import retrieve_anchor_file
from pypaimon.utils.range import Range
from pypaimon.utils.range_helper import RangeHelper
@@ -42,6 +45,121 @@ def _null_safe_partition_key(partition_values) -> tuple:
return tuple((v is None, v) for v in partition_values)
+@dataclass
+class _PhysicalRowSlice:
+ """A half-open physical row slice containing visible rows."""
+
+ start_inclusive: int
+ end_exclusive: int
+ 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,
+ )
+
+
+class _LiveRowRangeSlicer:
+ """Map requested live-row counts to contiguous physical row ranges.
+
+ Deleted positions must be sorted and unique. Each position is consumed
+ once, so slicing costs O(number of output ranges + DV cardinality)
+ instead of O(physical row count).
+ """
+
+ def __init__(
+ self,
+ physical_row_count: int,
+ deleted_positions: Iterator[int],
+ ):
+ if physical_row_count < 0:
+ raise ValueError(
+ f"physical_row_count must be non-negative, got
{physical_row_count}"
+ )
+ self._physical_row_count = physical_row_count
+ self._deleted_positions = iter(deleted_positions)
+ self._physical_position = 0
+ self._last_deleted_position = None
+ self._next_deleted_position = self._advance_deleted_position()
+
+ def take(self, expected_live_rows: int) -> Optional[_PhysicalRowSlice]:
+ if expected_live_rows <= 0:
+ raise ValueError(
+ f"expected_live_rows must be positive, got
{expected_live_rows}"
+ )
+ if self._physical_position >= self._physical_row_count:
+ return None
+
+ start = self._physical_position
+ live_rows = 0
+
+ while self._physical_position < self._physical_row_count:
+ if self._next_deleted_position is None:
+ take = min(
+ expected_live_rows - live_rows,
+ self._physical_row_count - self._physical_position,
+ )
+ 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
+ else:
+ 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
+ # immediately after the boundary to this range so the next
+ # range starts at a live row (or EOF).
+ self._skip_deleted_positions_at_cursor()
+ return _PhysicalRowSlice(
+ start,
+ self._physical_position,
+ live_rows,
+ )
+
+ # The current live run was insufficient, so the cursor must be at
+ # the next deleted position. Consume it and continue.
+ self._skip_deleted_positions_at_cursor()
+
+ if live_rows == 0:
+ return None
+ return _PhysicalRowSlice(start, self._physical_position, live_rows)
+
+ def _skip_deleted_positions_at_cursor(self) -> None:
+ while (
+ self._next_deleted_position is not None
+ and self._next_deleted_position == self._physical_position
+ ):
+ self._physical_position += 1
+ self._next_deleted_position = self._advance_deleted_position()
+
+ def _advance_deleted_position(self) -> Optional[int]:
+ position = next(self._deleted_positions, None)
+ if position is None:
+ return None
+ if position < 0 or position >= self._physical_row_count:
+ raise ValueError(
+ f"Deletion vector position {position} is outside physical row "
+ f"range [0, {self._physical_row_count})."
+ )
+ if (
+ self._last_deleted_position is not None
+ and position <= self._last_deleted_position
+ ):
+ raise ValueError(
+ "Deletion vector positions must be strictly increasing, but
found "
+ f"{position} after {self._last_deleted_position}."
+ )
+ self._last_deleted_position = position
+ return position
+
+
@dataclass
class _Chunk:
"""A unit of work for one DataLoader read. ``segments`` carries
@@ -68,10 +186,8 @@ class
ChunkShuffleSplitGeneratorBase(AbstractSplitGenerator):
6. If sharded, take this worker's slice via balanced
``_compute_shard_range``.
7. Map each chunk through :meth:`_chunk_to_split`.
- Subclasses implement the three abstract hooks. Reader paths
- (``RawFileSplitRead`` for append, ``DataEvolutionSplitRead`` for DE)
- are unchanged because chunks ride on existing wrappers
- (``SlicedSplit`` / ``IndexedSplit``).
+ Subclasses implement the three abstract hooks. Chunks ride on existing
+ reader wrappers (``SlicedSplit`` / ``IndexedSplit``).
"""
def __init__(
@@ -88,6 +204,8 @@ class ChunkShuffleSplitGeneratorBase(AbstractSplitGenerator):
self.chunk_size = chunk_size
def create_splits(self, file_entries: List[ManifestEntry]) -> List[Split]:
+ """TODO: Lazily initialize DataSplits to avoid creating too many
objects."""
+
if not file_entries:
return []
@@ -146,9 +264,56 @@ class
ChunkShuffleSplitGeneratorBase(AbstractSplitGenerator):
def _chunk_to_split(self, chunk: _Chunk) -> Split:
"""Wrap a chunk into a Split that the existing readers consume."""
+ def _deletion_file(
+ self,
+ partition: GenericRow,
+ bucket: int,
+ file_name: str,
+ ) -> Optional[DeletionFile]:
+ partition_key = (tuple(partition.values), bucket)
+ return self.deletion_files_map.get(partition_key, {}).get(file_name)
+
+ def _live_row_slicer(
+ self,
+ physical_row_count: int,
+ deletion_file: Optional[DeletionFile],
+ ) -> Optional[_LiveRowRangeSlicer]:
+ if deletion_file is None or deletion_file.cardinality == 0:
+ return _LiveRowRangeSlicer(physical_row_count, iter(()))
+
+ cardinality = deletion_file.cardinality
+ if cardinality is not None:
+ if cardinality < 0 or cardinality > physical_row_count:
+ raise ValueError(
+ f"Deletion vector cardinality {cardinality} is outside
valid "
+ f"range [0, {physical_row_count}]."
+ )
+ if cardinality == physical_row_count:
+ return None
+
+ deletion_vector = DeletionVector.read(self.table.file_io,
deletion_file)
+
+ actual_cardinality = deletion_vector.get_cardinality()
+ if cardinality is not None and cardinality != actual_cardinality:
+ raise ValueError(
+ f"Deletion vector cardinality mismatch, metadata is
{cardinality} "
+ f"but bitmap contains {actual_cardinality} positions."
+ )
+ if actual_cardinality > physical_row_count:
+ raise ValueError(
+ f"Deletion vector cardinality {actual_cardinality} exceeds
physical "
+ f"row count {physical_row_count}."
+ )
+ if actual_cardinality == physical_row_count:
+ return None
+ return _LiveRowRangeSlicer(
+ physical_row_count,
+ iter(deletion_vector.bit_map()),
+ )
+
# ---------------------------------------------------------------------------
-# Append (non-DE, non-DV) implementation
+# Append implementation
# ---------------------------------------------------------------------------
@@ -164,10 +329,11 @@ class _FileSegment:
file: DataFileMeta
start: Optional[int]
end: Optional[int]
+ live_row_count: int
class AppendChunkShuffleSplitGenerator(ChunkShuffleSplitGeneratorBase):
- """Chunk-shuffled splits for plain append tables (non-PK, non-DV,
non-DE)."""
+ """Chunk-shuffled splits for plain append tables (non-PK, non-DE)."""
def _sort_key(self, entry: ManifestEntry):
return (
@@ -180,8 +346,9 @@ class
AppendChunkShuffleSplitGenerator(ChunkShuffleSplitGeneratorBase):
self, entries: List[ManifestEntry]
) -> List[List[_FileSegment]]:
"""Cut a (partition, bucket) group into chunks of at most
- ``self.chunk_size`` rows. ``chunk_size`` is a hard upper bound:
- the last chunk may be smaller, but no chunk exceeds it.
+ ``self.chunk_size`` live rows. ``chunk_size`` is a hard upper bound:
+ the last chunk may be smaller, but no chunk exceeds it after DV
+ filtering.
"""
chunks: List[List[_FileSegment]] = []
current: List[_FileSegment] = []
@@ -189,9 +356,16 @@ class
AppendChunkShuffleSplitGenerator(ChunkShuffleSplitGeneratorBase):
for entry in entries:
file = entry.file
- offset = 0
- remaining = file.row_count
- while remaining > 0:
+ deletion_file = self._deletion_file(
+ entry.partition,
+ entry.bucket,
+ file.file_name,
+ )
+ slicer = self._live_row_slicer(file.row_count, deletion_file)
+ if slicer is None:
+ continue
+
+ while True:
avail = self.chunk_size - current_rows
if avail <= 0:
chunks.append(current)
@@ -199,16 +373,33 @@ class
AppendChunkShuffleSplitGenerator(ChunkShuffleSplitGeneratorBase):
current_rows = 0
avail = self.chunk_size
- take = min(remaining, avail)
-
- if take == file.row_count and offset == 0:
- current.append(_FileSegment(file, None, None))
+ physical_slice = slicer.take(avail)
+ 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, offset, offset + take))
+ current.append(
+ _FileSegment(
+ file,
+ physical_slice.start_inclusive,
+ physical_slice.end_exclusive,
+ physical_slice.live_row_count,
+ )
+ )
- current_rows += take
- offset += take
- remaining -= take
+ current_rows += physical_slice.live_row_count
if current:
chunks.append(current)
@@ -225,17 +416,32 @@ class
AppendChunkShuffleSplitGenerator(ChunkShuffleSplitGeneratorBase):
# set_file_path is already done once per unique file in
# ChunkShuffleSplitGeneratorBase.create_splits.
+ data_deletion_files = self._get_deletion_files_for_split(
+ files,
+ chunk.partition,
+ chunk.bucket,
+ )
data_split = DataSplit(
files=files,
partition=chunk.partition,
bucket=chunk.bucket,
raw_convertible=True,
- data_deletion_files=None,
+ data_deletion_files=data_deletion_files,
)
- if shard_file_idx_map:
- return SlicedSplit(data_split, shard_file_idx_map)
+ exact_merged_row_count = sum(
+ seg.live_row_count for seg in chunk.segments
+ )
+ 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
@@ -255,6 +461,7 @@ class _AlignedGroupSegment:
"""
files: List[DataFileMeta]
row_range: Range
+ live_row_count: int
class DataEvolutionChunkShuffleSplitGenerator(ChunkShuffleSplitGeneratorBase):
@@ -295,11 +502,43 @@ class
DataEvolutionChunkShuffleSplitGenerator(ChunkShuffleSplitGeneratorBase):
chunks: List[List[_AlignedGroupSegment]] = []
current: List[_AlignedGroupSegment] = []
current_rows = 0
+ partition = entries[0].partition
+ bucket = entries[0].bucket
for group_range, group_files in aligned_groups:
- offset = 0
- group_rows = group_range.count()
- while offset < group_rows:
+ anchor = None
+ deletion_file = None
+ if self.deletion_files_map:
+ anchor = retrieve_anchor_file(group_files)
+ deletion_file = self._deletion_file(
+ partition,
+ bucket,
+ anchor.file_name,
+ )
+
+ physical_row_count = group_range.count()
+ first_row_id = group_range.from_
+ if deletion_file is not None:
+ anchor_range = anchor.row_id_range()
+ if (
+ anchor_range.from_ > group_range.from_
+ or anchor_range.to < group_range.to
+ ):
+ raise ValueError(
+ f"Data evolution anchor range {anchor_range} does not
contain "
+ f"aligned group range {group_range}."
+ )
+ physical_row_count = anchor.row_count
+ first_row_id = anchor_range.from_
+
+ slicer = self._live_row_slicer(
+ physical_row_count,
+ deletion_file,
+ )
+ if slicer is None:
+ continue
+
+ while True:
avail = self.chunk_size - current_rows
if avail <= 0:
chunks.append(current)
@@ -307,14 +546,18 @@ class
DataEvolutionChunkShuffleSplitGenerator(ChunkShuffleSplitGeneratorBase):
current_rows = 0
avail = self.chunk_size
- take = min(group_rows - offset, avail)
- seg_range = Range(
- group_range.from_ + offset,
- group_range.from_ + offset + take - 1,
+ 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.live_row_count,
+ )
)
- current.append(_AlignedGroupSegment(group_files, seg_range))
- current_rows += take
- offset += take
+ current_rows += physical_slice.live_row_count
if current:
chunks.append(current)
@@ -334,14 +577,26 @@ class
DataEvolutionChunkShuffleSplitGenerator(ChunkShuffleSplitGeneratorBase):
row_ranges.append(seg.row_range)
row_ranges.sort(key=lambda r: r.from_)
+ data_deletion_files = self._get_deletion_files_for_split(
+ all_files,
+ chunk.partition,
+ chunk.bucket,
+ )
data_split = DataSplit(
files=all_files,
partition=chunk.partition,
bucket=chunk.bucket,
raw_convertible=False,
- data_deletion_files=None,
+ data_deletion_files=data_deletion_files,
+ )
+ return IndexedSplit(
+ data_split,
+ row_ranges,
+ scores=None,
+ exact_merged_row_count=sum(
+ seg.live_row_count for seg in segments
+ ),
)
- return IndexedSplit(data_split, row_ranges, scores=None)
@staticmethod
def _split_by_row_id_with_range(
diff --git a/paimon-python/pypaimon/read/scanner/file_scanner.py
b/paimon-python/pypaimon/read/scanner/file_scanner.py
index c2fcff26a9..3283a51bf9 100755
--- a/paimon-python/pypaimon/read/scanner/file_scanner.py
+++ b/paimon-python/pypaimon/read/scanner/file_scanner.py
@@ -614,8 +614,6 @@ class FileScanner:
def _validate_chunk_shuffle_compat(self) -> None:
if self.table.is_primary_key_table:
raise ValueError("chunk_shuffle only supports append tables")
- if self.deletion_vectors_enabled:
- raise ValueError("chunk_shuffle not supported with deletion
vectors")
if self.start_pos_of_this_subtask is not None:
raise ValueError("chunk_shuffle cannot combine with with_slice")
if self.limit is not None:
diff --git a/paimon-python/pypaimon/read/sliced_split.py
b/paimon-python/pypaimon/read/sliced_split.py
index 6921207cab..0bdf42ab5d 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 List, Dict, Tuple
+from typing import Dict, List, Optional, Tuple
from pypaimon.read.split import Split
@@ -40,10 +40,12 @@ class SlicedSplit(Split):
def __init__(
self,
data_split: 'Split',
- shard_file_idx_map: Dict[str, Tuple[int, int]]
+ 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
@@ -102,6 +104,8 @@ class SlicedSplit(Split):
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()
@@ -176,12 +180,20 @@ class SlicedSplit(Split):
def __eq__(self, other):
if not isinstance(other, SlicedSplit):
return False
- return (self._data_split == other._data_split and
- self._shard_file_idx_map == other._shard_file_idx_map)
+ 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()))))
+ 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"shard_file_idx_map={self._shard_file_idx_map}, "
+ f"exact_merged_row_count={self._exact_merged_row_count})")
diff --git a/paimon-python/pypaimon/read/split_read.py
b/paimon-python/pypaimon/read/split_read.py
index 77bbd44837..e3d7f3f16d 100644
--- a/paimon-python/pypaimon/read/split_read.py
+++ b/paimon-python/pypaimon/read/split_read.py
@@ -768,6 +768,11 @@ class RawFileSplitRead(SplitRead):
row_tracking_enabled=True)
dv = dv_factory() if dv_factory else None
if dv:
+ if file.file_name in shard_file_idx_map:
+ dv = PositionMappedDeletionVector(
+ dv,
+ file_offset=start_pos,
+ )
return
ApplyDeletionVectorReader(RowPositionReader(file_batch_reader), dv)
else:
return file_batch_reader
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 448a8ab1b0..4ef3748f78 100644
--- a/paimon-python/pypaimon/tests/data_evolution_deletion_vector_test.py
+++ b/paimon-python/pypaimon/tests/data_evolution_deletion_vector_test.py
@@ -16,6 +16,7 @@
# under the License.
import unittest
+from unittest.mock import Mock
import pyarrow as pa
@@ -32,7 +33,9 @@ from pypaimon.read.reader.concat_batch_reader import (
MergeAllBatchReader,
)
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.table.row.blob import Blob, BlobData
from pypaimon.table.row.generic_row import GenericRow
from pypaimon.table.source.deletion_file import DeletionFile
@@ -160,6 +163,39 @@ class DataEvolutionDeletionVectorTest(unittest.TestCase):
self.assertTrue(reader.deletion_vector().is_deleted(1))
self.assertFalse(reader.deletion_vector().is_deleted(2))
+ def
test_append_sliced_reader_maps_positions_to_original_file_offsets(self):
+ file = _file("slice.parquet", 0, 10, 1)
+ data_split = DataSplit(
+ files=[file],
+ partition=GenericRow([], []),
+ bucket=0,
+ raw_convertible=True,
+ data_deletion_files=None,
+ )
+ sliced_split = SlicedSplit(
+ data_split,
+ {"slice.parquet": (5, 10)},
+ )
+ split_read = RawFileSplitRead.__new__(RawFileSplitRead)
+ split_read.split = sliced_split
+ split_read._get_final_read_data_fields = Mock(return_value=[])
+ split_read.file_reader_supplier = Mock(
+ return_value=_OneBatchReader([5, 6, 7, 8, 9])
+ )
+ deletion_vector = BitmapDeletionVector()
+ deletion_vector.delete(7)
+
+ reader = split_read.raw_reader_supplier(
+ file,
+ dv_factory=lambda: deletion_vector,
+ )
+
+ self.assertEqual(
+ [5, 6, 8, 9],
+ reader.read_arrow_batch().column(0).to_pylist(),
+ )
+ self.assertIsNone(reader.read_arrow_batch())
+
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/scanner/chunk_shuffle_split_generator_test.py
b/paimon-python/pypaimon/tests/scanner/chunk_shuffle_split_generator_test.py
index ce69096d7c..735a0edab0 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
@@ -26,19 +26,22 @@ import os
import shutil
import tempfile
import unittest
-from unittest.mock import Mock
+from unittest.mock import Mock, patch
import pyarrow as pa
from pypaimon import CatalogFactory, Schema
+from pypaimon.deletionvectors.bitmap_deletion_vector import
BitmapDeletionVector
from pypaimon.globalindex.indexed_split import IndexedSplit
from pypaimon.manifest.schema.data_file_meta import DataFileMeta
from pypaimon.read.scanner.chunk_shuffle_split_generator import (
AppendChunkShuffleSplitGenerator,
DataEvolutionChunkShuffleSplitGenerator,
+ _LiveRowRangeSlicer,
)
from pypaimon.read.sliced_split import SlicedSplit
from pypaimon.read.split import DataSplit
+from pypaimon.table.source.deletion_file import DeletionFile
from pypaimon.utils.range import Range
@@ -63,33 +66,41 @@ def _mock_entry(partition_values, bucket, file_name,
row_count, file_size=1024):
return entry
-def _make_generator(seed, chunk_size, table=None):
+def _make_generator(seed, chunk_size, table=None, deletion_files_map=None):
if table is None:
table = _mock_table()
return AppendChunkShuffleSplitGenerator(
table,
target_split_size=128 * 1024 * 1024,
open_file_cost=4 * 1024 * 1024,
- deletion_files_map=None,
+ deletion_files_map=deletion_files_map,
seed=seed,
chunk_size=chunk_size,
)
-def _make_de_generator(seed, chunk_size, table=None):
+def _make_de_generator(seed, chunk_size, table=None, deletion_files_map=None):
if table is None:
table = _mock_table()
return DataEvolutionChunkShuffleSplitGenerator(
table,
target_split_size=128 * 1024 * 1024,
open_file_cost=4 * 1024 * 1024,
- deletion_files_map=None,
+ deletion_files_map=deletion_files_map,
seed=seed,
chunk_size=chunk_size,
)
-def _mock_de_entry(partition_values, bucket, file_name, first_row_id,
row_count, file_size=1024):
+def _mock_de_entry(
+ partition_values,
+ bucket,
+ file_name,
+ first_row_id,
+ row_count,
+ file_size=1024,
+ max_sequence_number=0,
+):
"""A DE-flavoured mock entry: file carries first_row_id and a real
Range so :meth:`row_id_range` and ``Range.overlaps`` work."""
entry = Mock()
@@ -101,6 +112,7 @@ def _mock_de_entry(partition_values, bucket, file_name,
first_row_id, row_count,
file.file_size = file_size
file.row_count = row_count
file.first_row_id = first_row_id
+ file.max_sequence_number = max_sequence_number
file.row_id_range = lambda f=first_row_id, c=row_count: Range(f, f + c - 1)
file.set_file_path = Mock()
entry.file = file
@@ -139,6 +151,65 @@ def _split_rows(split):
return split.row_count
+def _bitmap_deletion_vector(*positions):
+ deletion_vector = BitmapDeletionVector()
+ for position in positions:
+ deletion_vector.delete(position)
+ return deletion_vector
+
+
+class LiveRowRangeSlicerTest(unittest.TestCase):
+
+ def test_slices_by_live_rows_and_absorbs_deletion_runs(self):
+ slicer = _LiveRowRangeSlicer(10, iter([0, 3, 4, 9]))
+
+ first = slicer.take(3)
+ second = slicer.take(3)
+
+ self.assertEqual(
+ (
+ first.start_inclusive,
+ first.end_exclusive,
+ first.live_row_count,
+ ),
+ (0, 6, 3),
+ )
+ self.assertEqual(
+ (
+ second.start_inclusive,
+ second.end_exclusive,
+ second.live_row_count,
+ ),
+ (6, 10, 3),
+ )
+ self.assertEqual(first.to_closed_row_id_range(100), Range(100, 105))
+ self.assertIsNone(slicer.take(3))
+
+ 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),
+ )
+ self.assertIsNone(tail_slicer.take(1))
+
+ deleted_slicer = _LiveRowRangeSlicer(3, iter([0, 1, 2]))
+ self.assertIsNone(deleted_slicer.take(1))
+
+ def test_rejects_invalid_deletion_positions(self):
+ with self.assertRaisesRegex(ValueError, "outside physical row range"):
+ _LiveRowRangeSlicer(3, iter([-1]))
+
+ slicer = _LiveRowRangeSlicer(3, iter([0, 0]))
+ with self.assertRaisesRegex(ValueError, "strictly increasing"):
+ slicer.take(1)
+
+
class ChunkShuffleSplitGeneratorAlgoTest(unittest.TestCase):
def test_no_entries_returns_empty(self):
@@ -200,6 +271,133 @@ class
ChunkShuffleSplitGeneratorAlgoTest(unittest.TestCase):
self.assertIsInstance(splits[0], DataSplit)
self.assertEqual(_split_rows(splits[0]), 60)
+ def test_deletion_vector_slices_by_live_rows_and_is_attached(self):
+ entry = _mock_entry([], 0, 'f1', 10)
+ deletion_file = DeletionFile('dv.index', 10, 20, cardinality=4)
+ deletion_files_map = {((), 0): {'f1': deletion_file}}
+ deletion_vector = _bitmap_deletion_vector(0, 3, 4, 9)
+ gen = _make_generator(
+ seed=1,
+ chunk_size=3,
+ deletion_files_map=deletion_files_map,
+ )
+
+ with patch(
+
'pypaimon.read.scanner.chunk_shuffle_split_generator.DeletionVector.read',
+ return_value=deletion_vector,
+ ) as read:
+ splits = gen.create_splits([entry])
+
+ self.assertEqual(
+ sorted(s.shard_file_idx_map()['f1'] for s in splits),
+ [(0, 6), (6, 10)],
+ )
+ self.assertEqual(
+ sorted(s.merged_row_count() for s in splits),
+ [3, 3],
+ )
+ for split in splits:
+ self.assertEqual(
+ split.data_split().data_deletion_files,
+ [deletion_file],
+ )
+ read.assert_called_once_with(gen.table.file_io, deletion_file)
+
+ def test_whole_file_unknown_dv_cardinality_preserves_live_row_count(self):
+ entry = _mock_entry([], 0, 'f1', 10)
+ deletion_file = DeletionFile('dv.index', 10, 20, cardinality=None)
+ gen = _make_generator(
+ seed=1,
+ chunk_size=6,
+ deletion_files_map={((), 0): {'f1': deletion_file}},
+ )
+
+ with patch(
+
'pypaimon.read.scanner.chunk_shuffle_split_generator.DeletionVector.read',
+ return_value=_bitmap_deletion_vector(0, 3, 4, 9),
+ ) as read:
+ 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.assertEqual(splits[0].merged_row_count(), 6)
+ read.assert_called_once_with(gen.table.file_io, deletion_file)
+
+ def test_two_deletion_vector_files_share_one_live_row_chunk(self):
+ # Each file contributes 5 - 2 = 3 live rows, so both should fit
+ # exactly in one six-row chunk.
+ entries = [
+ _mock_entry([], 0, 'f1', 5),
+ _mock_entry([], 0, 'f2', 5),
+ ]
+ first_deletion_file = DeletionFile(
+ 'dv.index',
+ 10,
+ 20,
+ cardinality=2,
+ )
+ second_deletion_file = DeletionFile(
+ 'dv.index',
+ 30,
+ 20,
+ cardinality=2,
+ )
+ deletion_files_map = {
+ ((), 0): {
+ 'f1': first_deletion_file,
+ 'f2': second_deletion_file,
+ }
+ }
+ deletion_vectors = {
+ first_deletion_file: _bitmap_deletion_vector(1, 3),
+ second_deletion_file: _bitmap_deletion_vector(0, 4),
+ }
+ gen = _make_generator(
+ seed=1,
+ chunk_size=6,
+ deletion_files_map=deletion_files_map,
+ )
+
+ with patch(
+
'pypaimon.read.scanner.chunk_shuffle_split_generator.DeletionVector.read',
+ side_effect=lambda _, deletion_file:
deletion_vectors[deletion_file],
+ ) as read:
+ splits = gen.create_splits(entries)
+
+ self.assertEqual(len(splits), 1)
+ self.assertIsInstance(splits[0], DataSplit)
+ self.assertEqual(
+ [file.file_name for file in splits[0].files],
+ ['f1', 'f2'],
+ )
+ self.assertEqual(
+ splits[0].data_deletion_files,
+ [first_deletion_file, second_deletion_file],
+ )
+ self.assertEqual(read.call_count, 2)
+ self.assertEqual(
+ {mock_call.args[1] for mock_call in read.call_args_list},
+ {first_deletion_file, second_deletion_file},
+ )
+
+ def test_fully_deleted_file_does_not_create_a_segment(self):
+ entry = _mock_entry([], 0, 'f1', 5)
+ deletion_file = DeletionFile('dv.index', 10, 20, cardinality=5)
+ gen = _make_generator(
+ seed=1,
+ chunk_size=2,
+ deletion_files_map={((), 0): {'f1': deletion_file}},
+ )
+
+ with patch(
+
'pypaimon.read.scanner.chunk_shuffle_split_generator.DeletionVector.read',
+ ) as read:
+ self.assertEqual(gen.create_splits([entry]), [])
+
+ read.assert_not_called()
+
def test_deterministic_same_seed_same_order(self):
entries = [_mock_entry([], 0, f'f{i}', 100) for i in range(20)]
gen1 = _make_generator(seed=42, chunk_size=50)
@@ -453,12 +651,11 @@ class ChunkShuffleCompatibilityTest(unittest.TestCase):
with self.assertRaisesRegex(ValueError, "only supports append tables"):
scan.plan()
- def test_dv_table_rejected(self):
+ def test_dv_table_supported(self):
table = self._append_table('cs_dv',
options={'deletion-vectors.enabled': 'true'})
scan = table.new_read_builder().new_scan()
scan.with_chunk_shuffle(seed=1, chunk_size=100)
- with self.assertRaisesRegex(ValueError, "deletion vectors"):
- scan.plan()
+ self.assertEqual(scan.plan().splits(), [])
def test_with_slice_then_chunk_shuffle_rejected(self):
table = self._append_table('cs_slice')
@@ -616,6 +813,125 @@ class
DataEvolutionChunkShuffleAlgoTest(unittest.TestCase):
files = sorted(f.file_name for f in s.files)
self.assertEqual(files, ['g0.blob', 'g0.parquet'])
+ def test_deletion_vector_uses_anchor_and_slices_by_live_rows(self):
+ anchor = _mock_de_entry(
+ [],
+ 0,
+ 'anchor.parquet',
+ 100,
+ 10,
+ max_sequence_number=1,
+ )
+ blob = _mock_de_entry(
+ [],
+ 0,
+ 'field.blob',
+ 100,
+ 10,
+ max_sequence_number=2,
+ )
+ deletion_file = DeletionFile('dv.index', 10, 20, cardinality=4)
+ deletion_files_map = {((), 0): {'anchor.parquet': deletion_file}}
+ gen = _make_de_generator(
+ seed=1,
+ chunk_size=3,
+ deletion_files_map=deletion_files_map,
+ )
+
+ with patch(
+
'pypaimon.read.scanner.chunk_shuffle_split_generator.DeletionVector.read',
+ return_value=_bitmap_deletion_vector(1, 2, 7, 9),
+ ) as read:
+ splits = gen.create_splits([blob, anchor])
+
+ ranges = sorted(
+ (row_range.from_, row_range.to)
+ for split in splits
+ for row_range in split.row_ranges()
+ )
+ self.assertEqual(ranges, [(100, 104), (105, 109)])
+ self.assertEqual(
+ sorted(split.merged_row_count() for split in splits),
+ [3, 3],
+ )
+ for split in splits:
+ self.assertEqual(
+ sorted(f.file_name for f in split.files),
+ ['anchor.parquet', 'field.blob'],
+ )
+ deletion_by_name = dict(
+ (file.file_name, dv)
+ for file, dv in zip(split.files, split.data_deletion_files)
+ )
+ self.assertEqual(deletion_by_name['anchor.parquet'], deletion_file)
+ self.assertIsNone(deletion_by_name['field.blob'])
+ read.assert_called_once_with(gen.table.file_io, deletion_file)
+
+ def test_multiple_deletion_vector_groups_share_one_live_row_chunk(self):
+ # The groups contribute 3 live rows each despite spanning 5 and 6
+ # physical rows, so both should become segments of the same chunk.
+ entries = [
+ _mock_de_entry([], 0, 'g0.parquet', 100, 5),
+ _mock_de_entry([], 0, 'g1.parquet', 200, 6),
+ ]
+ first_deletion_file = DeletionFile(
+ 'dv.index',
+ 10,
+ 20,
+ cardinality=2,
+ )
+ second_deletion_file = DeletionFile(
+ 'dv.index',
+ 30,
+ 20,
+ cardinality=3,
+ )
+ deletion_files_map = {
+ ((), 0): {
+ 'g0.parquet': first_deletion_file,
+ 'g1.parquet': second_deletion_file,
+ }
+ }
+ deletion_vectors = {
+ first_deletion_file: _bitmap_deletion_vector(1, 3),
+ second_deletion_file: _bitmap_deletion_vector(0, 2, 5),
+ }
+ gen = _make_de_generator(
+ seed=1,
+ chunk_size=6,
+ deletion_files_map=deletion_files_map,
+ )
+
+ with patch(
+
'pypaimon.read.scanner.chunk_shuffle_split_generator.DeletionVector.read',
+ side_effect=lambda _, deletion_file:
deletion_vectors[deletion_file],
+ ) as read:
+ splits = gen.create_splits(entries)
+
+ self.assertEqual(len(splits), 1)
+ self.assertIsInstance(splits[0], IndexedSplit)
+ self.assertEqual(
+ [
+ (row_range.from_, row_range.to)
+ for row_range in splits[0].row_ranges()
+ ],
+ [(100, 104), (200, 205)],
+ )
+ self.assertEqual(
+ [file.file_name for file in splits[0].files],
+ ['g0.parquet', 'g1.parquet'],
+ )
+ self.assertEqual(
+ splits[0].data_deletion_files,
+ [first_deletion_file, second_deletion_file],
+ )
+ self.assertEqual(splits[0].merged_row_count(), 6)
+ self.assertEqual(read.call_count, 2)
+ self.assertEqual(
+ {mock_call.args[1] for mock_call in read.call_args_list},
+ {first_deletion_file, second_deletion_file},
+ )
+
def test_deterministic_same_seed(self):
entries = [_mock_de_entry([], 0, f'g{i:02d}.parquet', i * 100, 100)
for i in range(20)]
gen1 = _make_de_generator(seed=42, chunk_size=100)
@@ -705,19 +1021,22 @@ class
DataEvolutionChunkShuffleEndToEndTest(unittest.TestCase):
def tearDownClass(cls):
shutil.rmtree(cls.tempdir, ignore_errors=True)
- def _create_de_table(self, name):
+ def _create_de_table(self, name, deletion_vectors_enabled=False):
pa_schema = pa.schema([
('id', pa.int32()),
('value', pa.string()),
('payload', pa.large_binary()),
])
+ options = {
+ 'row-tracking.enabled': 'true',
+ 'data-evolution.enabled': 'true',
+ 'blob.target-file-size': '1 b',
+ }
+ if deletion_vectors_enabled:
+ options['deletion-vectors.enabled'] = 'true'
schema = Schema.from_pyarrow_schema(
pa_schema,
- options={
- 'row-tracking.enabled': 'true',
- 'data-evolution.enabled': 'true',
- 'blob.target-file-size': '1 b',
- },
+ options=options,
)
identifier = f'default.{name}'
self.catalog.create_table(identifier, schema, False)
@@ -744,6 +1063,14 @@ class
DataEvolutionChunkShuffleEndToEndTest(unittest.TestCase):
tc.close()
return commit_messages
+ @staticmethod
+ def _delete_by_row_id(table, row_ids):
+ write_builder = table.new_batch_write_builder()
+ commit_messages = write_builder.new_update().delete_by_row_id(row_ids)
+ table_commit = write_builder.new_commit()
+ table_commit.commit(commit_messages)
+ table_commit.close()
+
def _assert_commit_has_main_and_multiple_blob_files(self, commit_messages):
all_files = [f for msg in commit_messages for f in msg.new_files]
main_files = [f for f in all_files if not
DataFileMeta.is_blob_file(f.file_name)]
@@ -816,6 +1143,41 @@ class
DataEvolutionChunkShuffleEndToEndTest(unittest.TestCase):
for worker in range(4):
self.assertEqual(plan_sigs(worker), plan_sigs(worker))
+ def test_python_delete_then_chunk_shuffle_read(self):
+ table, pa_schema = self._create_de_table(
+ 'cs_de_python_delete',
+ deletion_vectors_enabled=True,
+ )
+ commit_messages = self._commit_full_rows(
+ table,
+ pa_schema,
+ list(range(12)),
+ )
+ self._assert_commit_has_main_and_multiple_blob_files(commit_messages)
+ self._delete_by_row_id(table, [0, 2, 3, 7, 11])
+
+ read_builder = table.new_read_builder()
+ splits = read_builder.new_scan() \
+ .with_chunk_shuffle(seed=42, chunk_size=3) \
+ .plan() \
+ .splits()
+ self._assert_splits_include_blob_files(splits)
+
+ table_read = read_builder.new_read()
+ rows_per_split = [
+ table_read.to_arrow([split]).num_rows
+ for split in splits
+ ]
+ self.assertEqual(sorted(rows_per_split), [1, 3, 3])
+
+ actual = table_read.to_arrow(splits).sort_by('id')
+ expected_ids = [1, 4, 5, 6, 8, 9, 10]
+ self.assertEqual(actual.column('id').to_pylist(), expected_ids)
+ self.assertEqual(
+ actual.column('payload').to_pylist(),
+ self._payloads(expected_ids),
+ )
+
if __name__ == '__main__':
unittest.main()