| # Licensed to the Apache Software Foundation (ASF) under one |
| # or more contributor license agreements. See the NOTICE file |
| # distributed with this work for additional information |
| # regarding copyright ownership. The ASF licenses this file |
| # to you under the Apache License, Version 2.0 (the |
| # "License"); you may not use this file except in compliance |
| # with the License. You may obtain a copy of the License at |
| # |
| # http://www.apache.org/licenses/LICENSE-2.0 |
| # |
| # Unless required by applicable law or agreed to in writing, |
| # software distributed under the License is distributed on an |
| # "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY |
| # KIND, either express or implied. See the License for the |
| # specific language governing permissions and limitations |
| # under the License. |
| |
| import unittest |
| from unittest.mock import Mock |
| |
| import pyarrow as pa |
| |
| from pypaimon.deletionvectors.apply_deletion_vector_reader import ( |
| ApplyDeletionVectorReader, |
| PositionMappedDeletionVector, |
| ) |
| 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.read.reader.concat_batch_reader import ( |
| BlobFallbackBatchReader, |
| DataEvolutionMergeReader, |
| 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 |
| from pypaimon.utils.range import Range |
| from pypaimon.utils.data_evolution_utils import retrieve_anchor_file |
| |
| |
| class _OneBatchReader(RecordBatchReader): |
| def __init__(self, values): |
| self._batch = pa.record_batch([pa.array(values, type=pa.int64())], names=["v"]) |
| self._returned = False |
| |
| def read_arrow_batch(self): |
| if self._returned: |
| return None |
| self._returned = True |
| return self._batch |
| |
| def close(self): |
| pass |
| |
| |
| class _BlobFallbackBatchReaderForTest(BlobFallbackBatchReader): |
| def __init__( |
| self, files, values_by_file_name, row_ranges=None, deletion_vector=None, |
| batch_size=1024 |
| ): |
| super().__init__( |
| [(file, lambda: None) for file in files], |
| "blob_col", |
| pa.binary(), |
| row_ranges=row_ranges, |
| blob_as_descriptor=False, |
| deletion_vector=deletion_vector, |
| batch_size=batch_size, |
| ) |
| self._values_by_file_name = values_by_file_name |
| |
| def _read_blob_values(self, state, batch_row_ids): |
| values = self._values_by_file_name[state.file.file_name] |
| return { |
| row_id: values[pos] |
| for pos, row_id in self._selected_positions_and_row_ids( |
| state, batch_row_ids |
| ) |
| } |
| |
| |
| def _file(name, first_row_id, row_count, max_sequence_number): |
| empty_row = GenericRow([], []) |
| return DataFileMeta( |
| file_name=name, |
| file_size=1, |
| row_count=row_count, |
| min_key=empty_row, |
| max_key=empty_row, |
| key_stats=SimpleStats.empty_stats(), |
| value_stats=SimpleStats.empty_stats(), |
| min_sequence_number=max_sequence_number, |
| max_sequence_number=max_sequence_number, |
| schema_id=0, |
| level=0, |
| extra_files=[], |
| first_row_id=first_row_id, |
| ) |
| |
| |
| class DataEvolutionDeletionVectorTest(unittest.TestCase): |
| def test_retrieve_anchor_file_uses_oldest_normal_file(self): |
| files = [ |
| _file("field-2.blob", 0, 5, 1), |
| _file("normal-b.parquet", 0, 5, 1), |
| _file("normal-a.parquet", 0, 5, 1), |
| _file("newer.parquet", 0, 5, 2), |
| ] |
| |
| self.assertEqual("normal-a.parquet", retrieve_anchor_file(files).file_name) |
| |
| def test_data_evolution_merged_row_count_subtracts_deletion_vectors(self): |
| split = DataSplit( |
| files=[ |
| _file("anchor-0.parquet", 0, 5, 1), |
| _file("blob-0.blob", 0, 5, 2), |
| _file("anchor-5.parquet", 5, 5, 3), |
| ], |
| partition=GenericRow([], []), |
| bucket=0, |
| raw_convertible=False, |
| data_deletion_files=[ |
| DeletionFile("dv", 0, 1, cardinality=2), |
| None, |
| DeletionFile("dv", 1, 1, cardinality=1), |
| ], |
| ) |
| |
| self.assertEqual(7, split.merged_row_count()) |
| |
| def test_data_evolution_merged_row_count_unknown_without_cardinality(self): |
| split = DataSplit( |
| files=[_file("anchor.parquet", 0, 5, 1)], |
| partition=GenericRow([], []), |
| bucket=0, |
| raw_convertible=False, |
| data_deletion_files=[DeletionFile("dv", 0, 1, cardinality=None)], |
| ) |
| |
| self.assertIsNone(split.merged_row_count()) |
| |
| def test_apply_deletion_vector_reader_uses_mapped_deletion_vector(self): |
| deletion_vector = BitmapDeletionVector() |
| deletion_vector.delete(12) |
| mapped_dv = PositionMappedDeletionVector( |
| deletion_vector, |
| file_offset=10, |
| row_positions=[0, 2, 4], |
| ) |
| |
| reader = ApplyDeletionVectorReader( |
| _OneBatchReader([0, 2, 4]), |
| mapped_dv, |
| ) |
| |
| batch = reader.read_arrow_batch() |
| self.assertEqual([0, 4], batch.column(0).to_pylist()) |
| 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) |
| deletion_vector.delete(1) |
| |
| field_reader = MergeAllBatchReader([ |
| lambda: ApplyDeletionVectorReader( |
| _OneBatchReader([0, 1]), |
| deletion_vector, |
| ) |
| ]) |
| reader = DataEvolutionMergeReader( |
| row_offsets=[0], |
| field_offsets=[0], |
| readers=[field_reader], |
| schema=pa.schema([pa.field("v", pa.int64())]), |
| ) |
| |
| self.assertIsNone(reader.read_arrow_batch()) |
| |
| def test_blob_fallback_batch_reader_applies_deletion_vector(self): |
| files = [ |
| _file("blob-old.blob", 0, 5, 1), |
| _file("blob-new.blob", 0, 5, 2), |
| ] |
| deletion_vector = BitmapDeletionVector() |
| deletion_vector.delete(1) |
| deletion_vector.delete(4) |
| |
| reader = _BlobFallbackBatchReaderForTest( |
| files, |
| { |
| "blob-old.blob": [ |
| BlobData(b"old-0"), |
| BlobData(b"old-1"), |
| BlobData(b"old-2"), |
| BlobData(b"old-3"), |
| BlobData(b"old-4"), |
| ], |
| "blob-new.blob": [ |
| Blob.PLACE_HOLDER, |
| BlobData(b"new-1"), |
| BlobData(b"new-2"), |
| Blob.PLACE_HOLDER, |
| BlobData(b"new-4"), |
| ], |
| }, |
| deletion_vector=(Range(0, 4), deletion_vector), |
| ) |
| |
| batch = reader.read_arrow_batch() |
| self.assertEqual( |
| [b"old-0", b"new-2", b"old-3"], |
| batch.column(0).to_pylist(), |
| ) |
| self.assertIsNone(reader.read_arrow_batch()) |
| |
| def test_blob_fallback_batch_reader_does_not_eof_on_row_id_gap(self): |
| reader = _BlobFallbackBatchReaderForTest( |
| [ |
| _file("blob-left.blob", 0, 1, 1), |
| _file("blob-right.blob", 2, 1, 1), |
| ], |
| { |
| "blob-left.blob": [BlobData(b"left")], |
| "blob-right.blob": [BlobData(b"right")], |
| }, |
| row_ranges=[Range(0, 2)], |
| batch_size=1, |
| ) |
| |
| values = [] |
| for batch in iter(reader.read_arrow_batch, None): |
| values.extend(batch.column(0).to_pylist()) |
| |
| self.assertEqual([b"left", b"right"], values) |
| |
| def test_blob_fallback_batch_reader_skips_states_outside_batch(self): |
| class _CountingBlobFallbackBatchReader(_BlobFallbackBatchReaderForTest): |
| def __init__(self, files, values_by_file_name, batch_size): |
| super().__init__(files, values_by_file_name, batch_size=batch_size) |
| self.selected_calls = [] |
| |
| def _selected_positions_and_row_ids(self, state, batch_row_ids): |
| self.selected_calls.append( |
| (state.file.file_name, list(batch_row_ids)) |
| ) |
| return BlobFallbackBatchReader._selected_positions_and_row_ids( |
| state, batch_row_ids |
| ) |
| |
| reader = _CountingBlobFallbackBatchReader( |
| [ |
| _file("blob-0.blob", 0, 1, 1), |
| _file("blob-10.blob", 10, 1, 1), |
| _file("blob-20.blob", 20, 1, 1), |
| ], |
| { |
| "blob-0.blob": [BlobData(b"0")], |
| "blob-10.blob": [BlobData(b"10")], |
| "blob-20.blob": [BlobData(b"20")], |
| }, |
| batch_size=1, |
| ) |
| try: |
| self.assertEqual([b"0"], reader.read_arrow_batch().column(0).to_pylist()) |
| self.assertEqual([("blob-0.blob", [0])], reader.selected_calls) |
| |
| self.assertEqual([b"10"], reader.read_arrow_batch().column(0).to_pylist()) |
| self.assertEqual( |
| [("blob-0.blob", [0]), ("blob-10.blob", [10])], |
| reader.selected_calls, |
| ) |
| finally: |
| reader.close() |
| |
| def test_blob_fallback_batch_reader_avoids_full_row_id_materialization(self): |
| class _LazyBlobFallbackBatchReaderForTest(BlobFallbackBatchReader): |
| def _read_blob_values(self, state, batch_row_ids): |
| return { |
| row_id: BlobData(str(row_id).encode("utf-8")) |
| for row_id in batch_row_ids |
| } |
| |
| reader = _LazyBlobFallbackBatchReaderForTest( |
| [(_file("large.blob", 0, 1_000_000, 1), lambda: None)], |
| "blob_col", |
| pa.binary(), |
| batch_size=2, |
| ) |
| try: |
| self.assertNotIn("_target_row_ids", reader.__dict__) |
| for state in reader._file_states: |
| self.assertFalse(hasattr(state, "selected_row_ids")) |
| self.assertFalse(hasattr(state, "row_id_to_pos")) |
| self.assertFalse(hasattr(state, "reader_uses_selected_positions")) |
| |
| batch = reader.read_arrow_batch() |
| self.assertEqual([b"0", b"1"], batch.column(0).to_pylist()) |
| finally: |
| reader.close() |
| |
| def test_blob_fallback_batch_reader_advances_sparse_range_cursor(self): |
| reader = _BlobFallbackBatchReaderForTest( |
| [_file("sparse.blob", 0, 100, 1)], |
| { |
| "sparse.blob": [ |
| BlobData(b"0"), |
| BlobData(b"10"), |
| BlobData(b"20"), |
| BlobData(b"30"), |
| ], |
| }, |
| row_ranges=[Range(0, 0), Range(10, 10), Range(20, 20), Range(30, 30)], |
| batch_size=2, |
| ) |
| try: |
| state = reader._file_states[0] |
| |
| first = reader.read_arrow_batch() |
| self.assertEqual([b"0", b"10"], first.column(0).to_pylist()) |
| self.assertEqual(1, state.selected_range_index) |
| self.assertEqual(1, state.selected_position_base) |
| |
| second = reader.read_arrow_batch() |
| self.assertEqual([b"20", b"30"], second.column(0).to_pylist()) |
| self.assertEqual(3, state.selected_range_index) |
| self.assertEqual(3, state.selected_position_base) |
| self.assertIsNone(reader.read_arrow_batch()) |
| finally: |
| reader.close() |
| |
| def test_data_evolution_merge_reader_aligns_blob_with_row_ranges_and_dv(self): |
| row_ranges = [Range(1, 4)] |
| deletion_vector = BitmapDeletionVector() |
| deletion_vector.delete(2) |
| deletion_vector.delete(4) |
| |
| normal_reader = ApplyDeletionVectorReader( |
| _OneBatchReader([1, 2, 3, 4]), |
| PositionMappedDeletionVector( |
| deletion_vector, |
| file_offset=0, |
| row_positions=[1, 2, 3, 4], |
| ), |
| ) |
| blob_reader = _BlobFallbackBatchReaderForTest( |
| [ |
| _file("blob-old.blob", 0, 6, 1), |
| _file("blob-new.blob", 0, 6, 2), |
| ], |
| { |
| "blob-old.blob": [ |
| BlobData(b"old-1"), |
| BlobData(b"old-2"), |
| BlobData(b"old-3"), |
| BlobData(b"old-4"), |
| ], |
| "blob-new.blob": [ |
| Blob.PLACE_HOLDER, |
| BlobData(b"new-2"), |
| BlobData(b"new-3"), |
| BlobData(b"new-4"), |
| ], |
| }, |
| row_ranges=row_ranges, |
| deletion_vector=(Range(0, 5), deletion_vector), |
| ) |
| reader = DataEvolutionMergeReader( |
| row_offsets=[0, 1], |
| field_offsets=[0, 0], |
| readers=[normal_reader, blob_reader], |
| schema=pa.schema([pa.field("id", pa.int64()), pa.field("blob_col", pa.binary())]), |
| ) |
| |
| batch = reader.read_arrow_batch() |
| self.assertEqual([1, 3], batch.column(0).to_pylist()) |
| self.assertEqual([b"old-1", b"new-3"], batch.column(1).to_pylist()) |
| self.assertIsNone(reader.read_arrow_batch()) |
| |
| |
| if __name__ == "__main__": |
| unittest.main() |