diff --git a/paimon-python/pypaimon/read/reader/format_blob_reader.py b/paimon-python/pypaimon/read/reader/format_blob_reader.py index 5b826c584202..897a81f67a60 100644 --- a/paimon-python/pypaimon/read/reader/format_blob_reader.py +++ b/paimon-python/pypaimon/read/reader/format_blob_reader.py @@ -46,7 +46,7 @@ class FormatBlobReader(RecordBatchReader): def __init__(self, file_io: FileIO, file_path: str, read_fields: List[str], full_fields: List[DataField], push_down_predicate: Any, blob_as_descriptor: bool, batch_size: int = 1024, row_indices: Optional[Any] = None, - blob_parallelism: int = 1): + blob_parallelism: int = 1, file_size: Optional[int] = None): self._file_io = file_io self._file_path = file_path self._push_down_predicate = push_down_predicate @@ -63,7 +63,11 @@ def __init__(self, file_io: FileIO, file_path: str, read_fields: List[str], self._blob_iterator = None self._current_batch = None try: - self._file_size = file_io.get_file_size(file_path) + self._file_size = ( + file_size + if file_size is not None and file_size > 0 + else file_io.get_file_size(file_path) + ) self._input_stream = file_io.new_input_stream(file_path) self._read_index() self._apply_row_indices(row_indices) diff --git a/paimon-python/pypaimon/read/split_read.py b/paimon-python/pypaimon/read/split_read.py index 2f6f09665f9b..1085802dd508 100644 --- a/paimon-python/pypaimon/read/split_read.py +++ b/paimon-python/pypaimon/read/split_read.py @@ -336,7 +336,8 @@ def file_reader_supplier(self, file: DataFileMeta, for_merge_read: bool, self.read_fields, read_arrow_predicate, blob_as_descriptor, batch_size=batch_size, row_indices=row_indices, - blob_parallelism=blob_parallelism) + blob_parallelism=blob_parallelism, + file_size=file.file_size) elif file_format == CoreOptions.FILE_FORMAT_LANCE: if has_nested: raise NotImplementedError( @@ -1460,6 +1461,7 @@ def _create_raw_blob_file_reader( batch_size=self.table.options.read_batch_size(), row_indices=row_indices, blob_parallelism=blob_parallelism, + file_size=file.file_size, ) def _split_field_bunches(self, need_merge_files: List[DataFileMeta]) -> List[FieldBunch]: diff --git a/paimon-python/pypaimon/tests/blob_test.py b/paimon-python/pypaimon/tests/blob_test.py index f661faebf281..19aacce9e044 100644 --- a/paimon-python/pypaimon/tests/blob_test.py +++ b/paimon-python/pypaimon/tests/blob_test.py @@ -16,6 +16,7 @@ # under the License. import datetime +import inspect import io import os import shutil @@ -46,11 +47,16 @@ class MockFileIO: """Mock FileIO for testing.""" - def __init__(self, file_io: FileIO): + def __init__(self, file_io: FileIO, fail_on_file_size=False): self._file_io = file_io + self.fail_on_file_size = fail_on_file_size + self.file_size_calls = 0 def get_file_size(self, path: str) -> int: """Get file size.""" + self.file_size_calls += 1 + if self.fail_on_file_size: + raise AssertionError("get_file_size should not be called") return self._file_io.get_file_size(path) def new_input_stream(self, path): @@ -1354,6 +1360,16 @@ def tearDown(self): except OSError: pass + @staticmethod + def _write_single_blob(path, field, value): + from pypaimon.write.blob_format_writer import BlobFormatWriter + + with open(path, 'wb') as output: + writer = BlobFormatWriter(output) + writer.add_element(GenericRow( + [BlobData(value)], [field], RowKind.INSERT)) + writer.close() + def test_blob_end_to_end(self): # Set up file I/O file_io = LocalFileIO(self.temp_dir, Options({})) @@ -1443,6 +1459,112 @@ def new_input_stream(self, path): self.assertEqual(counting_file_io.input_stream_count, 1) reader.close() + def test_blob_reader_uses_provided_file_size(self): + field = DataField(0, "blob_field", AtomicType("BLOB")) + path = os.path.join(self.temp_dir, "provided-size.blob") + file_io = LocalFileIO(self.temp_dir, Options({})) + self._write_single_blob(path, field, b"value") + + reader = FormatBlobReader( + MockFileIO(file_io, fail_on_file_size=True), + path, + [field.name], + [field], + None, + False, + file_size=os.path.getsize(path), + ) + try: + self.assertEqual( + [b"value"], reader.read_arrow_batch().column(0).to_pylist()) + finally: + reader.close() + + def test_blob_reader_falls_back_to_file_size_lookup(self): + field = DataField(0, "blob_field", AtomicType("BLOB")) + path = os.path.join(self.temp_dir, "fallback-size.blob") + file_io = LocalFileIO(self.temp_dir, Options({})) + self._write_single_blob(path, field, b"value") + + for file_size in [None, 0, -1]: + with self.subTest(file_size=file_size): + counting_file_io = MockFileIO(file_io) + reader = FormatBlobReader( + counting_file_io, + path, + [field.name], + [field], + None, + False, + file_size=file_size, + ) + try: + self.assertEqual( + [b"value"], + reader.read_arrow_batch().column(0).to_pylist()) + self.assertEqual(1, counting_file_io.file_size_calls) + finally: + reader.close() + + def test_split_read_passes_blob_file_size(self): + from pypaimon.read.split import DataSplit + from pypaimon.read.split_read import ( + DataEvolutionSplitRead, + RawFileSplitRead, + ) + + fields = [ + DataField(0, "id", AtomicType("INT")), + DataField(1, "blob_field", AtomicType("BLOB")), + ] + self.catalog.create_table( + "test_db.blob_file_size_forwarding", + Schema(fields, options={ + "blob.file-format": "blob", + "data-evolution.enabled": "true", + "row-tracking.enabled": "true", + }), + False, + ) + table = self.catalog.get_table("test_db.blob_file_size_forwarding") + field = fields[1] + file = DataFileMeta( + file_name="data.blob", + file_size=123, + row_count=1, + min_key=None, + max_key=None, + key_stats=None, + value_stats=None, + min_sequence_number=0, + max_sequence_number=0, + schema_id=0, + level=0, + extra_files=[], + first_row_id=0, + write_cols=[field.name], + file_path="data.blob", + ) + split = DataSplit([file], GenericRow([], []), 0) + raw_read = RawFileSplitRead(table, None, [field], split, False) + evolution_read = DataEvolutionSplitRead( + table, None, [field], split, False) + + with patch("pypaimon.read.split_read.FormatBlobReader") as reader_cls: + def assert_file_size(): + args, kwargs = reader_cls.call_args + arguments = inspect.signature(FormatBlobReader).bind_partial( + *args, **kwargs + ).arguments + self.assertEqual(123, arguments.get("file_size")) + + raw_read.file_reader_supplier(file, False, [field.name], False) + assert_file_size() + + reader_cls.reset_mock() + evolution_read._create_raw_blob_file_reader(file, [field.name]) + assert_file_size() + def test_blob_reader_row_indices_pushdown(self): file_io = LocalFileIO(self.temp_dir, Options({})) blob_field_name = "blob_field"