blob: 4ef3748f78c120dd30a067d145e2676f010e0367 [file]
# 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()