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 8c61c180ec [python] Cache decoded BLOB indexes (#9547)
8c61c180ec is described below
commit 8c61c180ec0fd1cbcccc19a36bbf40c10476b8ca
Author: Yann Byron <[email protected]>
AuthorDate: Thu Sep 3 10:13:49 2026 +0800
[python] Cache decoded BLOB indexes (#9547)
---
.../pypaimon/read/reader/format_blob_reader.py | 45 +++++++++++++++------
paimon-python/pypaimon/tests/blob_test.py | 47 +++++++++++++++++++++-
2 files changed, 79 insertions(+), 13 deletions(-)
diff --git a/paimon-python/pypaimon/read/reader/format_blob_reader.py
b/paimon-python/pypaimon/read/reader/format_blob_reader.py
index 8446bfcb84..4605da5804 100644
--- a/paimon-python/pypaimon/read/reader/format_blob_reader.py
+++ b/paimon-python/pypaimon/read/reader/format_blob_reader.py
@@ -16,10 +16,12 @@
# under the License.
import struct
+from threading import Lock
from typing import List, Optional, Any, Iterator, BinaryIO
import pyarrow as pa
import pyarrow.dataset as ds
+from cachetools import LRUCache
from pyarrow import RecordBatch
from pypaimon.common.delta_varint_compressor import DeltaVarintCompressor
@@ -43,6 +45,24 @@ from pypaimon.table.row.generic_row import GenericRow
from pypaimon.table.row.row_kind import RowKind
+_BLOB_INDEX_CACHE = LRUCache(maxsize=16)
+_BLOB_INDEX_CACHE_LOCK = Lock()
+
+
+def _decode_blob_index(index_bytes):
+ """Decode BLOB lengths and their relative file offsets."""
+ blob_lengths = tuple(DeltaVarintCompressor.decompress(index_bytes))
+ blob_offsets = []
+ offset = 0
+ for length in blob_lengths:
+ if length < 0:
+ blob_offsets.append(-1)
+ else:
+ blob_offsets.append(offset)
+ offset += length
+ return blob_lengths, tuple(blob_offsets)
+
+
class FormatBlobReader(RecordBatchReader):
NULL_LENGTH = -1
PLACE_HOLDER_LENGTH = -2
@@ -350,6 +370,14 @@ class FormatBlobReader(RecordBatchReader):
)
return
+ with _BLOB_INDEX_CACHE_LOCK:
+ cached_index = _BLOB_INDEX_CACHE.get(self.file_path)
+ if cached_index is not None:
+ blob_lengths, blob_offsets = cached_index
+ self.blob_lengths = list(blob_lengths)
+ self.blob_offsets = list(blob_offsets)
+ return
+
f = self._input_stream
# Seek to header: last 5 bytes
@@ -373,18 +401,11 @@ class FormatBlobReader(RecordBatchReader):
if len(index_bytes) != index_length:
raise IOError("Invalid blob file: cannot read index")
- # Decompress blob lengths and compute offsets
- blob_lengths = DeltaVarintCompressor.decompress(index_bytes)
- blob_offsets = []
- offset = 0
- for length in blob_lengths:
- if length < 0:
- blob_offsets.append(-1)
- else:
- blob_offsets.append(offset)
- offset += length
- self.blob_lengths = blob_lengths
- self.blob_offsets = blob_offsets
+ blob_lengths, blob_offsets = _decode_blob_index(index_bytes)
+ with _BLOB_INDEX_CACHE_LOCK:
+ _BLOB_INDEX_CACHE[self.file_path] = blob_lengths, blob_offsets
+ self.blob_lengths = list(blob_lengths)
+ self.blob_offsets = list(blob_offsets)
def _apply_row_indices(self, row_indices: Optional[Any]) -> None:
if row_indices is None:
diff --git a/paimon-python/pypaimon/tests/blob_test.py
b/paimon-python/pypaimon/tests/blob_test.py
index 544e00cdee..c0ba12bdfe 100644
--- a/paimon-python/pypaimon/tests/blob_test.py
+++ b/paimon-python/pypaimon/tests/blob_test.py
@@ -28,7 +28,7 @@ import unittest
import zlib
from decimal import Decimal
from pathlib import Path
-from unittest.mock import patch
+from unittest.mock import Mock, patch
import pyarrow as pa
@@ -1562,6 +1562,51 @@ class BlobEndToEndTest(unittest.TestCase):
finally:
reader.close()
+ def test_blob_readers_reuse_index_by_path(self):
+ from pypaimon.read.reader.format_blob_reader import _BLOB_INDEX_CACHE
+
+ field = DataField(0, "blob_field", AtomicType("BLOB"))
+ path = os.path.join(self.temp_dir, "cached-index.blob")
+ file_io = LocalFileIO(self.temp_dir, Options({}))
+ self._write_single_blob(path, field, b"cached-value")
+ _BLOB_INDEX_CACHE.clear()
+ input_streams = []
+ new_input_stream = file_io.new_input_stream
+
+ def counting_input_stream(file_path):
+ stream = Mock(wraps=new_input_stream(file_path))
+ input_streams.append(stream)
+ return stream
+
+ try:
+ with patch.object(
+ file_io,
+ "new_input_stream",
+ side_effect=counting_input_stream,
+ ), patch.object(
+ DeltaVarintCompressor,
+ "decompress",
+ wraps=DeltaVarintCompressor.decompress,
+ ) as decompress:
+ for _ in range(2):
+ reader = FormatBlobReader(
+ file_io,
+ path,
+ [field.name],
+ [field],
+ None,
+ True,
+ )
+ reader.close()
+
+ self.assertEqual(1, decompress.call_count)
+ self.assertEqual(
+ [2, 0],
+ [stream.read.call_count for stream in input_streams],
+ )
+ finally:
+ _BLOB_INDEX_CACHE.clear()
+
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")