From 1022af4344fbe9bbbb8b6ec75f843cd5da177921 Mon Sep 17 00:00:00 2001 From: "xiaohongbo.xhb" Date: Sat, 26 Sep 2026 18:21:55 -0700 Subject: [PATCH 01/10] [python] Cache BLOB metadata ranges separately from data --- docs/docs/program-api/file-cache.mdx | 4 +- docs/docs/pypaimon/blob.md | 8 + .../pypaimon/common/options/core_options.py | 4 +- .../pypaimon/filesystem/caching_file_io.py | 59 ++++++-- .../read/reader/format_blob_reader.py | 25 ++-- .../pypaimon/tests/blob_meta_cache_test.py | 141 ++++++++++++++++++ .../pypaimon/tests/caching_file_io_test.py | 2 +- paimon-python/pypaimon/utils/file_type.py | 5 +- 8 files changed, 219 insertions(+), 29 deletions(-) create mode 100644 paimon-python/pypaimon/tests/blob_meta_cache_test.py diff --git a/docs/docs/program-api/file-cache.mdx b/docs/docs/program-api/file-cache.mdx index 8c58a398bcab..5f47cf653c59 100644 --- a/docs/docs/program-api/file-cache.mdx +++ b/docs/docs/program-api/file-cache.mdx @@ -52,7 +52,9 @@ The cache classifies files by path. Set `local-cache.whitelist` to select the fi | DATA | data | Data files (ORC, Parquet, etc.) | No | | FILE_INDEX | file-index | Data-file level bloom filter, bitmap | No | -The default whitelist is `meta,global-index`. The mutable `LATEST` and `EARLIEST` hint files bypass +The default whitelist is `meta,global-index`. The Python reader additionally includes +`blob-meta`, which caches exact BLOB metadata ranges rather than whole data blocks; see +[Python BLOB metadata cache](../pypaimon/blob.md#blob-metadata-cache). The mutable `LATEST` and `EARLIEST` hint files bypass the cache even though they are classified as metadata. ## Enable Cache diff --git a/docs/docs/pypaimon/blob.md b/docs/docs/pypaimon/blob.md index 56126ef5b2da..eed1f8c9f091 100644 --- a/docs/docs/pypaimon/blob.md +++ b/docs/docs/pypaimon/blob.md @@ -176,6 +176,14 @@ The factory auto-dispatches based on the bytes content (`BLOBDESC`, `VIDEOFRM`, or blob-view magic header). This mirrors Java's `Blob.fromBytes(...)`. +## BLOB metadata cache + +With `local-cache.enabled=true`, the Python reader's default whitelist is +`meta,global-index,blob-meta`. `blob-meta` caches exact ranges for BLOB footers, +row indexes, ARRAY headers/indexes, and MAP headers/keys/indexes, excluding value +bodies. It shares the local cache's memory or disk budget. Add `data` to cache +normal data blocks as well. This option applies to the Python reader, not Rust native reads. + ## See Also - [Blob Storage](../multimodal-table/blob) — concept, storage modes, diff --git a/paimon-python/pypaimon/common/options/core_options.py b/paimon-python/pypaimon/common/options/core_options.py index a1fb6b2abfbf..4c3ab71f36da 100644 --- a/paimon-python/pypaimon/common/options/core_options.py +++ b/paimon-python/pypaimon/common/options/core_options.py @@ -1138,10 +1138,10 @@ class CoreOptions: LOCAL_CACHE_WHITELIST: ConfigOption[str] = ( ConfigOptions.key("local-cache.whitelist") .string_type() - .default_value("meta,global-index") + .default_value("meta,global-index,blob-meta") .with_description( "Comma-separated list of file types to cache. " - "Supported values: meta, global-index, bucket-index, data, file-index." + "Supported values: meta, global-index, bucket-index, data, file-index, blob-meta." ) ) diff --git a/paimon-python/pypaimon/filesystem/caching_file_io.py b/paimon-python/pypaimon/filesystem/caching_file_io.py index 3fad543d8d8c..182eb9067ee9 100644 --- a/paimon-python/pypaimon/filesystem/caching_file_io.py +++ b/paimon-python/pypaimon/filesystem/caching_file_io.py @@ -22,21 +22,24 @@ at block granularity. If a cache directory is configured, disk cache is used; otherwise an in-memory LRU cache is used. Files are classified by FileType and only cacheable types in the whitelist are cached; others are read directly from -the delegate FileIO. +the delegate FileIO. BLOB metadata can be cached as exact reader-selected ranges. """ import hashlib import os import threading from collections import OrderedDict -from typing import Optional +from typing import Optional, Tuple, Union from pypaimon.common.file_io import FileIO, supports_pread, pread from pypaimon.utils.file_type import FileType +_CacheEntryKey = Union[int, Tuple[str, int, int]] + + class LocalMemoryCacheManager: - """Block-level in-memory cache with LRU eviction.""" + """In-memory block/range cache with LRU eviction.""" def __init__(self, max_size_bytes: int, block_size: int = 1 * 1024 * 1024): self._max_size_bytes = max_size_bytes @@ -50,7 +53,7 @@ def __init__(self, max_size_bytes: int, block_size: int = 1 * 1024 * 1024): def block_size(self) -> int: return self._block_size - def get_block(self, file_path: str, block_index: int) -> Optional[bytes]: + def get_block(self, file_path: str, block_index: _CacheEntryKey) -> Optional[bytes]: key = (file_path, block_index) with self._lock: data = self._cache.get(key) @@ -58,7 +61,7 @@ def get_block(self, file_path: str, block_index: int) -> Optional[bytes]: self._cache.move_to_end(key) return data - def put_block(self, file_path: str, block_index: int, data: bytes) -> None: + def put_block(self, file_path: str, block_index: _CacheEntryKey, data: bytes) -> None: key = (file_path, block_index) with self._lock: if key in self._cache: @@ -79,7 +82,7 @@ def put_file_size(self, file_path: str, size: int) -> None: class LocalDiskCacheManager: - """Block-level local disk cache with LRU eviction.""" + """Local disk block/range cache with LRU eviction.""" def __init__(self, cache_dir: str, max_size_bytes: int, block_size: int = 1 * 1024 * 1024): @@ -98,14 +101,14 @@ def __init__(self, cache_dir: str, max_size_bytes: int, def block_size(self) -> int: return self._block_size - def _cache_path(self, file_path: str, block_index: int) -> str: + def _cache_path(self, file_path: str, block_index: _CacheEntryKey) -> str: key = f"{file_path}:{block_index}" h = hashlib.sha256(key.encode('utf-8')).hexdigest() prefix = h[:2] sub_dir = os.path.join(self._cache_dir, prefix) return os.path.join(sub_dir, h) - def get_block(self, file_path: str, block_index: int) -> Optional[bytes]: + def get_block(self, file_path: str, block_index: _CacheEntryKey) -> Optional[bytes]: path = self._cache_path(file_path, block_index) with self._lock: if path not in self._entry_index: @@ -121,7 +124,7 @@ def get_block(self, file_path: str, block_index: int) -> Optional[bytes]: self._current_size -= size return None - def put_block(self, file_path: str, block_index: int, data: bytes) -> None: + def put_block(self, file_path: str, block_index: _CacheEntryKey, data: bytes) -> None: path = self._cache_path(file_path, block_index) with self._lock: @@ -325,6 +328,33 @@ def __exit__(self, exc_type, exc_val, exc_tb): return False +class BlobMetadataInputStream(CachingInputStream): + """Cache only explicitly marked ranges; ordinary reads bypass the cache.""" + + def read(self, size=-1) -> bytes: + if size is None or size < 0: + size = self._get_file_size() - self._pos + data = self.read_at(size, self._pos) + self._pos += len(data) + return data + + def read_at(self, nbytes: int, offset: int) -> bytes: + return self._read_remote(offset, nbytes) if nbytes > 0 else b'' + + def read_blob_metadata(self, size: int) -> bytes: + if size <= 0: + return b'' + # Tuple keys cannot collide with integer block indexes, in memory or on disk. + key = ('blob-meta', self._pos, size) + data = self._cache.get_block(self._file_path, key) + if data is None or len(data) != size: + data = self._read_remote(self._pos, size) + if len(data) == size: + self._cache.put_block(self._file_path, key, data) + self._pos += len(data) + return data + + class CachingFileIO(FileIO): """FileIO wrapper that caches reads at block granularity. @@ -337,7 +367,7 @@ def __init__(self, delegate: FileIO, cache, whitelist=None): self._delegate = delegate self._cache = cache if whitelist is None: - self._whitelist = {FileType.META, FileType.GLOBAL_INDEX} + self._whitelist = {FileType.META, FileType.GLOBAL_INDEX, FileType.BLOB_META} else: self._whitelist = whitelist @@ -398,9 +428,12 @@ def properties(self): def new_input_stream(self, path: str): file_type = FileType.classify(path) - if self._cache is None or file_type not in self._whitelist or FileType.is_mutable(path): - return self._delegate.new_input_stream(path) - return CachingInputStream(self._delegate, path, self._cache) + if self._cache is not None and not FileType.is_mutable(path): + if file_type in self._whitelist: + return CachingInputStream(self._delegate, path, self._cache) + if FileType.BLOB_META in self._whitelist and path.endswith('.blob'): + return BlobMetadataInputStream(self._delegate, path, self._cache) + return self._delegate.new_input_stream(path) def new_output_stream(self, path: str): return self._delegate.new_output_stream(path) diff --git a/paimon-python/pypaimon/read/reader/format_blob_reader.py b/paimon-python/pypaimon/read/reader/format_blob_reader.py index e39e08b7b229..bd0dbd5647a2 100644 --- a/paimon-python/pypaimon/read/reader/format_blob_reader.py +++ b/paimon-python/pypaimon/read/reader/format_blob_reader.py @@ -387,7 +387,7 @@ def _read_index(self) -> None: # Seek to header: last 5 bytes f.seek(self._file_size - 5) - header = f.read(5) + header = BlobRecordIterator._read_fully_from(f, 5, metadata=True) if len(header) != 5: raise IOError("Invalid blob file: cannot read header") @@ -401,7 +401,7 @@ def _read_index(self) -> None: # Read index data f.seek(self._file_size - 5 - index_length) - index_bytes = f.read(index_length) + index_bytes = BlobRecordIterator._read_fully_from(f, index_length, metadata=True) if len(index_bytes) != index_length: raise IOError("Invalid blob file: cannot read index") @@ -543,7 +543,7 @@ def _read_blob_array(self, position: int, length: int): close_stream = True try: stream.seek(position) - header = self._read_fully_from(stream, self.ARRAY_HEADER_SIZE) + header = self._read_fully_from(stream, self.ARRAY_HEADER_SIZE, metadata=True) if len(header) != self.ARRAY_HEADER_SIZE: raise IOError("Invalid ARRAY payload: cannot read header") magic, version, element_count = struct.unpack(' payload: cannot read index length") index_length = struct.unpack(' payload: cannot read element index") self._validate_array_element_index(index_bytes) @@ -645,7 +645,7 @@ def _read_blob_map(self, position: int, length: int): close_stream = True try: stream.seek(position) - header = self._read_fully_from(stream, self.MAP_HEADER_SIZE) + header = self._read_fully_from(stream, self.MAP_HEADER_SIZE, metadata=True) if len(header) != self.MAP_HEADER_SIZE: raise IOError("Invalid MAP payload: cannot read header") magic, version, entry_count = struct.unpack(' payload: cannot read index lengths") @@ -678,7 +678,7 @@ def _read_blob_map(self, position: int, length: int): stream.seek(key_index_start) # The two indexes are adjacent; read both without touching BLOB values. index_bytes = self._read_fully_from( - stream, key_index_length + value_index_length) + stream, key_index_length + value_index_length, metadata=True) key_index_bytes = index_bytes[:key_index_length] if len(key_index_bytes) != key_index_length: raise IOError("Invalid MAP payload: cannot read key index") @@ -707,7 +707,7 @@ def _read_blob_map(self, position: int, length: int): ) stream.seek(data_start) - key_data = self._read_fully_from(stream, key_data_length) + key_data = self._read_fully_from(stream, key_data_length, metadata=True) if len(key_data) != key_data_length: raise IOError("Invalid MAP payload: cannot read key data") keys = [] @@ -866,10 +866,13 @@ def _read_fully(self, length: int) -> bytes: return self._read_fully_from(self.input_stream, length) @staticmethod - def _read_fully_from(stream, length: int) -> bytes: + def _read_fully_from(stream, length: int, metadata: bool = False) -> bytes: + read = stream.read + if metadata and callable(getattr(type(stream), 'read_blob_metadata', None)): + read = stream.read_blob_metadata data = bytearray() while len(data) < length: - chunk = stream.read(length - len(data)) + chunk = read(length - len(data)) if not chunk: break data.extend(chunk) diff --git a/paimon-python/pypaimon/tests/blob_meta_cache_test.py b/paimon-python/pypaimon/tests/blob_meta_cache_test.py new file mode 100644 index 000000000000..9d41a3b47806 --- /dev/null +++ b/paimon-python/pypaimon/tests/blob_meta_cache_test.py @@ -0,0 +1,141 @@ +# 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 io +from unittest.mock import patch + +import pytest + +from pypaimon.common.options.options import Options +from pypaimon.filesystem.caching_file_io import CachingFileIO, CachingInputStream +from pypaimon.filesystem.local_file_io import LocalFileIO +from pypaimon.read.reader.format_blob_reader import ( + FormatBlobReader, _BLOB_INDEX_CACHE, _BLOB_INDEX_CACHE_LOCK, +) +from pypaimon.schema.data_types import ArrayType, AtomicType, DataField, MapType +from pypaimon.table.row.blob import BlobData, BlobDescriptor +from pypaimon.table.row.generic_row import GenericRow +from pypaimon.utils.file_type import FileType +from pypaimon.write.blob_format_writer import BlobFormatWriter + + +@pytest.mark.parametrize('disk', [False, True], ids=['memory', 'disk']) +@pytest.mark.parametrize('kind', ['scalar', 'array', 'map']) +def test_blob_metadata_ranges_exclude_values(tmp_path, disk, kind): + payload = b'body' * 4096 + blob = AtomicType('BLOB') + types = {'scalar': blob, 'array': ArrayType(True, blob), + 'map': MapType(True, AtomicType('STRING'), blob)} + values = {'scalar': BlobData(payload), + 'array': [BlobData(payload), None, BlobData(b'')], + 'map': [('key', BlobData(payload)), ('missing', None), ('empty', BlobData(b''))]} + field = DataField(0, 'value', types[kind]) + path = str(tmp_path / 'data.blob') + writer = BlobFormatWriter(open(path, 'wb')) + writer.add_element(GenericRow([values[kind]], [field])) + writer.close() + contents = (tmp_path / 'data.blob').read_bytes() + options = Options({'local-cache.enabled': 'true', 'local-cache.max-size': '1 mb', + **({'local-cache.dir': str(tmp_path / 'cache')} if disk else {})}) + delegate = LocalFileIO(str(tmp_path), Options({})) + cache = CachingFileIO.create_cache_manager(options) + file_io = CachingFileIO.wrap_with_caching_if_needed(delegate, options, cache) + reads = [] + + class CountingStream(io.BytesIO): + def read(self, size=-1): + start = self.tell() + result = super().read(size) + reads.append((start, start + len(result))) + return result + + def read(descriptors=True): + # Isolate the byte cache from the existing parsed-index cache. + with _BLOB_INDEX_CACHE_LOCK: + _BLOB_INDEX_CACHE.clear() + reader = FormatBlobReader(file_io, path, ['value'], [field], None, descriptors, + file_size=len(contents)) + try: + return reader.read_arrow_batch().column(0)[0].as_py() + finally: + reader.close() + + with patch.object(delegate, 'new_input_stream', side_effect=lambda _: CountingStream(contents)): + first = read() + serialized = ([first] if kind == 'scalar' else first if kind == 'array' + else [value for _, value in first]) + descriptors = [BlobDescriptor.deserialize(value) for value in serialized if value is not None] + assert reads + for descriptor in descriptors: + if descriptor.length: + assert all(end <= descriptor.offset or start >= descriptor.offset + descriptor.length + for start, end in reads) + retained = cache._current_size + assert 0 < retained < len(payload) + reads.clear() + if disk: + # Reopen the disk manager to verify persistence, not just memory hits. + cache = CachingFileIO.create_cache_manager(options) + file_io = CachingFileIO.wrap_with_caching_if_needed(delegate, options, cache) + assert read() == first + assert reads == [] + materialized = read(False) + expected = {'scalar': payload, 'array': [payload, None, b''], + 'map': [('key', payload), ('missing', None), ('empty', b'')]}[kind] + assert materialized == expected + assert reads # Values still come from the delegate. + assert cache._current_size == retained + + +def test_blob_meta_is_a_read_category_not_a_file_type(tmp_path): + delegate = LocalFileIO(str(tmp_path), Options({})) + options = Options({'local-cache.enabled': 'true', 'local-cache.whitelist': 'blob-meta'}) + cache = CachingFileIO.create_cache_manager(options) + file_io = CachingFileIO.wrap_with_caching_if_needed(delegate, options, cache) + assert FileType.classify('data.blob') == FileType.DATA + assert FileType.parse_whitelist('blob-meta') == {FileType.BLOB_META} + for name in ('data.parquet', '.data.blob.uuid.tmp'): + (tmp_path / name).write_bytes(b'data') + with file_io.new_input_stream(str(tmp_path / name)) as stream: + assert stream.read() == b'data' + assert cache._current_size == 0 + + +@pytest.mark.parametrize('disk', [False, True]) +def test_metadata_cache_budget_and_data_cache_coexist(tmp_path, disk): + path = str(tmp_path / 'data.blob') + (tmp_path / 'data.blob').write_bytes(b'abcdefgh') + delegate = LocalFileIO(str(tmp_path), Options({})) + options = Options({'local-cache.enabled': 'true', 'local-cache.whitelist': 'blob-meta', + 'local-cache.max-size': '4 b', 'local-cache.block-size': '4 b', + **({'local-cache.dir': str(tmp_path / 'cache')} if disk else {})}) + cache = CachingFileIO.create_cache_manager(options) + file_io = CachingFileIO.wrap_with_caching_if_needed(delegate, options, cache) + with file_io.new_input_stream(path) as stream: + assert stream.read_blob_metadata(4) == b'abcd' + assert stream.read_blob_metadata(4) == b'efgh' + assert cache._current_size == 4 + assert stream.read_blob_metadata(4) == b'' # Never cache a short read. + assert cache._current_size == 4 + stream.seek(0) + assert stream.read_blob_metadata(4) == b'abcd' + assert cache._current_size == 4 + # Reusing a manager for normal data blocks cannot hit metadata range entries. + data_io = CachingFileIO(delegate, cache, {FileType.DATA, FileType.BLOB_META}) + with data_io.new_input_stream(path) as stream: + assert type(stream) is CachingInputStream + assert stream.read() == b'abcdefgh' diff --git a/paimon-python/pypaimon/tests/caching_file_io_test.py b/paimon-python/pypaimon/tests/caching_file_io_test.py index 940a4f4ab29a..d721343561a0 100644 --- a/paimon-python/pypaimon/tests/caching_file_io_test.py +++ b/paimon-python/pypaimon/tests/caching_file_io_test.py @@ -492,7 +492,7 @@ def test_local_cache_options_defaults(self): self.assertIsNone(opts.local_cache_dir()) self.assertIsNone(opts.local_cache_max_size()) self.assertEqual(1 * 1024 * 1024, opts.local_cache_block_size().get_bytes()) - self.assertEqual("meta,global-index", opts.local_cache_whitelist()) + self.assertEqual("meta,global-index,blob-meta", opts.local_cache_whitelist()) def test_local_cache_options_custom(self): from pypaimon.common.options import Options diff --git a/paimon-python/pypaimon/utils/file_type.py b/paimon-python/pypaimon/utils/file_type.py index 7dadc93af631..61791ba6b795 100644 --- a/paimon-python/pypaimon/utils/file_type.py +++ b/paimon-python/pypaimon/utils/file_type.py @@ -34,12 +34,14 @@ class FileType(Enum): - BUCKET_INDEX: bucket level index files (Hash, DV) - GLOBAL_INDEX: table level global index files (btree, lumina, full-text) - FILE_INDEX: data-file index files (bloom filter, bitmap, etc.) + - BLOB_META: reader-selected metadata ranges within BLOB files """ META = "META" DATA = "DATA" BUCKET_INDEX = "BUCKET_INDEX" GLOBAL_INDEX = "GLOBAL_INDEX" FILE_INDEX = "FILE_INDEX" + BLOB_META = "BLOB_META" def is_index(self) -> bool: return self in (FileType.BUCKET_INDEX, FileType.GLOBAL_INDEX, FileType.FILE_INDEX) @@ -94,6 +96,7 @@ def parse_whitelist(whitelist_str: str) -> set: "bucket-index": FileType.BUCKET_INDEX, "data": FileType.DATA, "file-index": FileType.FILE_INDEX, + "blob-meta": FileType.BLOB_META, } result = set() for name in whitelist_str.split(","): @@ -103,7 +106,7 @@ def parse_whitelist(whitelist_str: str) -> set: elif name: logger.warning( "Unknown local-cache.whitelist value '%s'. " - "Supported values: meta, global-index, bucket-index, data, file-index.", + "Supported values: meta, global-index, bucket-index, data, file-index, blob-meta.", name, ) return result From 37db156d438519044a297658b65e274bb165803e Mon Sep 17 00:00:00 2001 From: "xiaohongbo.xhb" Date: Sat, 26 Sep 2026 18:37:02 -0700 Subject: [PATCH 02/10] [docs] Shorten BLOB metadata cache documentation --- docs/docs/pypaimon/blob.md | 8 +++----- 1 file changed, 3 insertions(+), 5 deletions(-) diff --git a/docs/docs/pypaimon/blob.md b/docs/docs/pypaimon/blob.md index eed1f8c9f091..76edc63ced4f 100644 --- a/docs/docs/pypaimon/blob.md +++ b/docs/docs/pypaimon/blob.md @@ -178,11 +178,9 @@ The factory auto-dispatches based on the bytes content (`BLOBDESC`, ## BLOB metadata cache -With `local-cache.enabled=true`, the Python reader's default whitelist is -`meta,global-index,blob-meta`. `blob-meta` caches exact ranges for BLOB footers, -row indexes, ARRAY headers/indexes, and MAP headers/keys/indexes, excluding value -bodies. It shares the local cache's memory or disk budget. Add `data` to cache -normal data blocks as well. This option applies to the Python reader, not Rust native reads. +Enable `local-cache.enabled=true` to cache BLOB metadata needed for descriptors, +including MAP keys, without caching value bodies. `blob-meta` is included in the +default whitelist and shares the local cache budget. Python reader only. ## See Also From c6f8aeb59bae36ce4227c6200b00590194f2de56 Mon Sep 17 00:00:00 2001 From: "xiaohongbo.xhb" Date: Sat, 26 Sep 2026 20:28:17 -0700 Subject: [PATCH 03/10] [python] Bound metadata cache costs and support readinto --- docs/docs/program-api/file-cache.mdx | 9 +- .../pypaimon/filesystem/caching_file_io.py | 134 ++++++++++-------- .../pypaimon/tests/blob_meta_cache_test.py | 110 +++++++++++++- .../pypaimon/tests/caching_file_io_test.py | 3 +- 4 files changed, 189 insertions(+), 67 deletions(-) diff --git a/docs/docs/program-api/file-cache.mdx b/docs/docs/program-api/file-cache.mdx index 5f47cf653c59..be7746c5d3c4 100644 --- a/docs/docs/program-api/file-cache.mdx +++ b/docs/docs/program-api/file-cache.mdx @@ -51,8 +51,9 @@ The cache classifies files by path. Set `local-cache.whitelist` to select the fi | BUCKET_INDEX | bucket-index | Hash, deletion vector index files | No | | DATA | data | Data files (ORC, Parquet, etc.) | No | | FILE_INDEX | file-index | Data-file level bloom filter, bitmap | No | +| BLOB_META | blob-meta | BLOB metadata read ranges (Python only) | Python only | -The default whitelist is `meta,global-index`. The Python reader additionally includes +The Java default whitelist is `meta,global-index`. The Python reader additionally includes `blob-meta`, which caches exact BLOB metadata ranges rather than whole data blocks; see [Python BLOB metadata cache](../pypaimon/blob.md#blob-metadata-cache). The mutable `LATEST` and `EARLIEST` hint files bypass the cache even though they are classified as metadata. @@ -128,7 +129,11 @@ The Java snippet is a method-body fragment; declare or handle the catalog except | `local-cache.dir` | String | (none) | Directory for storing cached blocks on disk. If not configured, memory cache is used. | | `local-cache.max-size` | MemorySize | See below | Maximum cached block bytes per cache manager. Least recently used blocks are evicted when the limit is exceeded. | | `local-cache.block-size` | MemorySize | 1 mb | Block size for caching. Files are logically divided into fixed-size blocks and cached independently. | -| `local-cache.whitelist` | String | meta,global-index | Comma-separated list of file types to cache. Supported values: `meta`, `global-index`, `bucket-index`, `data`, `file-index`. | +| `local-cache.whitelist` | String | Java: `meta,global-index`; Python: `meta,global-index,blob-meta` | Comma-separated cache categories: `meta`, `global-index`, `bucket-index`, `data`, `file-index`, and Python-only `blob-meta`. | + +Python charges estimated object overhead for memory entries and allocated disk bytes (at least 4 KiB) +for disk entries, and limits each cache manager to 65,536 entries. These limits do not represent a +strict process RSS or total filesystem usage limit. When `local-cache.max-size` is omitted, Java has no configured size limit. PyPaimon uses a 256 MiB memory limit or a 10 GiB disk limit. Set the option explicitly when you want the same limit across diff --git a/paimon-python/pypaimon/filesystem/caching_file_io.py b/paimon-python/pypaimon/filesystem/caching_file_io.py index 182eb9067ee9..f35fe333b034 100644 --- a/paimon-python/pypaimon/filesystem/caching_file_io.py +++ b/paimon-python/pypaimon/filesystem/caching_file_io.py @@ -27,6 +27,7 @@ import hashlib import os +import sys import threading from collections import OrderedDict from typing import Optional, Tuple, Union @@ -36,6 +37,17 @@ _CacheEntryKey = Union[int, Tuple[str, int, int]] +# Bound Python index overhead and the number of cache files, even with an unlimited byte budget. +_MAX_CACHE_ENTRIES = 65536 + + +def _memory_entry_size(key, data): + # Conservatively count referenced objects plus OrderedDict bookkeeping. + path, entry_key = key + size = sys.getsizeof(key) + sys.getsizeof(path) + sys.getsizeof(entry_key) + if isinstance(entry_key, tuple): + size += sum(sys.getsizeof(part) for part in entry_key) + return max(512, size + sys.getsizeof(data) + 128) class LocalMemoryCacheManager: @@ -53,26 +65,25 @@ def __init__(self, max_size_bytes: int, block_size: int = 1 * 1024 * 1024): def block_size(self) -> int: return self._block_size - def get_block(self, file_path: str, block_index: _CacheEntryKey) -> Optional[bytes]: - key = (file_path, block_index) + def get_block(self, file_path: str, cache_key: _CacheEntryKey) -> Optional[bytes]: + key = (file_path, cache_key) with self._lock: data = self._cache.get(key) if data is not None: self._cache.move_to_end(key) return data - def put_block(self, file_path: str, block_index: _CacheEntryKey, data: bytes) -> None: - key = (file_path, block_index) + def put_block(self, file_path: str, cache_key: _CacheEntryKey, data: bytes) -> None: + key = (file_path, cache_key) with self._lock: if key in self._cache: return - self._current_size += len(data) + self._current_size += _memory_entry_size(key, data) self._cache[key] = data - while (self._max_size_bytes < (2 ** 63 - 1) - and self._current_size > self._max_size_bytes - and self._cache): - _, evicted = self._cache.popitem(last=False) - self._current_size -= len(evicted) + while self._cache and (self._current_size > self._max_size_bytes + or len(self._cache) > _MAX_CACHE_ENTRIES): + evicted_key, evicted = self._cache.popitem(last=False) + self._current_size -= _memory_entry_size(evicted_key, evicted) def get_file_size(self, file_path: str) -> int: return self._file_size_cache.get(file_path, -1) @@ -95,21 +106,21 @@ def __init__(self, cache_dir: str, max_size_bytes: int, # LRU-ordered index: cache_path -> size. OrderedDict with move_to_end for access order. self._entry_index: OrderedDict = OrderedDict() os.makedirs(cache_dir, exist_ok=True) - self._current_size = self._scan_and_populate_index() + self._scan_and_populate_index() @property def block_size(self) -> int: return self._block_size - def _cache_path(self, file_path: str, block_index: _CacheEntryKey) -> str: - key = f"{file_path}:{block_index}" + def _cache_path(self, file_path: str, cache_key: _CacheEntryKey) -> str: + key = f"{file_path}:{cache_key}" h = hashlib.sha256(key.encode('utf-8')).hexdigest() prefix = h[:2] sub_dir = os.path.join(self._cache_dir, prefix) return os.path.join(sub_dir, h) - def get_block(self, file_path: str, block_index: _CacheEntryKey) -> Optional[bytes]: - path = self._cache_path(file_path, block_index) + def get_block(self, file_path: str, cache_key: _CacheEntryKey) -> Optional[bytes]: + path = self._cache_path(file_path, cache_key) with self._lock: if path not in self._entry_index: return None @@ -124,69 +135,66 @@ def get_block(self, file_path: str, block_index: _CacheEntryKey) -> Optional[byt self._current_size -= size return None - def put_block(self, file_path: str, block_index: _CacheEntryKey, data: bytes) -> None: - path = self._cache_path(file_path, block_index) - + def put_block(self, file_path: str, cache_key: _CacheEntryKey, data: bytes) -> None: + path = self._cache_path(file_path, cache_key) + # Keep publication and eviction atomic with respect to other writers. with self._lock: if path in self._entry_index: return - - sub_dir = os.path.dirname(path) - os.makedirs(sub_dir, exist_ok=True) - - tmp_path = path + f".tmp.{os.getpid()}.{threading.get_ident()}" - try: - with open(tmp_path, 'wb') as f: - f.write(data) - os.rename(tmp_path, path) - except Exception: + sub_dir = os.path.dirname(path) + os.makedirs(sub_dir, exist_ok=True) + tmp_path = path + f".tmp.{os.getpid()}.{threading.get_ident()}" try: - os.unlink(tmp_path) - except OSError: - pass - return - - need_evict = False - with self._lock: - self._entry_index[path] = len(data) - self._current_size += len(data) - need_evict = (self._max_size_bytes < (2 ** 63 - 1) - and self._current_size > self._max_size_bytes) - if need_evict: - self._evict() - - def _evict(self) -> None: - to_delete = [] - with self._lock: - if self._current_size <= self._max_size_bytes: + with open(tmp_path, 'wb') as f: + f.write(data) + size = self._disk_entry_size(tmp_path) + os.rename(tmp_path, path) + except Exception: + try: + os.unlink(tmp_path) + except OSError: + pass return - while self._entry_index and self._current_size > self._max_size_bytes: - path, size = self._entry_index.popitem(last=False) - self._current_size -= size - to_delete.append((path, size)) + self._entry_index[path] = size + self._current_size += size + self._evict_locked() - for path, size in to_delete: + @staticmethod + def _disk_entry_size(path: str) -> int: + stat = os.stat(path) + # st_blocks is in 512-byte units. Keep a minimum charge on platforms + # without allocation information and for sparse/empty files. + return max(stat.st_size, getattr(stat, 'st_blocks', 0) * 512, + getattr(stat, 'st_blksize', 4096), 4096) + + def _evict_locked(self) -> None: + while (self._entry_index + and (self._current_size > self._max_size_bytes + or len(self._entry_index) > _MAX_CACHE_ENTRIES)): + path, size = self._entry_index.popitem(last=False) try: os.unlink(path) + except FileNotFoundError: + pass except OSError: - with self._lock: - self._entry_index[path] = size - self._current_size += size + self._entry_index[path] = size + break + self._current_size -= size - def _scan_and_populate_index(self) -> int: - total = 0 + def _scan_and_populate_index(self) -> None: for dirpath, _, filenames in os.walk(self._cache_dir): for fn in filenames: if '.tmp.' in fn: continue fp = os.path.join(dirpath, fn) try: - size = os.path.getsize(fp) + size = self._disk_entry_size(fp) self._entry_index[fp] = size - total += size + self._current_size += size + # Bound the index while scanning an existing cache directory. + self._evict_locked() except OSError: pass - return total def get_file_size(self, file_path: str) -> int: return self._file_size_cache.get(file_path, -1) @@ -254,6 +262,14 @@ def read(self, size=-1) -> bytes: self._pos = end return bytes(result) + def readinto(self, buffer) -> int: + view = memoryview(buffer).cast('B') + if view.readonly: + raise TypeError("readinto() requires a writable buffer") + data = self.read(len(view)) + view[:len(data)] = data + return len(data) + def read_at(self, nbytes: int, offset: int) -> bytes: """Position-based read. Does not change the cursor. Thread-safe.""" if nbytes <= 0 or offset >= self._get_file_size(): diff --git a/paimon-python/pypaimon/tests/blob_meta_cache_test.py b/paimon-python/pypaimon/tests/blob_meta_cache_test.py index 9d41a3b47806..82b7b0f03324 100644 --- a/paimon-python/pypaimon/tests/blob_meta_cache_test.py +++ b/paimon-python/pypaimon/tests/blob_meta_cache_test.py @@ -85,7 +85,7 @@ def read(descriptors=True): assert all(end <= descriptor.offset or start >= descriptor.offset + descriptor.length for start, end in reads) retained = cache._current_size - assert 0 < retained < len(payload) + assert retained > 0 reads.clear() if disk: # Reopen the disk manager to verify persistence, not just memory hits. @@ -121,21 +121,121 @@ def test_metadata_cache_budget_and_data_cache_coexist(tmp_path, disk): (tmp_path / 'data.blob').write_bytes(b'abcdefgh') delegate = LocalFileIO(str(tmp_path), Options({})) options = Options({'local-cache.enabled': 'true', 'local-cache.whitelist': 'blob-meta', - 'local-cache.max-size': '4 b', 'local-cache.block-size': '4 b', + 'local-cache.max-size': '4096 b', 'local-cache.block-size': '4 b', **({'local-cache.dir': str(tmp_path / 'cache')} if disk else {})}) cache = CachingFileIO.create_cache_manager(options) file_io = CachingFileIO.wrap_with_caching_if_needed(delegate, options, cache) with file_io.new_input_stream(path) as stream: assert stream.read_blob_metadata(4) == b'abcd' assert stream.read_blob_metadata(4) == b'efgh' - assert cache._current_size == 4 + assert 0 < cache._current_size <= 4096 + retained = cache._current_size assert stream.read_blob_metadata(4) == b'' # Never cache a short read. - assert cache._current_size == 4 + assert cache._current_size == retained stream.seek(0) assert stream.read_blob_metadata(4) == b'abcd' - assert cache._current_size == 4 + assert 0 < cache._current_size <= 4096 # Reusing a manager for normal data blocks cannot hit metadata range entries. data_io = CachingFileIO(delegate, cache, {FileType.DATA, FileType.BLOB_META}) with data_io.new_input_stream(path) as stream: assert type(stream) is CachingInputStream assert stream.read() == b'abcdefgh' + + +@pytest.mark.parametrize('disk', [False, True], ids=['memory', 'disk']) +@pytest.mark.parametrize('kind', ['array', 'map']) +def test_many_blob_rows_respect_cache_budget(tmp_path, disk, kind): + blob = AtomicType('BLOB') + field = DataField(0, 'value', ArrayType(True, blob) if kind == 'array' + else MapType(True, AtomicType('STRING'), blob)) + value = [BlobData(b'x')] if kind == 'array' else [('key', BlobData(b'x'))] + path = str(tmp_path / 'data.blob') + writer = BlobFormatWriter(open(path, 'wb')) + for _ in range(1000): + writer.add_element(GenericRow([value], [field])) + writer.close() + budget = 64 * 1024 + options = Options({'local-cache.enabled': 'true', 'local-cache.max-size': '64 kb', + **({'local-cache.dir': str(tmp_path / 'cache')} if disk else {})}) + cache = CachingFileIO.create_cache_manager(options) + delegate = LocalFileIO(str(tmp_path), Options({})) + file_io = CachingFileIO.wrap_with_caching_if_needed(delegate, options, cache) + reader = FormatBlobReader(file_io, path, ['value'], [field], None, True) + try: + batch = reader.read_arrow_batch() + assert batch.num_rows == 1000 + for row in batch.column(0).to_pylist(): + raw = row[0] if kind == 'array' else row[0][1] + descriptor = BlobDescriptor.deserialize(raw) + with open(path, 'rb') as stream: + stream.seek(descriptor.offset) + assert stream.read(descriptor.length) == b'x' + finally: + reader.close() + assert 0 < cache._current_size <= budget + entries = cache._entry_index if disk else cache._cache + assert 0 < len(entries) <= budget // (4096 if disk else 512) + if disk: + files = [p for p in (tmp_path / 'cache').rglob('*') if p.is_file()] + assert len(files) == len(entries) + assert sum(p.stat().st_blocks * 512 for p in files) <= budget + reopened = CachingFileIO.create_cache_manager(options) + assert reopened._current_size == cache._current_size + assert len(reopened._entry_index) == len(entries) + + +@pytest.mark.parametrize('disk', [False, True], ids=['memory', 'disk']) +def test_tiny_ranges_respect_entry_limit_and_lru(tmp_path, disk): + options = Options({'local-cache.enabled': 'true', 'local-cache.max-size': '1 gb', + **({'local-cache.dir': str(tmp_path / 'cache')} if disk else {})}) + cache = CachingFileIO.create_cache_manager(options) + with patch('pypaimon.filesystem.caching_file_io._MAX_CACHE_ENTRIES', 32): + for offset in range(32): + cache.put_block('data.blob', ('blob-meta', offset, 1), b'x') + assert cache.get_block('data.blob', ('blob-meta', 0, 1)) == b'x' + cache.put_block('data.blob', ('blob-meta', 32, 1), b'x') + assert cache.get_block('data.blob', ('blob-meta', 0, 1)) == b'x' + assert cache.get_block('data.blob', ('blob-meta', 1, 1)) is None + for offset in range(33, 1000): + cache.put_block('data.blob', ('blob-meta', offset, 1), b'x') + entries = cache._entry_index if disk else cache._cache + assert len(entries) == 32 + assert cache._current_size >= 32 * (4096 if disk else 512) + if disk: + # Opening an existing directory must enforce both limits immediately. + with patch('pypaimon.filesystem.caching_file_io._MAX_CACHE_ENTRIES', 8): + reopened = CachingFileIO.create_cache_manager(options) + assert len(reopened._entry_index) == 8 + options = Options({'local-cache.enabled': 'true', 'local-cache.max-size': '4 kb', + 'local-cache.dir': str(tmp_path / 'cache')}) + reopened = CachingFileIO.create_cache_manager(options) + assert reopened._current_size <= 4096 + assert len(reopened._entry_index) <= 1 + assert len([p for p in (tmp_path / 'cache').rglob('*') if p.is_file()]) <= 1 + + +@pytest.mark.parametrize('whitelist', ['blob-meta', 'data']) +def test_blob_readinto(tmp_path, whitelist): + from pypaimon.table.row.blob import BlobRef + + path = str(tmp_path / 'data.blob') + (tmp_path / 'data.blob').write_bytes(b'abcdefgh') + delegate = LocalFileIO(str(tmp_path), Options({})) + options = Options({'local-cache.enabled': 'true', 'local-cache.whitelist': whitelist}) + cache = CachingFileIO.create_cache_manager(options) + file_io = CachingFileIO.wrap_with_caching_if_needed(delegate, options, cache) + blob = BlobRef(file_io, BlobDescriptor(path, 1, 4)) + with blob.new_input_stream() as stream: + assert stream.readinto(bytearray()) == 0 + buf = bytearray(b'------') + assert stream.readinto(buf) == 4 + assert buf == b'bcde--' + assert stream.readinto(buf) == 0 + stream.seek(0) + assert stream.readinto(memoryview(buf)[1:3]) == 2 + assert buf == b'bbce--' + assert stream.tell() == 2 + if whitelist == 'blob-meta': + assert cache._current_size == 0 # readinto must not cache value bodies. + else: + assert cache._current_size > 0 diff --git a/paimon-python/pypaimon/tests/caching_file_io_test.py b/paimon-python/pypaimon/tests/caching_file_io_test.py index d721343561a0..46541009f4e6 100644 --- a/paimon-python/pypaimon/tests/caching_file_io_test.py +++ b/paimon-python/pypaimon/tests/caching_file_io_test.py @@ -86,7 +86,8 @@ def test_scan_size_on_restart(self): # Simulate restart: new cache instance on same directory cache2 = LocalDiskCacheManager(self.cache_dir, 2 ** 63 - 1, block_size=64) - self.assertEqual(300, cache2._current_size) + self.assertEqual(cache1._current_size, cache2._current_size) + self.assertGreaterEqual(cache2._current_size, 8192) self.assertEqual(b"x" * 100, cache2.get_block("f", 0)) self.assertEqual(b"y" * 200, cache2.get_block("f", 1)) From 43213e94af82dbaaf8a07964d150efc8b9175b55 Mon Sep 17 00:00:00 2001 From: "xiaohongbo.xhb" Date: Sat, 26 Sep 2026 20:41:59 -0700 Subject: [PATCH 04/10] [python] Write disk cache entries outside the index lock --- .../pypaimon/filesystem/caching_file_io.py | 38 +++++++------ .../pypaimon/tests/caching_file_io_test.py | 56 ++++++++++++++++++- 2 files changed, 76 insertions(+), 18 deletions(-) diff --git a/paimon-python/pypaimon/filesystem/caching_file_io.py b/paimon-python/pypaimon/filesystem/caching_file_io.py index f35fe333b034..c9264e68336e 100644 --- a/paimon-python/pypaimon/filesystem/caching_file_io.py +++ b/paimon-python/pypaimon/filesystem/caching_file_io.py @@ -137,27 +137,31 @@ def get_block(self, file_path: str, cache_key: _CacheEntryKey) -> Optional[bytes def put_block(self, file_path: str, cache_key: _CacheEntryKey, data: bytes) -> None: path = self._cache_path(file_path, cache_key) - # Keep publication and eviction atomic with respect to other writers. with self._lock: if path in self._entry_index: return - sub_dir = os.path.dirname(path) - os.makedirs(sub_dir, exist_ok=True) - tmp_path = path + f".tmp.{os.getpid()}.{threading.get_ident()}" - try: - with open(tmp_path, 'wb') as f: - f.write(data) - size = self._disk_entry_size(tmp_path) + + tmp_path = path + f".tmp.{os.getpid()}.{threading.get_ident()}" + try: + os.makedirs(os.path.dirname(path), exist_ok=True) + with open(tmp_path, 'wb') as f: + f.write(data) + size = self._disk_entry_size(tmp_path) + with self._lock: + # Another writer may have published this key during the write. + if path in self._entry_index: + return os.rename(tmp_path, path) - except Exception: - try: - os.unlink(tmp_path) - except OSError: - pass - return - self._entry_index[path] = size - self._current_size += size - self._evict_locked() + self._entry_index[path] = size + self._current_size += size + self._evict_locked() + except OSError: + return + finally: + try: + os.unlink(tmp_path) + except OSError: + pass @staticmethod def _disk_entry_size(path: str) -> int: diff --git a/paimon-python/pypaimon/tests/caching_file_io_test.py b/paimon-python/pypaimon/tests/caching_file_io_test.py index 46541009f4e6..3f6450330fcc 100644 --- a/paimon-python/pypaimon/tests/caching_file_io_test.py +++ b/paimon-python/pypaimon/tests/caching_file_io_test.py @@ -25,7 +25,8 @@ import tempfile import threading import unittest -from unittest.mock import MagicMock +from concurrent.futures import ThreadPoolExecutor +from unittest.mock import MagicMock, patch from pypaimon.filesystem.caching_file_io import ( LocalDiskCacheManager, @@ -127,6 +128,59 @@ def reader(idx): self.assertEqual([], errors) + def test_slow_write_does_not_block_unrelated_hit(self): + cache = LocalDiskCacheManager(self.cache_dir, 1024 * 1024) + cache.put_block("existing", 0, b"cached") + writing = threading.Event() + release = threading.Event() + real_open = open + + class SlowFile(io.FileIO): + def write(self, data): + writing.set() + if not release.wait(5): + raise AssertionError("Timed out waiting to release write") + return super().write(data) + + def open_file(path, mode): + return SlowFile(path, mode) if mode == 'wb' else real_open(path, mode) + + with patch('pypaimon.filesystem.caching_file_io.open', side_effect=open_file): + with ThreadPoolExecutor(max_workers=2) as executor: + writer = executor.submit(cache.put_block, "new", 0, b"new data") + try: + self.assertTrue(writing.wait(5)) + hit = executor.submit(cache.get_block, "existing", 0) + self.assertEqual(b"cached", hit.result(timeout=2)) + self.assertFalse(writer.done()) + finally: + release.set() + writer.result(timeout=5) + self.assertEqual(b"new data", cache.get_block("new", 0)) + + def test_concurrent_same_key_is_published_once(self): + cache = LocalDiskCacheManager(self.cache_dir, 1024 * 1024) + ready = threading.Barrier(2) + entry_size = cache._disk_entry_size + + def wait_for_writers(path): + size = entry_size(path) + ready.wait(timeout=5) + return size + + with patch.object(cache, '_disk_entry_size', side_effect=wait_for_writers): + with ThreadPoolExecutor(max_workers=2) as executor: + writers = [executor.submit(cache.put_block, "same", 0, data) + for data in (b"first", b"second")] + for writer in writers: + writer.result(timeout=5) + self.assertIn(cache.get_block("same", 0), (b"first", b"second")) + self.assertEqual(1, len(cache._entry_index)) + path = cache._cache_path("same", 0) + self.assertEqual(entry_size(path), cache._current_size) + files = [name for _, _, names in os.walk(self.cache_dir) for name in names] + self.assertEqual([os.path.basename(path)], files) + def test_unlimited_cache_skips_eviction(self): cache = LocalDiskCacheManager(self.cache_dir, max_size_bytes=2 ** 63 - 1, block_size=64) for i in range(50): From 27fe6272d281c0522ad0a4fefb8977b4b5d13741 Mon Sep 17 00:00:00 2001 From: "xiaohongbo.xhb" Date: Sat, 26 Sep 2026 21:01:27 -0700 Subject: [PATCH 05/10] [python] Read BLOB bodies directly into caller buffers --- .../pypaimon/filesystem/caching_file_io.py | 22 +++++++++ .../pypaimon/tests/blob_meta_cache_test.py | 47 +++++++++++++++++++ 2 files changed, 69 insertions(+) diff --git a/paimon-python/pypaimon/filesystem/caching_file_io.py b/paimon-python/pypaimon/filesystem/caching_file_io.py index c9264e68336e..3ded9fa3fe6a 100644 --- a/paimon-python/pypaimon/filesystem/caching_file_io.py +++ b/paimon-python/pypaimon/filesystem/caching_file_io.py @@ -30,6 +30,7 @@ import sys import threading from collections import OrderedDict +from io import UnsupportedOperation from typing import Optional, Tuple, Union from pypaimon.common.file_io import FileIO, supports_pread, pread @@ -358,6 +359,27 @@ def read(self, size=-1) -> bytes: self._pos += len(data) return data + def readinto(self, buffer) -> Optional[int]: + view = memoryview(buffer).cast('B') + if view.readonly: + raise TypeError("readinto() requires a writable buffer") + if not view: + return 0 + with self._io_lock: + stream = self._get_remote_stream() + readinto = getattr(stream, 'readinto', None) + if callable(readinto): + stream.seek(self._pos) + try: + count = readinto(view) + except (UnsupportedOperation, NotImplementedError): + pass + else: + if count is not None: + self._pos += count + return count + return super().readinto(view) + def read_at(self, nbytes: int, offset: int) -> bytes: return self._read_remote(offset, nbytes) if nbytes > 0 else b'' diff --git a/paimon-python/pypaimon/tests/blob_meta_cache_test.py b/paimon-python/pypaimon/tests/blob_meta_cache_test.py index 82b7b0f03324..2f637830bc8c 100644 --- a/paimon-python/pypaimon/tests/blob_meta_cache_test.py +++ b/paimon-python/pypaimon/tests/blob_meta_cache_test.py @@ -239,3 +239,50 @@ def test_blob_readinto(tmp_path, whitelist): assert cache._current_size == 0 # readinto must not cache value bodies. else: assert cache._current_size > 0 + + +@pytest.mark.parametrize('support', ['direct', 'missing', 'unsupported']) +def test_blob_readinto_delegates_without_caching(tmp_path, support): + from pypaimon.table.row.blob import BlobRef + + calls = {'read': 0, 'readinto': 0} + + class CountingStream(io.BytesIO): + def read(self, size=-1): + calls['read'] += 1 + return super().read(size) + + def readinto(self, buffer): + calls['readinto'] += 1 + if support == 'unsupported': + raise io.UnsupportedOperation('readinto') + # Exercise a short read and cursor updates. + return super().readinto(memoryview(buffer)[:2]) + + remote = CountingStream(b'abcdefgh') + if support == 'missing': + remote.readinto = None + delegate = LocalFileIO(str(tmp_path), Options({})) + options = Options({'local-cache.enabled': 'true'}) + cache = CachingFileIO.create_cache_manager(options) + file_io = CachingFileIO.wrap_with_caching_if_needed(delegate, options, cache) + blob = BlobRef(file_io, BlobDescriptor('data.blob', 1, 4)) + with patch.object(delegate, 'new_input_stream', return_value=remote): + with blob.new_input_stream() as stream: + buf = bytearray(b'----') + if support == 'direct': + assert stream.readinto(buf) == 2 + assert buf == b'bc--' + assert stream.tell() == 2 + assert stream.readinto(memoryview(buf)[2:]) == 2 + assert calls == {'read': 0, 'readinto': 2} + else: + assert stream.readinto(buf) == 4 + assert calls['read'] > 0 + assert buf == b'bcde' + assert stream.tell() == 4 + assert stream.readinto(buf) == 0 + stream.seek(0) + assert stream.readinto(memoryview(buf)[:2]) == 2 + assert stream.tell() == 2 + assert cache._current_size == 0 From 126e519b6fca624804a81370836db04febd44ede Mon Sep 17 00:00:00 2001 From: "xiaohongbo.xhb" Date: Sat, 26 Sep 2026 21:25:03 -0700 Subject: [PATCH 06/10] [python] Bypass oversized cache entries and unlink outside lock --- .../pypaimon/filesystem/caching_file_io.py | 38 ++++++++++++++++--- .../pypaimon/tests/blob_meta_cache_test.py | 32 ++++++++++++++++ .../pypaimon/tests/caching_file_io_test.py | 34 +++++++++++++++++ 3 files changed, 99 insertions(+), 5 deletions(-) diff --git a/paimon-python/pypaimon/filesystem/caching_file_io.py b/paimon-python/pypaimon/filesystem/caching_file_io.py index 3ded9fa3fe6a..5e23c5c08741 100644 --- a/paimon-python/pypaimon/filesystem/caching_file_io.py +++ b/paimon-python/pypaimon/filesystem/caching_file_io.py @@ -29,6 +29,7 @@ import os import sys import threading +import uuid from collections import OrderedDict from io import UnsupportedOperation from typing import Optional, Tuple, Union @@ -76,10 +77,13 @@ def get_block(self, file_path: str, cache_key: _CacheEntryKey) -> Optional[bytes def put_block(self, file_path: str, cache_key: _CacheEntryKey, data: bytes) -> None: key = (file_path, cache_key) + size = _memory_entry_size(key, data) + if size > self._max_size_bytes: + return with self._lock: if key in self._cache: return - self._current_size += _memory_entry_size(key, data) + self._current_size += size self._cache[key] = data while self._cache and (self._current_size > self._max_size_bytes or len(self._cache) > _MAX_CACHE_ENTRIES): @@ -137,6 +141,8 @@ def get_block(self, file_path: str, cache_key: _CacheEntryKey) -> Optional[bytes return None def put_block(self, file_path: str, cache_key: _CacheEntryKey, data: bytes) -> None: + if max(len(data), 4096) > self._max_size_bytes: + return path = self._cache_path(file_path, cache_key) with self._lock: if path in self._entry_index: @@ -148,6 +154,8 @@ def put_block(self, file_path: str, cache_key: _CacheEntryKey, data: bytes) -> N with open(tmp_path, 'wb') as f: f.write(data) size = self._disk_entry_size(tmp_path) + if size > self._max_size_bytes: + return with self._lock: # Another writer may have published this key during the write. if path in self._entry_index: @@ -155,7 +163,8 @@ def put_block(self, file_path: str, cache_key: _CacheEntryKey, data: bytes) -> N os.rename(tmp_path, path) self._entry_index[path] = size self._current_size += size - self._evict_locked() + evicted = self._evict_locked() + self._delete_evicted(evicted) except OSError: return finally: @@ -172,19 +181,38 @@ def _disk_entry_size(path: str) -> int: return max(stat.st_size, getattr(stat, 'st_blocks', 0) * 512, getattr(stat, 'st_blksize', 4096), 4096) - def _evict_locked(self) -> None: + def _evict_locked(self): + evicted = [] while (self._entry_index and (self._current_size > self._max_size_bytes or len(self._entry_index) > _MAX_CACHE_ENTRIES)): path, size = self._entry_index.popitem(last=False) + # Detach the old file before unlocking, so deletion cannot remove a + # newly published entry for the same key. + retired = os.path.join(os.path.dirname(path), uuid.uuid4().hex + '.evicted') try: - os.unlink(path) + os.rename(path, retired) except FileNotFoundError: pass except OSError: self._entry_index[path] = size break + else: + evicted.append((retired, size)) self._current_size -= size + return evicted + + def _delete_evicted(self, evicted) -> None: + for path, size in evicted: + try: + os.unlink(path) + except FileNotFoundError: + pass + except OSError: + # Keep failed deletions accounted for and eligible for retry. + with self._lock: + self._entry_index[path] = size + self._current_size += size def _scan_and_populate_index(self) -> None: for dirpath, _, filenames in os.walk(self._cache_dir): @@ -197,7 +225,7 @@ def _scan_and_populate_index(self) -> None: self._entry_index[fp] = size self._current_size += size # Bound the index while scanning an existing cache directory. - self._evict_locked() + self._delete_evicted(self._evict_locked()) except OSError: pass diff --git a/paimon-python/pypaimon/tests/blob_meta_cache_test.py b/paimon-python/pypaimon/tests/blob_meta_cache_test.py index 2f637830bc8c..f54839280973 100644 --- a/paimon-python/pypaimon/tests/blob_meta_cache_test.py +++ b/paimon-python/pypaimon/tests/blob_meta_cache_test.py @@ -286,3 +286,35 @@ def readinto(self, buffer): assert stream.readinto(memoryview(buf)[:2]) == 2 assert stream.tell() == 2 assert cache._current_size == 0 + + +@pytest.mark.parametrize('disk', [False, True], ids=['memory', 'disk']) +def test_oversized_map_keys_preserve_cached_manifest(tmp_path, disk): + field = DataField(0, 'value', MapType(True, AtomicType('STRING'), AtomicType('BLOB'))) + key = 'k' * (128 * 1024) + path = str(tmp_path / 'data.blob') + writer = BlobFormatWriter(open(path, 'wb')) + writer.add_element(GenericRow([[(key, BlobData(b'value'))]], [field])) + writer.close() + options = Options({'local-cache.enabled': 'true', 'local-cache.max-size': '64 kb', + **({'local-cache.dir': str(tmp_path / 'cache')} if disk else {})}) + cache = CachingFileIO.create_cache_manager(options) + cache.put_block('manifest', 0, b'manifest') + delegate = LocalFileIO(str(tmp_path), Options({})) + file_io = CachingFileIO.wrap_with_caching_if_needed(delegate, options, cache) + reader = FormatBlobReader(file_io, path, ['value'], [field], None, True) + try: + result = reader.read_arrow_batch().column(0)[0].as_py() + assert result[0][0] == key + descriptor = BlobDescriptor.deserialize(result[0][1]) + with open(path, 'rb') as stream: + stream.seek(descriptor.offset) + assert stream.read(descriptor.length) == b'value' + finally: + reader.close() + assert cache.get_block('manifest', 0) == b'manifest' + assert 0 < cache._current_size <= 64 * 1024 + if disk: + # Reject obviously oversized entries before opening a temporary file. + with patch('pypaimon.filesystem.caching_file_io.open', side_effect=AssertionError): + cache.put_block('oversized', 0, b'x' * (128 * 1024)) diff --git a/paimon-python/pypaimon/tests/caching_file_io_test.py b/paimon-python/pypaimon/tests/caching_file_io_test.py index 3f6450330fcc..8b945304a7c0 100644 --- a/paimon-python/pypaimon/tests/caching_file_io_test.py +++ b/paimon-python/pypaimon/tests/caching_file_io_test.py @@ -181,6 +181,40 @@ def wait_for_writers(path): files = [name for _, _, names in os.walk(self.cache_dir) for name in names] self.assertEqual([os.path.basename(path)], files) + def test_slow_eviction_allows_hits_and_same_key_republication(self): + cache = LocalDiskCacheManager(self.cache_dir, 8192) + cache.put_block("old", 0, b"old") + cache.put_block("keep", 0, b"cached") + deleting = threading.Event() + release = threading.Event() + unlink = os.unlink + + def slow_unlink(path): + if path.endswith('.evicted') and not deleting.is_set(): + deleting.set() + if not release.wait(5): + raise AssertionError("Timed out waiting to release deletion") + return unlink(path) + + with patch('pypaimon.filesystem.caching_file_io.os.unlink', side_effect=slow_unlink): + with ThreadPoolExecutor(max_workers=2) as executor: + writer = executor.submit(cache.put_block, "new", 0, b"new") + try: + self.assertTrue(deleting.wait(5)) + hit = executor.submit(cache.get_block, "keep", 0) + self.assertEqual(b"cached", hit.result(timeout=2)) + replacement = executor.submit(cache.put_block, "old", 0, b"replacement") + replacement.result(timeout=2) + self.assertFalse(writer.done()) + finally: + release.set() + writer.result(timeout=5) + self.assertEqual(b"replacement", cache.get_block("old", 0)) + self.assertLessEqual(cache._current_size, 8192) + files = [name for _, _, names in os.walk(self.cache_dir) for name in names] + self.assertEqual(2, len(files)) + self.assertFalse(any(name.endswith('.evicted') for name in files)) + def test_unlimited_cache_skips_eviction(self): cache = LocalDiskCacheManager(self.cache_dir, max_size_bytes=2 ** 63 - 1, block_size=64) for i in range(50): From b4a3ad82a6be311517531adb212fd75aed3cdc08 Mon Sep 17 00:00:00 2001 From: "xiaohongbo.xhb" Date: Sun, 27 Sep 2026 06:01:56 -0700 Subject: [PATCH 07/10] [python] Replace BLOB metadata caching with Parquet-only selection --- docs/docs/program-api/file-cache.mdx | 14 +- docs/docs/pypaimon/blob.md | 6 - .../pypaimon/common/options/core_options.py | 4 +- .../pypaimon/filesystem/caching_file_io.py | 217 ++++-------- .../read/reader/format_blob_reader.py | 25 +- .../pypaimon/tests/blob_meta_cache_test.py | 320 ------------------ .../pypaimon/tests/caching_file_io_test.py | 120 ++----- .../pypaimon/tests/file_type_test.py | 2 +- paimon-python/pypaimon/utils/file_type.py | 11 +- 9 files changed, 112 insertions(+), 607 deletions(-) delete mode 100644 paimon-python/pypaimon/tests/blob_meta_cache_test.py diff --git a/docs/docs/program-api/file-cache.mdx b/docs/docs/program-api/file-cache.mdx index be7746c5d3c4..ae7b7fe8d3e9 100644 --- a/docs/docs/program-api/file-cache.mdx +++ b/docs/docs/program-api/file-cache.mdx @@ -51,11 +51,8 @@ The cache classifies files by path. Set `local-cache.whitelist` to select the fi | BUCKET_INDEX | bucket-index | Hash, deletion vector index files | No | | DATA | data | Data files (ORC, Parquet, etc.) | No | | FILE_INDEX | file-index | Data-file level bloom filter, bitmap | No | -| BLOB_META | blob-meta | BLOB metadata read ranges (Python only) | Python only | -The Java default whitelist is `meta,global-index`. The Python reader additionally includes -`blob-meta`, which caches exact BLOB metadata ranges rather than whole data blocks; see -[Python BLOB metadata cache](../pypaimon/blob.md#blob-metadata-cache). The mutable `LATEST` and `EARLIEST` hint files bypass +The default whitelist is `meta,global-index`. The mutable `LATEST` and `EARLIEST` hint files bypass the cache even though they are classified as metadata. ## Enable Cache @@ -129,16 +126,15 @@ The Java snippet is a method-body fragment; declare or handle the catalog except | `local-cache.dir` | String | (none) | Directory for storing cached blocks on disk. If not configured, memory cache is used. | | `local-cache.max-size` | MemorySize | See below | Maximum cached block bytes per cache manager. Least recently used blocks are evicted when the limit is exceeded. | | `local-cache.block-size` | MemorySize | 1 mb | Block size for caching. Files are logically divided into fixed-size blocks and cached independently. | -| `local-cache.whitelist` | String | Java: `meta,global-index`; Python: `meta,global-index,blob-meta` | Comma-separated cache categories: `meta`, `global-index`, `bucket-index`, `data`, `file-index`, and Python-only `blob-meta`. | - -Python charges estimated object overhead for memory entries and allocated disk bytes (at least 4 KiB) -for disk entries, and limits each cache manager to 65,536 entries. These limits do not represent a -strict process RSS or total filesystem usage limit. +| `local-cache.whitelist` | String | meta,global-index | Comma-separated list of file types to cache. Supported values: `meta`, `global-index`, `bucket-index`, `data`, `file-index`; Python also supports `parquet-data`. | When `local-cache.max-size` is omitted, Java has no configured size limit. PyPaimon uses a 256 MiB memory limit or a 10 GiB disk limit. Set the option explicitly when you want the same limit across clients. The limit accounts for cached block bytes, not all reader buffers or process memory. +For PyPaimon, use `meta,global-index,parquet-data` to cache Parquet without BLOB bodies. +`data` still includes all data formats. + ## How It Works 1. The reader requests bytes from an eligible file. The cache maps the byte range to fixed-size diff --git a/docs/docs/pypaimon/blob.md b/docs/docs/pypaimon/blob.md index 76edc63ced4f..56126ef5b2da 100644 --- a/docs/docs/pypaimon/blob.md +++ b/docs/docs/pypaimon/blob.md @@ -176,12 +176,6 @@ The factory auto-dispatches based on the bytes content (`BLOBDESC`, `VIDEOFRM`, or blob-view magic header). This mirrors Java's `Blob.fromBytes(...)`. -## BLOB metadata cache - -Enable `local-cache.enabled=true` to cache BLOB metadata needed for descriptors, -including MAP keys, without caching value bodies. `blob-meta` is included in the -default whitelist and shares the local cache budget. Python reader only. - ## See Also - [Blob Storage](../multimodal-table/blob) — concept, storage modes, diff --git a/paimon-python/pypaimon/common/options/core_options.py b/paimon-python/pypaimon/common/options/core_options.py index 4c3ab71f36da..a1fb6b2abfbf 100644 --- a/paimon-python/pypaimon/common/options/core_options.py +++ b/paimon-python/pypaimon/common/options/core_options.py @@ -1138,10 +1138,10 @@ class CoreOptions: LOCAL_CACHE_WHITELIST: ConfigOption[str] = ( ConfigOptions.key("local-cache.whitelist") .string_type() - .default_value("meta,global-index,blob-meta") + .default_value("meta,global-index") .with_description( "Comma-separated list of file types to cache. " - "Supported values: meta, global-index, bucket-index, data, file-index, blob-meta." + "Supported values: meta, global-index, bucket-index, data, file-index." ) ) diff --git a/paimon-python/pypaimon/filesystem/caching_file_io.py b/paimon-python/pypaimon/filesystem/caching_file_io.py index 5e23c5c08741..83b8a2cca8b2 100644 --- a/paimon-python/pypaimon/filesystem/caching_file_io.py +++ b/paimon-python/pypaimon/filesystem/caching_file_io.py @@ -22,38 +22,21 @@ at block granularity. If a cache directory is configured, disk cache is used; otherwise an in-memory LRU cache is used. Files are classified by FileType and only cacheable types in the whitelist are cached; others are read directly from -the delegate FileIO. BLOB metadata can be cached as exact reader-selected ranges. +the delegate FileIO. """ import hashlib import os -import sys import threading -import uuid from collections import OrderedDict -from io import UnsupportedOperation -from typing import Optional, Tuple, Union +from typing import Optional from pypaimon.common.file_io import FileIO, supports_pread, pread from pypaimon.utils.file_type import FileType -_CacheEntryKey = Union[int, Tuple[str, int, int]] -# Bound Python index overhead and the number of cache files, even with an unlimited byte budget. -_MAX_CACHE_ENTRIES = 65536 - - -def _memory_entry_size(key, data): - # Conservatively count referenced objects plus OrderedDict bookkeeping. - path, entry_key = key - size = sys.getsizeof(key) + sys.getsizeof(path) + sys.getsizeof(entry_key) - if isinstance(entry_key, tuple): - size += sum(sys.getsizeof(part) for part in entry_key) - return max(512, size + sys.getsizeof(data) + 128) - - class LocalMemoryCacheManager: - """In-memory block/range cache with LRU eviction.""" + """Block-level in-memory cache with LRU eviction.""" def __init__(self, max_size_bytes: int, block_size: int = 1 * 1024 * 1024): self._max_size_bytes = max_size_bytes @@ -67,28 +50,26 @@ def __init__(self, max_size_bytes: int, block_size: int = 1 * 1024 * 1024): def block_size(self) -> int: return self._block_size - def get_block(self, file_path: str, cache_key: _CacheEntryKey) -> Optional[bytes]: - key = (file_path, cache_key) + def get_block(self, file_path: str, block_index: int) -> Optional[bytes]: + key = (file_path, block_index) with self._lock: data = self._cache.get(key) if data is not None: self._cache.move_to_end(key) return data - def put_block(self, file_path: str, cache_key: _CacheEntryKey, data: bytes) -> None: - key = (file_path, cache_key) - size = _memory_entry_size(key, data) - if size > self._max_size_bytes: - return + def put_block(self, file_path: str, block_index: int, data: bytes) -> None: + key = (file_path, block_index) with self._lock: if key in self._cache: return - self._current_size += size + self._current_size += len(data) self._cache[key] = data - while self._cache and (self._current_size > self._max_size_bytes - or len(self._cache) > _MAX_CACHE_ENTRIES): - evicted_key, evicted = self._cache.popitem(last=False) - self._current_size -= _memory_entry_size(evicted_key, evicted) + while (self._max_size_bytes < (2 ** 63 - 1) + and self._current_size > self._max_size_bytes + and self._cache): + _, evicted = self._cache.popitem(last=False) + self._current_size -= len(evicted) def get_file_size(self, file_path: str) -> int: return self._file_size_cache.get(file_path, -1) @@ -98,7 +79,7 @@ def put_file_size(self, file_path: str, size: int) -> None: class LocalDiskCacheManager: - """Local disk block/range cache with LRU eviction.""" + """Block-level local disk cache with LRU eviction.""" def __init__(self, cache_dir: str, max_size_bytes: int, block_size: int = 1 * 1024 * 1024): @@ -111,21 +92,21 @@ def __init__(self, cache_dir: str, max_size_bytes: int, # LRU-ordered index: cache_path -> size. OrderedDict with move_to_end for access order. self._entry_index: OrderedDict = OrderedDict() os.makedirs(cache_dir, exist_ok=True) - self._scan_and_populate_index() + self._current_size = self._scan_and_populate_index() @property def block_size(self) -> int: return self._block_size - def _cache_path(self, file_path: str, cache_key: _CacheEntryKey) -> str: - key = f"{file_path}:{cache_key}" + def _cache_path(self, file_path: str, block_index: int) -> str: + key = f"{file_path}:{block_index}" h = hashlib.sha256(key.encode('utf-8')).hexdigest() prefix = h[:2] sub_dir = os.path.join(self._cache_dir, prefix) return os.path.join(sub_dir, h) - def get_block(self, file_path: str, cache_key: _CacheEntryKey) -> Optional[bytes]: - path = self._cache_path(file_path, cache_key) + def get_block(self, file_path: str, block_index: int) -> Optional[bytes]: + path = self._cache_path(file_path, block_index) with self._lock: if path not in self._entry_index: return None @@ -140,94 +121,69 @@ def get_block(self, file_path: str, cache_key: _CacheEntryKey) -> Optional[bytes self._current_size -= size return None - def put_block(self, file_path: str, cache_key: _CacheEntryKey, data: bytes) -> None: - if max(len(data), 4096) > self._max_size_bytes: - return - path = self._cache_path(file_path, cache_key) + def put_block(self, file_path: str, block_index: int, data: bytes) -> None: + path = self._cache_path(file_path, block_index) + with self._lock: if path in self._entry_index: return + sub_dir = os.path.dirname(path) + os.makedirs(sub_dir, exist_ok=True) + tmp_path = path + f".tmp.{os.getpid()}.{threading.get_ident()}" try: - os.makedirs(os.path.dirname(path), exist_ok=True) with open(tmp_path, 'wb') as f: f.write(data) - size = self._disk_entry_size(tmp_path) - if size > self._max_size_bytes: - return - with self._lock: - # Another writer may have published this key during the write. - if path in self._entry_index: - return - os.rename(tmp_path, path) - self._entry_index[path] = size - self._current_size += size - evicted = self._evict_locked() - self._delete_evicted(evicted) - except OSError: - return - finally: + os.rename(tmp_path, path) + except Exception: try: os.unlink(tmp_path) except OSError: pass + return - @staticmethod - def _disk_entry_size(path: str) -> int: - stat = os.stat(path) - # st_blocks is in 512-byte units. Keep a minimum charge on platforms - # without allocation information and for sparse/empty files. - return max(stat.st_size, getattr(stat, 'st_blocks', 0) * 512, - getattr(stat, 'st_blksize', 4096), 4096) - - def _evict_locked(self): - evicted = [] - while (self._entry_index - and (self._current_size > self._max_size_bytes - or len(self._entry_index) > _MAX_CACHE_ENTRIES)): - path, size = self._entry_index.popitem(last=False) - # Detach the old file before unlocking, so deletion cannot remove a - # newly published entry for the same key. - retired = os.path.join(os.path.dirname(path), uuid.uuid4().hex + '.evicted') - try: - os.rename(path, retired) - except FileNotFoundError: - pass - except OSError: - self._entry_index[path] = size - break - else: - evicted.append((retired, size)) - self._current_size -= size - return evicted + need_evict = False + with self._lock: + self._entry_index[path] = len(data) + self._current_size += len(data) + need_evict = (self._max_size_bytes < (2 ** 63 - 1) + and self._current_size > self._max_size_bytes) + if need_evict: + self._evict() + + def _evict(self) -> None: + to_delete = [] + with self._lock: + if self._current_size <= self._max_size_bytes: + return + while self._entry_index and self._current_size > self._max_size_bytes: + path, size = self._entry_index.popitem(last=False) + self._current_size -= size + to_delete.append((path, size)) - def _delete_evicted(self, evicted) -> None: - for path, size in evicted: + for path, size in to_delete: try: os.unlink(path) - except FileNotFoundError: - pass except OSError: - # Keep failed deletions accounted for and eligible for retry. with self._lock: self._entry_index[path] = size self._current_size += size - def _scan_and_populate_index(self) -> None: + def _scan_and_populate_index(self) -> int: + total = 0 for dirpath, _, filenames in os.walk(self._cache_dir): for fn in filenames: if '.tmp.' in fn: continue fp = os.path.join(dirpath, fn) try: - size = self._disk_entry_size(fp) + size = os.path.getsize(fp) self._entry_index[fp] = size - self._current_size += size - # Bound the index while scanning an existing cache directory. - self._delete_evicted(self._evict_locked()) + total += size except OSError: pass + return total def get_file_size(self, file_path: str) -> int: return self._file_size_cache.get(file_path, -1) @@ -295,14 +251,6 @@ def read(self, size=-1) -> bytes: self._pos = end return bytes(result) - def readinto(self, buffer) -> int: - view = memoryview(buffer).cast('B') - if view.readonly: - raise TypeError("readinto() requires a writable buffer") - data = self.read(len(view)) - view[:len(data)] = data - return len(data) - def read_at(self, nbytes: int, offset: int) -> bytes: """Position-based read. Does not change the cursor. Thread-safe.""" if nbytes <= 0 or offset >= self._get_file_size(): @@ -377,54 +325,6 @@ def __exit__(self, exc_type, exc_val, exc_tb): return False -class BlobMetadataInputStream(CachingInputStream): - """Cache only explicitly marked ranges; ordinary reads bypass the cache.""" - - def read(self, size=-1) -> bytes: - if size is None or size < 0: - size = self._get_file_size() - self._pos - data = self.read_at(size, self._pos) - self._pos += len(data) - return data - - def readinto(self, buffer) -> Optional[int]: - view = memoryview(buffer).cast('B') - if view.readonly: - raise TypeError("readinto() requires a writable buffer") - if not view: - return 0 - with self._io_lock: - stream = self._get_remote_stream() - readinto = getattr(stream, 'readinto', None) - if callable(readinto): - stream.seek(self._pos) - try: - count = readinto(view) - except (UnsupportedOperation, NotImplementedError): - pass - else: - if count is not None: - self._pos += count - return count - return super().readinto(view) - - def read_at(self, nbytes: int, offset: int) -> bytes: - return self._read_remote(offset, nbytes) if nbytes > 0 else b'' - - def read_blob_metadata(self, size: int) -> bytes: - if size <= 0: - return b'' - # Tuple keys cannot collide with integer block indexes, in memory or on disk. - key = ('blob-meta', self._pos, size) - data = self._cache.get_block(self._file_path, key) - if data is None or len(data) != size: - data = self._read_remote(self._pos, size) - if len(data) == size: - self._cache.put_block(self._file_path, key, data) - self._pos += len(data) - return data - - class CachingFileIO(FileIO): """FileIO wrapper that caches reads at block granularity. @@ -437,7 +337,7 @@ def __init__(self, delegate: FileIO, cache, whitelist=None): self._delegate = delegate self._cache = cache if whitelist is None: - self._whitelist = {FileType.META, FileType.GLOBAL_INDEX, FileType.BLOB_META} + self._whitelist = {FileType.META, FileType.GLOBAL_INDEX} else: self._whitelist = whitelist @@ -498,12 +398,11 @@ def properties(self): def new_input_stream(self, path: str): file_type = FileType.classify(path) - if self._cache is not None and not FileType.is_mutable(path): - if file_type in self._whitelist: - return CachingInputStream(self._delegate, path, self._cache) - if FileType.BLOB_META in self._whitelist and path.endswith('.blob'): - return BlobMetadataInputStream(self._delegate, path, self._cache) - return self._delegate.new_input_stream(path) + eligible = (file_type in self._whitelist + or (file_type == FileType.PARQUET_DATA and FileType.DATA in self._whitelist)) + if self._cache is None or not eligible or FileType.is_mutable(path): + return self._delegate.new_input_stream(path) + return CachingInputStream(self._delegate, path, self._cache) def new_output_stream(self, path: str): return self._delegate.new_output_stream(path) diff --git a/paimon-python/pypaimon/read/reader/format_blob_reader.py b/paimon-python/pypaimon/read/reader/format_blob_reader.py index bd0dbd5647a2..e39e08b7b229 100644 --- a/paimon-python/pypaimon/read/reader/format_blob_reader.py +++ b/paimon-python/pypaimon/read/reader/format_blob_reader.py @@ -387,7 +387,7 @@ def _read_index(self) -> None: # Seek to header: last 5 bytes f.seek(self._file_size - 5) - header = BlobRecordIterator._read_fully_from(f, 5, metadata=True) + header = f.read(5) if len(header) != 5: raise IOError("Invalid blob file: cannot read header") @@ -401,7 +401,7 @@ def _read_index(self) -> None: # Read index data f.seek(self._file_size - 5 - index_length) - index_bytes = BlobRecordIterator._read_fully_from(f, index_length, metadata=True) + index_bytes = f.read(index_length) if len(index_bytes) != index_length: raise IOError("Invalid blob file: cannot read index") @@ -543,7 +543,7 @@ def _read_blob_array(self, position: int, length: int): close_stream = True try: stream.seek(position) - header = self._read_fully_from(stream, self.ARRAY_HEADER_SIZE, metadata=True) + header = self._read_fully_from(stream, self.ARRAY_HEADER_SIZE) if len(header) != self.ARRAY_HEADER_SIZE: raise IOError("Invalid ARRAY payload: cannot read header") magic, version, element_count = struct.unpack(' payload: cannot read index length") index_length = struct.unpack(' payload: cannot read element index") self._validate_array_element_index(index_bytes) @@ -645,7 +645,7 @@ def _read_blob_map(self, position: int, length: int): close_stream = True try: stream.seek(position) - header = self._read_fully_from(stream, self.MAP_HEADER_SIZE, metadata=True) + header = self._read_fully_from(stream, self.MAP_HEADER_SIZE) if len(header) != self.MAP_HEADER_SIZE: raise IOError("Invalid MAP payload: cannot read header") magic, version, entry_count = struct.unpack(' payload: cannot read index lengths") @@ -678,7 +678,7 @@ def _read_blob_map(self, position: int, length: int): stream.seek(key_index_start) # The two indexes are adjacent; read both without touching BLOB values. index_bytes = self._read_fully_from( - stream, key_index_length + value_index_length, metadata=True) + stream, key_index_length + value_index_length) key_index_bytes = index_bytes[:key_index_length] if len(key_index_bytes) != key_index_length: raise IOError("Invalid MAP payload: cannot read key index") @@ -707,7 +707,7 @@ def _read_blob_map(self, position: int, length: int): ) stream.seek(data_start) - key_data = self._read_fully_from(stream, key_data_length, metadata=True) + key_data = self._read_fully_from(stream, key_data_length) if len(key_data) != key_data_length: raise IOError("Invalid MAP payload: cannot read key data") keys = [] @@ -866,13 +866,10 @@ def _read_fully(self, length: int) -> bytes: return self._read_fully_from(self.input_stream, length) @staticmethod - def _read_fully_from(stream, length: int, metadata: bool = False) -> bytes: - read = stream.read - if metadata and callable(getattr(type(stream), 'read_blob_metadata', None)): - read = stream.read_blob_metadata + def _read_fully_from(stream, length: int) -> bytes: data = bytearray() while len(data) < length: - chunk = read(length - len(data)) + chunk = stream.read(length - len(data)) if not chunk: break data.extend(chunk) diff --git a/paimon-python/pypaimon/tests/blob_meta_cache_test.py b/paimon-python/pypaimon/tests/blob_meta_cache_test.py deleted file mode 100644 index f54839280973..000000000000 --- a/paimon-python/pypaimon/tests/blob_meta_cache_test.py +++ /dev/null @@ -1,320 +0,0 @@ -# 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 io -from unittest.mock import patch - -import pytest - -from pypaimon.common.options.options import Options -from pypaimon.filesystem.caching_file_io import CachingFileIO, CachingInputStream -from pypaimon.filesystem.local_file_io import LocalFileIO -from pypaimon.read.reader.format_blob_reader import ( - FormatBlobReader, _BLOB_INDEX_CACHE, _BLOB_INDEX_CACHE_LOCK, -) -from pypaimon.schema.data_types import ArrayType, AtomicType, DataField, MapType -from pypaimon.table.row.blob import BlobData, BlobDescriptor -from pypaimon.table.row.generic_row import GenericRow -from pypaimon.utils.file_type import FileType -from pypaimon.write.blob_format_writer import BlobFormatWriter - - -@pytest.mark.parametrize('disk', [False, True], ids=['memory', 'disk']) -@pytest.mark.parametrize('kind', ['scalar', 'array', 'map']) -def test_blob_metadata_ranges_exclude_values(tmp_path, disk, kind): - payload = b'body' * 4096 - blob = AtomicType('BLOB') - types = {'scalar': blob, 'array': ArrayType(True, blob), - 'map': MapType(True, AtomicType('STRING'), blob)} - values = {'scalar': BlobData(payload), - 'array': [BlobData(payload), None, BlobData(b'')], - 'map': [('key', BlobData(payload)), ('missing', None), ('empty', BlobData(b''))]} - field = DataField(0, 'value', types[kind]) - path = str(tmp_path / 'data.blob') - writer = BlobFormatWriter(open(path, 'wb')) - writer.add_element(GenericRow([values[kind]], [field])) - writer.close() - contents = (tmp_path / 'data.blob').read_bytes() - options = Options({'local-cache.enabled': 'true', 'local-cache.max-size': '1 mb', - **({'local-cache.dir': str(tmp_path / 'cache')} if disk else {})}) - delegate = LocalFileIO(str(tmp_path), Options({})) - cache = CachingFileIO.create_cache_manager(options) - file_io = CachingFileIO.wrap_with_caching_if_needed(delegate, options, cache) - reads = [] - - class CountingStream(io.BytesIO): - def read(self, size=-1): - start = self.tell() - result = super().read(size) - reads.append((start, start + len(result))) - return result - - def read(descriptors=True): - # Isolate the byte cache from the existing parsed-index cache. - with _BLOB_INDEX_CACHE_LOCK: - _BLOB_INDEX_CACHE.clear() - reader = FormatBlobReader(file_io, path, ['value'], [field], None, descriptors, - file_size=len(contents)) - try: - return reader.read_arrow_batch().column(0)[0].as_py() - finally: - reader.close() - - with patch.object(delegate, 'new_input_stream', side_effect=lambda _: CountingStream(contents)): - first = read() - serialized = ([first] if kind == 'scalar' else first if kind == 'array' - else [value for _, value in first]) - descriptors = [BlobDescriptor.deserialize(value) for value in serialized if value is not None] - assert reads - for descriptor in descriptors: - if descriptor.length: - assert all(end <= descriptor.offset or start >= descriptor.offset + descriptor.length - for start, end in reads) - retained = cache._current_size - assert retained > 0 - reads.clear() - if disk: - # Reopen the disk manager to verify persistence, not just memory hits. - cache = CachingFileIO.create_cache_manager(options) - file_io = CachingFileIO.wrap_with_caching_if_needed(delegate, options, cache) - assert read() == first - assert reads == [] - materialized = read(False) - expected = {'scalar': payload, 'array': [payload, None, b''], - 'map': [('key', payload), ('missing', None), ('empty', b'')]}[kind] - assert materialized == expected - assert reads # Values still come from the delegate. - assert cache._current_size == retained - - -def test_blob_meta_is_a_read_category_not_a_file_type(tmp_path): - delegate = LocalFileIO(str(tmp_path), Options({})) - options = Options({'local-cache.enabled': 'true', 'local-cache.whitelist': 'blob-meta'}) - cache = CachingFileIO.create_cache_manager(options) - file_io = CachingFileIO.wrap_with_caching_if_needed(delegate, options, cache) - assert FileType.classify('data.blob') == FileType.DATA - assert FileType.parse_whitelist('blob-meta') == {FileType.BLOB_META} - for name in ('data.parquet', '.data.blob.uuid.tmp'): - (tmp_path / name).write_bytes(b'data') - with file_io.new_input_stream(str(tmp_path / name)) as stream: - assert stream.read() == b'data' - assert cache._current_size == 0 - - -@pytest.mark.parametrize('disk', [False, True]) -def test_metadata_cache_budget_and_data_cache_coexist(tmp_path, disk): - path = str(tmp_path / 'data.blob') - (tmp_path / 'data.blob').write_bytes(b'abcdefgh') - delegate = LocalFileIO(str(tmp_path), Options({})) - options = Options({'local-cache.enabled': 'true', 'local-cache.whitelist': 'blob-meta', - 'local-cache.max-size': '4096 b', 'local-cache.block-size': '4 b', - **({'local-cache.dir': str(tmp_path / 'cache')} if disk else {})}) - cache = CachingFileIO.create_cache_manager(options) - file_io = CachingFileIO.wrap_with_caching_if_needed(delegate, options, cache) - with file_io.new_input_stream(path) as stream: - assert stream.read_blob_metadata(4) == b'abcd' - assert stream.read_blob_metadata(4) == b'efgh' - assert 0 < cache._current_size <= 4096 - retained = cache._current_size - assert stream.read_blob_metadata(4) == b'' # Never cache a short read. - assert cache._current_size == retained - stream.seek(0) - assert stream.read_blob_metadata(4) == b'abcd' - assert 0 < cache._current_size <= 4096 - # Reusing a manager for normal data blocks cannot hit metadata range entries. - data_io = CachingFileIO(delegate, cache, {FileType.DATA, FileType.BLOB_META}) - with data_io.new_input_stream(path) as stream: - assert type(stream) is CachingInputStream - assert stream.read() == b'abcdefgh' - - -@pytest.mark.parametrize('disk', [False, True], ids=['memory', 'disk']) -@pytest.mark.parametrize('kind', ['array', 'map']) -def test_many_blob_rows_respect_cache_budget(tmp_path, disk, kind): - blob = AtomicType('BLOB') - field = DataField(0, 'value', ArrayType(True, blob) if kind == 'array' - else MapType(True, AtomicType('STRING'), blob)) - value = [BlobData(b'x')] if kind == 'array' else [('key', BlobData(b'x'))] - path = str(tmp_path / 'data.blob') - writer = BlobFormatWriter(open(path, 'wb')) - for _ in range(1000): - writer.add_element(GenericRow([value], [field])) - writer.close() - budget = 64 * 1024 - options = Options({'local-cache.enabled': 'true', 'local-cache.max-size': '64 kb', - **({'local-cache.dir': str(tmp_path / 'cache')} if disk else {})}) - cache = CachingFileIO.create_cache_manager(options) - delegate = LocalFileIO(str(tmp_path), Options({})) - file_io = CachingFileIO.wrap_with_caching_if_needed(delegate, options, cache) - reader = FormatBlobReader(file_io, path, ['value'], [field], None, True) - try: - batch = reader.read_arrow_batch() - assert batch.num_rows == 1000 - for row in batch.column(0).to_pylist(): - raw = row[0] if kind == 'array' else row[0][1] - descriptor = BlobDescriptor.deserialize(raw) - with open(path, 'rb') as stream: - stream.seek(descriptor.offset) - assert stream.read(descriptor.length) == b'x' - finally: - reader.close() - assert 0 < cache._current_size <= budget - entries = cache._entry_index if disk else cache._cache - assert 0 < len(entries) <= budget // (4096 if disk else 512) - if disk: - files = [p for p in (tmp_path / 'cache').rglob('*') if p.is_file()] - assert len(files) == len(entries) - assert sum(p.stat().st_blocks * 512 for p in files) <= budget - reopened = CachingFileIO.create_cache_manager(options) - assert reopened._current_size == cache._current_size - assert len(reopened._entry_index) == len(entries) - - -@pytest.mark.parametrize('disk', [False, True], ids=['memory', 'disk']) -def test_tiny_ranges_respect_entry_limit_and_lru(tmp_path, disk): - options = Options({'local-cache.enabled': 'true', 'local-cache.max-size': '1 gb', - **({'local-cache.dir': str(tmp_path / 'cache')} if disk else {})}) - cache = CachingFileIO.create_cache_manager(options) - with patch('pypaimon.filesystem.caching_file_io._MAX_CACHE_ENTRIES', 32): - for offset in range(32): - cache.put_block('data.blob', ('blob-meta', offset, 1), b'x') - assert cache.get_block('data.blob', ('blob-meta', 0, 1)) == b'x' - cache.put_block('data.blob', ('blob-meta', 32, 1), b'x') - assert cache.get_block('data.blob', ('blob-meta', 0, 1)) == b'x' - assert cache.get_block('data.blob', ('blob-meta', 1, 1)) is None - for offset in range(33, 1000): - cache.put_block('data.blob', ('blob-meta', offset, 1), b'x') - entries = cache._entry_index if disk else cache._cache - assert len(entries) == 32 - assert cache._current_size >= 32 * (4096 if disk else 512) - if disk: - # Opening an existing directory must enforce both limits immediately. - with patch('pypaimon.filesystem.caching_file_io._MAX_CACHE_ENTRIES', 8): - reopened = CachingFileIO.create_cache_manager(options) - assert len(reopened._entry_index) == 8 - options = Options({'local-cache.enabled': 'true', 'local-cache.max-size': '4 kb', - 'local-cache.dir': str(tmp_path / 'cache')}) - reopened = CachingFileIO.create_cache_manager(options) - assert reopened._current_size <= 4096 - assert len(reopened._entry_index) <= 1 - assert len([p for p in (tmp_path / 'cache').rglob('*') if p.is_file()]) <= 1 - - -@pytest.mark.parametrize('whitelist', ['blob-meta', 'data']) -def test_blob_readinto(tmp_path, whitelist): - from pypaimon.table.row.blob import BlobRef - - path = str(tmp_path / 'data.blob') - (tmp_path / 'data.blob').write_bytes(b'abcdefgh') - delegate = LocalFileIO(str(tmp_path), Options({})) - options = Options({'local-cache.enabled': 'true', 'local-cache.whitelist': whitelist}) - cache = CachingFileIO.create_cache_manager(options) - file_io = CachingFileIO.wrap_with_caching_if_needed(delegate, options, cache) - blob = BlobRef(file_io, BlobDescriptor(path, 1, 4)) - with blob.new_input_stream() as stream: - assert stream.readinto(bytearray()) == 0 - buf = bytearray(b'------') - assert stream.readinto(buf) == 4 - assert buf == b'bcde--' - assert stream.readinto(buf) == 0 - stream.seek(0) - assert stream.readinto(memoryview(buf)[1:3]) == 2 - assert buf == b'bbce--' - assert stream.tell() == 2 - if whitelist == 'blob-meta': - assert cache._current_size == 0 # readinto must not cache value bodies. - else: - assert cache._current_size > 0 - - -@pytest.mark.parametrize('support', ['direct', 'missing', 'unsupported']) -def test_blob_readinto_delegates_without_caching(tmp_path, support): - from pypaimon.table.row.blob import BlobRef - - calls = {'read': 0, 'readinto': 0} - - class CountingStream(io.BytesIO): - def read(self, size=-1): - calls['read'] += 1 - return super().read(size) - - def readinto(self, buffer): - calls['readinto'] += 1 - if support == 'unsupported': - raise io.UnsupportedOperation('readinto') - # Exercise a short read and cursor updates. - return super().readinto(memoryview(buffer)[:2]) - - remote = CountingStream(b'abcdefgh') - if support == 'missing': - remote.readinto = None - delegate = LocalFileIO(str(tmp_path), Options({})) - options = Options({'local-cache.enabled': 'true'}) - cache = CachingFileIO.create_cache_manager(options) - file_io = CachingFileIO.wrap_with_caching_if_needed(delegate, options, cache) - blob = BlobRef(file_io, BlobDescriptor('data.blob', 1, 4)) - with patch.object(delegate, 'new_input_stream', return_value=remote): - with blob.new_input_stream() as stream: - buf = bytearray(b'----') - if support == 'direct': - assert stream.readinto(buf) == 2 - assert buf == b'bc--' - assert stream.tell() == 2 - assert stream.readinto(memoryview(buf)[2:]) == 2 - assert calls == {'read': 0, 'readinto': 2} - else: - assert stream.readinto(buf) == 4 - assert calls['read'] > 0 - assert buf == b'bcde' - assert stream.tell() == 4 - assert stream.readinto(buf) == 0 - stream.seek(0) - assert stream.readinto(memoryview(buf)[:2]) == 2 - assert stream.tell() == 2 - assert cache._current_size == 0 - - -@pytest.mark.parametrize('disk', [False, True], ids=['memory', 'disk']) -def test_oversized_map_keys_preserve_cached_manifest(tmp_path, disk): - field = DataField(0, 'value', MapType(True, AtomicType('STRING'), AtomicType('BLOB'))) - key = 'k' * (128 * 1024) - path = str(tmp_path / 'data.blob') - writer = BlobFormatWriter(open(path, 'wb')) - writer.add_element(GenericRow([[(key, BlobData(b'value'))]], [field])) - writer.close() - options = Options({'local-cache.enabled': 'true', 'local-cache.max-size': '64 kb', - **({'local-cache.dir': str(tmp_path / 'cache')} if disk else {})}) - cache = CachingFileIO.create_cache_manager(options) - cache.put_block('manifest', 0, b'manifest') - delegate = LocalFileIO(str(tmp_path), Options({})) - file_io = CachingFileIO.wrap_with_caching_if_needed(delegate, options, cache) - reader = FormatBlobReader(file_io, path, ['value'], [field], None, True) - try: - result = reader.read_arrow_batch().column(0)[0].as_py() - assert result[0][0] == key - descriptor = BlobDescriptor.deserialize(result[0][1]) - with open(path, 'rb') as stream: - stream.seek(descriptor.offset) - assert stream.read(descriptor.length) == b'value' - finally: - reader.close() - assert cache.get_block('manifest', 0) == b'manifest' - assert 0 < cache._current_size <= 64 * 1024 - if disk: - # Reject obviously oversized entries before opening a temporary file. - with patch('pypaimon.filesystem.caching_file_io.open', side_effect=AssertionError): - cache.put_block('oversized', 0, b'x' * (128 * 1024)) diff --git a/paimon-python/pypaimon/tests/caching_file_io_test.py b/paimon-python/pypaimon/tests/caching_file_io_test.py index 8b945304a7c0..9eeff96b2e8b 100644 --- a/paimon-python/pypaimon/tests/caching_file_io_test.py +++ b/paimon-python/pypaimon/tests/caching_file_io_test.py @@ -25,8 +25,7 @@ import tempfile import threading import unittest -from concurrent.futures import ThreadPoolExecutor -from unittest.mock import MagicMock, patch +from unittest.mock import MagicMock from pypaimon.filesystem.caching_file_io import ( LocalDiskCacheManager, @@ -87,8 +86,7 @@ def test_scan_size_on_restart(self): # Simulate restart: new cache instance on same directory cache2 = LocalDiskCacheManager(self.cache_dir, 2 ** 63 - 1, block_size=64) - self.assertEqual(cache1._current_size, cache2._current_size) - self.assertGreaterEqual(cache2._current_size, 8192) + self.assertEqual(300, cache2._current_size) self.assertEqual(b"x" * 100, cache2.get_block("f", 0)) self.assertEqual(b"y" * 200, cache2.get_block("f", 1)) @@ -128,93 +126,6 @@ def reader(idx): self.assertEqual([], errors) - def test_slow_write_does_not_block_unrelated_hit(self): - cache = LocalDiskCacheManager(self.cache_dir, 1024 * 1024) - cache.put_block("existing", 0, b"cached") - writing = threading.Event() - release = threading.Event() - real_open = open - - class SlowFile(io.FileIO): - def write(self, data): - writing.set() - if not release.wait(5): - raise AssertionError("Timed out waiting to release write") - return super().write(data) - - def open_file(path, mode): - return SlowFile(path, mode) if mode == 'wb' else real_open(path, mode) - - with patch('pypaimon.filesystem.caching_file_io.open', side_effect=open_file): - with ThreadPoolExecutor(max_workers=2) as executor: - writer = executor.submit(cache.put_block, "new", 0, b"new data") - try: - self.assertTrue(writing.wait(5)) - hit = executor.submit(cache.get_block, "existing", 0) - self.assertEqual(b"cached", hit.result(timeout=2)) - self.assertFalse(writer.done()) - finally: - release.set() - writer.result(timeout=5) - self.assertEqual(b"new data", cache.get_block("new", 0)) - - def test_concurrent_same_key_is_published_once(self): - cache = LocalDiskCacheManager(self.cache_dir, 1024 * 1024) - ready = threading.Barrier(2) - entry_size = cache._disk_entry_size - - def wait_for_writers(path): - size = entry_size(path) - ready.wait(timeout=5) - return size - - with patch.object(cache, '_disk_entry_size', side_effect=wait_for_writers): - with ThreadPoolExecutor(max_workers=2) as executor: - writers = [executor.submit(cache.put_block, "same", 0, data) - for data in (b"first", b"second")] - for writer in writers: - writer.result(timeout=5) - self.assertIn(cache.get_block("same", 0), (b"first", b"second")) - self.assertEqual(1, len(cache._entry_index)) - path = cache._cache_path("same", 0) - self.assertEqual(entry_size(path), cache._current_size) - files = [name for _, _, names in os.walk(self.cache_dir) for name in names] - self.assertEqual([os.path.basename(path)], files) - - def test_slow_eviction_allows_hits_and_same_key_republication(self): - cache = LocalDiskCacheManager(self.cache_dir, 8192) - cache.put_block("old", 0, b"old") - cache.put_block("keep", 0, b"cached") - deleting = threading.Event() - release = threading.Event() - unlink = os.unlink - - def slow_unlink(path): - if path.endswith('.evicted') and not deleting.is_set(): - deleting.set() - if not release.wait(5): - raise AssertionError("Timed out waiting to release deletion") - return unlink(path) - - with patch('pypaimon.filesystem.caching_file_io.os.unlink', side_effect=slow_unlink): - with ThreadPoolExecutor(max_workers=2) as executor: - writer = executor.submit(cache.put_block, "new", 0, b"new") - try: - self.assertTrue(deleting.wait(5)) - hit = executor.submit(cache.get_block, "keep", 0) - self.assertEqual(b"cached", hit.result(timeout=2)) - replacement = executor.submit(cache.put_block, "old", 0, b"replacement") - replacement.result(timeout=2) - self.assertFalse(writer.done()) - finally: - release.set() - writer.result(timeout=5) - self.assertEqual(b"replacement", cache.get_block("old", 0)) - self.assertLessEqual(cache._current_size, 8192) - files = [name for _, _, names in os.walk(self.cache_dir) for name in names] - self.assertEqual(2, len(files)) - self.assertFalse(any(name.endswith('.evicted') for name in files)) - def test_unlimited_cache_skips_eviction(self): cache = LocalDiskCacheManager(self.cache_dir, max_size_bytes=2 ** 63 - 1, block_size=64) for i in range(50): @@ -425,6 +336,31 @@ def test_custom_whitelist_meta_only(self): result = caching_io.new_input_stream("global-index-uuid.index") self.assertNotIsInstance(result, CachingInputStream) + def test_parquet_only_cache_and_data_compatibility(self): + from pypaimon.filesystem.caching_file_io import LocalMemoryCacheManager + from pypaimon.utils.file_type import FileType + + for disk in (False, True): + for whitelist in ('meta,global-index', 'parquet-data', 'data'): + for name in ('data.parquet', 'data.blob', 'data.orc', 'data.parquet.index'): + with self.subTest(disk=disk, whitelist=whitelist, name=name): + with tempfile.TemporaryDirectory() as directory: + cache = (LocalDiskCacheManager(directory, 1024, block_size=4) if disk + else LocalMemoryCacheManager(1024, block_size=4)) + delegate = self._make_delegate({name: b'abcdefgh'}) + caching_io = CachingFileIO( + delegate, cache, FileType.parse_whitelist(whitelist)) + cached = ((whitelist == 'parquet-data' and name.endswith('.parquet')) + or (whitelist == 'data' and not name.endswith('.index'))) + for _ in range(2): + with caching_io.new_input_stream(name) as stream: + stream.seek(1) + self.assertEqual(b'bcdefg', stream.read(6)) + self.assertEqual(1 if cached else 2, + delegate.new_input_stream.call_count) + if not cached: + delegate.get_file_size.assert_not_called() + def test_data_file_not_cached(self): data = b"data content" delegate = self._make_delegate({"data-abc.orc": data}) @@ -581,7 +517,7 @@ def test_local_cache_options_defaults(self): self.assertIsNone(opts.local_cache_dir()) self.assertIsNone(opts.local_cache_max_size()) self.assertEqual(1 * 1024 * 1024, opts.local_cache_block_size().get_bytes()) - self.assertEqual("meta,global-index,blob-meta", opts.local_cache_whitelist()) + self.assertEqual("meta,global-index", opts.local_cache_whitelist()) def test_local_cache_options_custom(self): from pypaimon.common.options import Options diff --git a/paimon-python/pypaimon/tests/file_type_test.py b/paimon-python/pypaimon/tests/file_type_test.py index 0efabd1833ee..f17293184111 100644 --- a/paimon-python/pypaimon/tests/file_type_test.py +++ b/paimon-python/pypaimon/tests/file_type_test.py @@ -84,7 +84,7 @@ def test_bucket_index(self): def test_data_files(self): self.assertEqual(FileType.DATA, FileType.classify("data-abc.orc")) - self.assertEqual(FileType.DATA, FileType.classify("data-abc.parquet")) + self.assertEqual(FileType.PARQUET_DATA, FileType.classify("data-abc.parquet")) self.assertEqual(FileType.DATA, FileType.classify("unknown-file")) def test_temp_file_unwrap(self): diff --git a/paimon-python/pypaimon/utils/file_type.py b/paimon-python/pypaimon/utils/file_type.py index 61791ba6b795..6f74152c4d1c 100644 --- a/paimon-python/pypaimon/utils/file_type.py +++ b/paimon-python/pypaimon/utils/file_type.py @@ -31,17 +31,17 @@ class FileType(Enum): - META: snapshot, schema, manifest, manifest sidecar, statistics, tag, changelog metadata, hint files, _SUCCESS, consumer, service files - DATA: data files and any unrecognized files (default) + - PARQUET_DATA: Parquet data files - BUCKET_INDEX: bucket level index files (Hash, DV) - GLOBAL_INDEX: table level global index files (btree, lumina, full-text) - FILE_INDEX: data-file index files (bloom filter, bitmap, etc.) - - BLOB_META: reader-selected metadata ranges within BLOB files """ META = "META" DATA = "DATA" + PARQUET_DATA = "PARQUET_DATA" BUCKET_INDEX = "BUCKET_INDEX" GLOBAL_INDEX = "GLOBAL_INDEX" FILE_INDEX = "FILE_INDEX" - BLOB_META = "BLOB_META" def is_index(self) -> bool: return self in (FileType.BUCKET_INDEX, FileType.GLOBAL_INDEX, FileType.FILE_INDEX) @@ -86,6 +86,9 @@ def classify(file_path: str) -> 'FileType': if parent == "changelog": return FileType.META + if name.endswith(".parquet"): + return FileType.PARQUET_DATA + return FileType.DATA @staticmethod @@ -95,8 +98,8 @@ def parse_whitelist(whitelist_str: str) -> set: "global-index": FileType.GLOBAL_INDEX, "bucket-index": FileType.BUCKET_INDEX, "data": FileType.DATA, + "parquet-data": FileType.PARQUET_DATA, "file-index": FileType.FILE_INDEX, - "blob-meta": FileType.BLOB_META, } result = set() for name in whitelist_str.split(","): @@ -106,7 +109,7 @@ def parse_whitelist(whitelist_str: str) -> set: elif name: logger.warning( "Unknown local-cache.whitelist value '%s'. " - "Supported values: meta, global-index, bucket-index, data, file-index, blob-meta.", + "Supported values: meta, global-index, bucket-index, data, parquet-data, file-index.", name, ) return result From eb1092c44f9192b968b681ce60a36270f4159a36 Mon Sep 17 00:00:00 2001 From: "xiaohongbo.xhb" Date: Mon, 28 Sep 2026 02:23:07 -0700 Subject: [PATCH 08/10] [python] Cache data reads with extension exclusions --- docs/docs/program-api/file-cache.mdx | 7 +- .../pypaimon/common/options/core_options.py | 10 ++ .../pypaimon/filesystem/caching_file_io.py | 46 ++++++-- .../filesystem/caching_file_system.py | 78 +++++++++++++ .../read/reader/format_pyarrow_reader.py | 30 +++-- .../pypaimon/tests/caching_file_io_test.py | 23 +++- .../pypaimon/tests/file_type_test.py | 2 +- .../tests/parquet_local_cache_test.py | 105 ++++++++++++++++++ paimon-python/pypaimon/utils/file_type.py | 8 +- 9 files changed, 274 insertions(+), 35 deletions(-) create mode 100644 paimon-python/pypaimon/filesystem/caching_file_system.py create mode 100644 paimon-python/pypaimon/tests/parquet_local_cache_test.py diff --git a/docs/docs/program-api/file-cache.mdx b/docs/docs/program-api/file-cache.mdx index ae7b7fe8d3e9..e6f194230d0a 100644 --- a/docs/docs/program-api/file-cache.mdx +++ b/docs/docs/program-api/file-cache.mdx @@ -126,14 +126,15 @@ The Java snippet is a method-body fragment; declare or handle the catalog except | `local-cache.dir` | String | (none) | Directory for storing cached blocks on disk. If not configured, memory cache is used. | | `local-cache.max-size` | MemorySize | See below | Maximum cached block bytes per cache manager. Least recently used blocks are evicted when the limit is exceeded. | | `local-cache.block-size` | MemorySize | 1 mb | Block size for caching. Files are logically divided into fixed-size blocks and cached independently. | -| `local-cache.whitelist` | String | meta,global-index | Comma-separated list of file types to cache. Supported values: `meta`, `global-index`, `bucket-index`, `data`, `file-index`; Python also supports `parquet-data`. | +| `local-cache.whitelist` | String | meta,global-index | Comma-separated list of file types to cache. Supported values: `meta`, `global-index`, `bucket-index`, `data`, `file-index`. | +| `local-cache.exclude-extensions` | String | (empty) | Comma-separated file extensions to bypass, even when whitelisted. Supported in PyPaimon and paimon-rust. | When `local-cache.max-size` is omitted, Java has no configured size limit. PyPaimon uses a 256 MiB memory limit or a 10 GiB disk limit. Set the option explicitly when you want the same limit across clients. The limit accounts for cached block bytes, not all reader buffers or process memory. -For PyPaimon, use `meta,global-index,parquet-data` to cache Parquet without BLOB bodies. -`data` still includes all data formats. +Use `local-cache.whitelist=meta,global-index,data` with `local-cache.exclude-extensions=blob` +to cache data files except BLOB bodies. ## How It Works diff --git a/paimon-python/pypaimon/common/options/core_options.py b/paimon-python/pypaimon/common/options/core_options.py index a1fb6b2abfbf..193b80b6d7dd 100644 --- a/paimon-python/pypaimon/common/options/core_options.py +++ b/paimon-python/pypaimon/common/options/core_options.py @@ -1145,6 +1145,13 @@ class CoreOptions: ) ) + LOCAL_CACHE_EXCLUDE_EXTENSIONS: ConfigOption[str] = ( + ConfigOptions.key("local-cache.exclude-extensions") + .string_type() + .default_value("") + .with_description("Comma-separated file extensions to bypass in the local cache.") + ) + READ_BATCH_SIZE: ConfigOption[int] = ( ConfigOptions.key("read.batch-size") .int_type() @@ -1926,6 +1933,9 @@ def local_cache_block_size(self) -> MemorySize: def local_cache_whitelist(self) -> str: return self.options.get(CoreOptions.LOCAL_CACHE_WHITELIST) + def local_cache_exclude_extensions(self) -> str: + return self.options.get(CoreOptions.LOCAL_CACHE_EXCLUDE_EXTENSIONS) + def read_batch_size(self, default=None) -> int: return self.options.get(CoreOptions.READ_BATCH_SIZE, default or 1024) diff --git a/paimon-python/pypaimon/filesystem/caching_file_io.py b/paimon-python/pypaimon/filesystem/caching_file_io.py index 83b8a2cca8b2..efc055c1e104 100644 --- a/paimon-python/pypaimon/filesystem/caching_file_io.py +++ b/paimon-python/pypaimon/filesystem/caching_file_io.py @@ -26,6 +26,7 @@ """ import hashlib +import io import os import threading from collections import OrderedDict @@ -192,19 +193,33 @@ def put_file_size(self, file_path: str, size: int) -> None: self._file_size_cache[file_path] = size -class CachingInputStream: +class CachingInputStream(io.RawIOBase): """Wraps a remote stream with block-level caching.""" - def __init__(self, file_io, file_path: str, cache): + def __init__(self, file_io, file_path: str, cache, file_size=None): self._file_io = file_io self._stream = None self._file_path = file_path - self._file_size = -1 + self._file_size = -1 if file_size is None else file_size self._cache = cache + if file_size is not None: + cache.put_file_size(file_path, file_size) self._pos = 0 self._io_lock = threading.Lock() self._remote_supports_pread = None + def readable(self): + return True + + def seekable(self): + return True + + def readinto(self, buffer): + view = memoryview(buffer).cast('B') + data = self.read(len(view)) + view[:len(data)] = data + return len(data) + def _get_file_size(self) -> int: if self._file_size == -1: cached = self._cache.get_file_size(self._file_path) @@ -316,6 +331,7 @@ def close(self): if self._stream is not None: self._stream.close() self._stream = None + super().close() def __enter__(self): return self @@ -333,9 +349,13 @@ class CachingFileIO(FileIO): fall through to the delegate directly. """ - def __init__(self, delegate: FileIO, cache, whitelist=None): + def __init__(self, delegate: FileIO, cache, whitelist=None, excluded_extensions=None): self._delegate = delegate self._cache = cache + self._excluded_extensions = { + value.strip().lstrip('.').lower() + for value in (excluded_extensions or '').split(',') if value.strip().lstrip('.') + } if whitelist is None: self._whitelist = {FileType.META, FileType.GLOBAL_INDEX} else: @@ -390,19 +410,25 @@ def wrap_with_caching_if_needed(file_io, options, cache=None): whitelist = FileType.parse_whitelist(opts.local_cache_whitelist()) if not whitelist: return file_io - return CachingFileIO(file_io, cache, whitelist) + return CachingFileIO( + file_io, cache, whitelist, opts.local_cache_exclude_extensions()) @property def properties(self): return self._delegate.properties - def new_input_stream(self, path: str): + def _is_cacheable(self, path: str): file_type = FileType.classify(path) - eligible = (file_type in self._whitelist - or (file_type == FileType.PARQUET_DATA and FileType.DATA in self._whitelist)) - if self._cache is None or not eligible or FileType.is_mutable(path): + extension = os.path.basename(path).rsplit('.', 1)[-1].lower() + return (self._cache is not None + and file_type in self._whitelist + and extension not in self._excluded_extensions + and not FileType.is_mutable(path)) + + def new_input_stream(self, path: str, file_size=None): + if not self._is_cacheable(path): return self._delegate.new_input_stream(path) - return CachingInputStream(self._delegate, path, self._cache) + return CachingInputStream(self._delegate, path, self._cache, file_size) def new_output_stream(self, path: str): return self._delegate.new_output_stream(path) diff --git a/paimon-python/pypaimon/filesystem/caching_file_system.py b/paimon-python/pypaimon/filesystem/caching_file_system.py new file mode 100644 index 000000000000..9205c416c7cf --- /dev/null +++ b/paimon-python/pypaimon/filesystem/caching_file_system.py @@ -0,0 +1,78 @@ +# 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 pyarrow as pa +import pyarrow.fs as pafs + + +class CachedFileSystemHandler(pafs.FileSystemHandler): + """Route one fragment's input through FileIO while preserving filesystem operations.""" + + def __init__(self, file_io, path, file_size=None): + self._file_io = file_io + self._path = path + self._delegate = file_io.filesystem + self._filesystem_path = self._delegate.normalize_path(file_io.to_filesystem_path(path)) + self._file_size = file_size + + def get_type_name(self): + return 'paimon-cached-file' + + def open_input_file(self, path): + if self._delegate.normalize_path(path) != self._filesystem_path: + return self._delegate.open_input_file(path) + return pa.PythonFile( + self._file_io.new_input_stream(self._path, file_size=self._file_size), mode='r') + + def open_input_stream(self, path): + return self.open_input_file(path) + + def normalize_path(self, path): + return self._delegate.normalize_path(path) + + def get_file_info(self, paths): + return self._delegate.get_file_info(paths) + + def get_file_info_selector(self, selector): + return self._delegate.get_file_info(selector) + + def create_dir(self, path, recursive): + return self._delegate.create_dir(path, recursive=recursive) + + def delete_dir(self, path): + return self._delegate.delete_dir(path) + + def delete_dir_contents(self, path, missing_dir_ok): + return self._delegate.delete_dir_contents(path, missing_dir_ok=missing_dir_ok) + + def delete_root_dir_contents(self): + return self._delegate.delete_dir_contents("", accept_root_dir=True) + + def delete_file(self, path): + return self._delegate.delete_file(path) + + def move(self, src, dest): + return self._delegate.move(src, dest) + + def copy_file(self, src, dest): + return self._delegate.copy_file(src, dest) + + def open_output_stream(self, path, metadata): + return self._delegate.open_output_stream(path, metadata=metadata) + + def open_append_stream(self, path, metadata): + return self._delegate.open_append_stream(path, metadata=metadata) diff --git a/paimon-python/pypaimon/read/reader/format_pyarrow_reader.py b/paimon-python/pypaimon/read/reader/format_pyarrow_reader.py index ecb7db7e7416..83af71899bf0 100644 --- a/paimon-python/pypaimon/read/reader/format_pyarrow_reader.py +++ b/paimon-python/pypaimon/read/reader/format_pyarrow_reader.py @@ -272,6 +272,17 @@ def _file_format_metadata_cache_max_size(file_io: FileIO) -> int: CatalogOptions.FILE_FORMAT_METADATA_CACHE_MAX_SIZE).get_bytes() +def _uses_block_cache(file_io: FileIO, file_path: str) -> bool: + from pypaimon.filesystem.caching_file_io import CachingFileIO + return isinstance(file_io, CachingFileIO) and file_io._is_cacheable(file_path) + + +def _open_pyarrow_input_file(file_io: FileIO, file_path: str, file_size=None): + if _uses_block_cache(file_io, file_path): + return pa.PythonFile(file_io.new_input_stream(file_path, file_size=file_size), mode='r') + return file_io.filesystem.open_input_file(file_io.to_filesystem_path(file_path)) + + def _file_format_dataset(file_io: FileIO, file_format: str, file_path: str, cache_max_size: int, file_size: Optional[int] = None): @@ -288,6 +299,11 @@ def _file_format_dataset(file_io: FileIO, file_format: str, file_path: str, if known_size is not None and register_file_size is not None: register_file_size(file_path_for_pyarrow, known_size) + if _uses_block_cache(file_io, file_path): + from pypaimon.filesystem.caching_file_system import CachedFileSystemHandler + import pyarrow.fs as pafs + filesystem = pafs.PyFileSystem(CachedFileSystemHandler(file_io, file_path, known_size)) + def load(): if file_format == 'parquet': parquet_format = ds.ParquetFileFormat() @@ -312,7 +328,8 @@ def load(): file_path_for_pyarrow, format=file_format, filesystem=filesystem) key = ( - _FilesystemIdentity(filesystem), file_format, file_path_for_pyarrow, + _FilesystemIdentity(file_io if _uses_block_cache(file_io, file_path) else filesystem), + file_format, file_path_for_pyarrow, known_size) if cache_max_size <= 0: _reset_file_format_dataset_cache() @@ -332,8 +349,7 @@ def _orc_schema_with_field_metadata(file_io: FileIO, file_path: str, return fallback import pyarrow.orc as orc - source = file_io.filesystem.open_input_file( - file_io.to_filesystem_path(file_path)) + source = _open_pyarrow_input_file(file_io, file_path) try: metadata = orc.ORCFile(source).metadata arrow_schema = metadata.get(b"ARROW:schema") @@ -380,7 +396,7 @@ def __init__(self, file_io: FileIO, file_format: str, file_path: str, file_path_for_pyarrow = file_io.to_filesystem_path(file_path) self._row_group_cache = row_group_cache self._row_group_cache_filesystem = _FilesystemIdentity( - file_io.filesystem) + file_io if _uses_block_cache(file_io, file_path) else file_io.filesystem) self._row_group_cache_path = file_path_for_pyarrow cache_max_size = _file_format_metadata_cache_max_size(file_io) self.dataset = _file_format_dataset( @@ -560,8 +576,7 @@ def __init__(self, file_io: FileIO, file_format: str, file_path: str, and self._selected_shared_map_paths)): import pyarrow.parquet as pq # ParquetFile(filesystem=...) is unavailable in PyArrow 6. - self._parquet_source = file_io.filesystem.open_input_file( - file_path_for_pyarrow) + self._parquet_source = _open_pyarrow_input_file(file_io, file_path) try: self._parquet_file = pq.ParquetFile(self._parquet_source) if (self._selected_parquet_row_groups is not None @@ -580,8 +595,7 @@ def __init__(self, file_io: FileIO, file_format: str, file_path: str, raise if file_format == 'orc' and self._selected_shared_map_paths: import pyarrow.orc as orc - self._orc_source = file_io.filesystem.open_input_file( - file_path_for_pyarrow) + self._orc_source = _open_pyarrow_input_file(file_io, file_path) self._orc_file = orc.ORCFile(self._orc_source) if self._exhausted: self._raw_batches = iter(()) diff --git a/paimon-python/pypaimon/tests/caching_file_io_test.py b/paimon-python/pypaimon/tests/caching_file_io_test.py index 9eeff96b2e8b..42beef5d064f 100644 --- a/paimon-python/pypaimon/tests/caching_file_io_test.py +++ b/paimon-python/pypaimon/tests/caching_file_io_test.py @@ -336,22 +336,23 @@ def test_custom_whitelist_meta_only(self): result = caching_io.new_input_stream("global-index-uuid.index") self.assertNotIsInstance(result, CachingInputStream) - def test_parquet_only_cache_and_data_compatibility(self): + def test_data_cache_excludes_blob_extension(self): from pypaimon.filesystem.caching_file_io import LocalMemoryCacheManager from pypaimon.utils.file_type import FileType for disk in (False, True): - for whitelist in ('meta,global-index', 'parquet-data', 'data'): - for name in ('data.parquet', 'data.blob', 'data.orc', 'data.parquet.index'): + for whitelist in ('meta,global-index', 'data'): + for name in ('data.parquet', 'data.blob', 'data.orc', 'data.avro', + 'data.parquet.index', 'snapshot-1.blob'): with self.subTest(disk=disk, whitelist=whitelist, name=name): with tempfile.TemporaryDirectory() as directory: cache = (LocalDiskCacheManager(directory, 1024, block_size=4) if disk else LocalMemoryCacheManager(1024, block_size=4)) delegate = self._make_delegate({name: b'abcdefgh'}) caching_io = CachingFileIO( - delegate, cache, FileType.parse_whitelist(whitelist)) - cached = ((whitelist == 'parquet-data' and name.endswith('.parquet')) - or (whitelist == 'data' and not name.endswith('.index'))) + delegate, cache, FileType.parse_whitelist(whitelist), ' .BLOB, ') + cached = (whitelist == 'data' and name not in + ('data.blob', 'data.parquet.index', 'snapshot-1.blob')) for _ in range(2): with caching_io.new_input_stream(name) as stream: stream.seek(1) @@ -361,6 +362,16 @@ def test_parquet_only_cache_and_data_compatibility(self): if not cached: delegate.get_file_size.assert_not_called() + def test_data_whitelist_without_exclusion_still_caches_blob(self): + from pypaimon.utils.file_type import FileType + delegate = self._make_delegate({"data.blob": b"abcdefgh"}) + cache = LocalDiskCacheManager(self.cache_dir, 1024, block_size=4) + caching_io = CachingFileIO(delegate, cache, {FileType.DATA}) + for _ in range(2): + with caching_io.new_input_stream("data.blob") as stream: + self.assertEqual(b"abcdefgh", stream.read()) + delegate.new_input_stream.assert_called_once() + def test_data_file_not_cached(self): data = b"data content" delegate = self._make_delegate({"data-abc.orc": data}) diff --git a/paimon-python/pypaimon/tests/file_type_test.py b/paimon-python/pypaimon/tests/file_type_test.py index f17293184111..0efabd1833ee 100644 --- a/paimon-python/pypaimon/tests/file_type_test.py +++ b/paimon-python/pypaimon/tests/file_type_test.py @@ -84,7 +84,7 @@ def test_bucket_index(self): def test_data_files(self): self.assertEqual(FileType.DATA, FileType.classify("data-abc.orc")) - self.assertEqual(FileType.PARQUET_DATA, FileType.classify("data-abc.parquet")) + self.assertEqual(FileType.DATA, FileType.classify("data-abc.parquet")) self.assertEqual(FileType.DATA, FileType.classify("unknown-file")) def test_temp_file_unwrap(self): diff --git a/paimon-python/pypaimon/tests/parquet_local_cache_test.py b/paimon-python/pypaimon/tests/parquet_local_cache_test.py new file mode 100644 index 000000000000..b46b940cc4d6 --- /dev/null +++ b/paimon-python/pypaimon/tests/parquet_local_cache_test.py @@ -0,0 +1,105 @@ +# 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. + + +from unittest.mock import patch + +import pyarrow as pa +import pyarrow.orc as orc +import pyarrow.parquet as pq +import pytest + +from pypaimon.common.options import Options +from pypaimon.filesystem.caching_file_io import CachingFileIO +from pypaimon.filesystem.local_file_io import LocalFileIO +from pypaimon.read.reader.format_pyarrow_reader import ( + FormatPyArrowReader, _reset_file_format_dataset_cache, +) +from pypaimon.schema.data_types import AtomicType, DataField + + +def _read(file_io, path, file_format, rows, file_size=None): + reader = FormatPyArrowReader( + file_io, file_format, path, [DataField(0, 'v', AtomicType('INT'))], + None, row_indices=rows, batch_size=2, file_size=file_size) + values = [] + try: + while True: + batch = reader.read_arrow_batch() + if batch is None: + break + values.extend(batch.column(0).to_pylist()) + finally: + reader.close() + return values + + +@pytest.mark.parametrize('disk', [False, True]) +@pytest.mark.parametrize('metadata_cache', ['0 b', '1 mb']) +@pytest.mark.parametrize('file_format,rows', [('parquet', None), ('parquet', [1]), ('orc', None)]) +def test_format_reader_uses_block_cache(tmp_path, disk, metadata_cache, file_format, rows): + path = tmp_path / ('data.' + file_format) + table = pa.table({'v': pa.array([1, 2, 3], type=pa.int32())}) + if file_format == 'parquet': + pq.write_table(table, path, row_group_size=1) + else: + orc.write_table(table, path) + opts = Options({'local-cache.enabled': 'true', 'local-cache.block-size': '64 b', + 'local-cache.whitelist': 'data', + 'local-cache.exclude-extensions': 'blob', + 'file-format.metadata-cache.max-size': metadata_cache}) + if disk: + opts.data['local-cache.dir'] = str(tmp_path / 'cache') + delegate = LocalFileIO(catalog_options=opts) + cache = CachingFileIO.create_cache_manager(opts) + file_io = CachingFileIO.wrap_with_caching_if_needed(delegate, opts, cache) + expected = [2] if rows is not None else [1, 2, 3] + _reset_file_format_dataset_cache() + try: + opened = [] + original_open = delegate.new_input_stream + + def open_source(location): + stream = original_open(location) + opened.append(stream) + return stream + + with patch.object(delegate, 'new_input_stream', side_effect=open_source) as opens: + assert _read(file_io, path.as_uri(), file_format, rows) == expected + assert opens.call_count > 0 + assert all(stream.closed for stream in opened) + if file_format == 'parquet': + path.unlink() + # Drop the parsed-footer cache, proving that raw cached blocks are sufficient. + _reset_file_format_dataset_cache() + with patch.object(delegate, 'new_input_stream', side_effect=AssertionError('source read')): + with patch.object(delegate, 'get_file_size', side_effect=AssertionError('source stat')): + assert _read(file_io, path.as_uri(), file_format, rows) == expected + finally: + _reset_file_format_dataset_cache() + + +def test_known_parquet_size_avoids_source_stat(tmp_path): + path = tmp_path / 'data.parquet' + pq.write_table(pa.table({'v': pa.array([1, 2, 3], type=pa.int32())}), path) + opts = Options({'local-cache.enabled': 'true', 'local-cache.whitelist': 'data', + 'local-cache.exclude-extensions': 'blob', + 'file-format.metadata-cache.max-size': '0 b'}) + delegate = LocalFileIO(catalog_options=opts) + file_io = CachingFileIO.wrap_with_caching_if_needed( + delegate, opts, CachingFileIO.create_cache_manager(opts)) + with patch.object(delegate, 'get_file_size', side_effect=AssertionError('source stat')): + assert _read(file_io, str(path), 'parquet', None, path.stat().st_size) == [1, 2, 3] diff --git a/paimon-python/pypaimon/utils/file_type.py b/paimon-python/pypaimon/utils/file_type.py index 6f74152c4d1c..7dadc93af631 100644 --- a/paimon-python/pypaimon/utils/file_type.py +++ b/paimon-python/pypaimon/utils/file_type.py @@ -31,14 +31,12 @@ class FileType(Enum): - META: snapshot, schema, manifest, manifest sidecar, statistics, tag, changelog metadata, hint files, _SUCCESS, consumer, service files - DATA: data files and any unrecognized files (default) - - PARQUET_DATA: Parquet data files - BUCKET_INDEX: bucket level index files (Hash, DV) - GLOBAL_INDEX: table level global index files (btree, lumina, full-text) - FILE_INDEX: data-file index files (bloom filter, bitmap, etc.) """ META = "META" DATA = "DATA" - PARQUET_DATA = "PARQUET_DATA" BUCKET_INDEX = "BUCKET_INDEX" GLOBAL_INDEX = "GLOBAL_INDEX" FILE_INDEX = "FILE_INDEX" @@ -86,9 +84,6 @@ def classify(file_path: str) -> 'FileType': if parent == "changelog": return FileType.META - if name.endswith(".parquet"): - return FileType.PARQUET_DATA - return FileType.DATA @staticmethod @@ -98,7 +93,6 @@ def parse_whitelist(whitelist_str: str) -> set: "global-index": FileType.GLOBAL_INDEX, "bucket-index": FileType.BUCKET_INDEX, "data": FileType.DATA, - "parquet-data": FileType.PARQUET_DATA, "file-index": FileType.FILE_INDEX, } result = set() @@ -109,7 +103,7 @@ def parse_whitelist(whitelist_str: str) -> set: elif name: logger.warning( "Unknown local-cache.whitelist value '%s'. " - "Supported values: meta, global-index, bucket-index, data, parquet-data, file-index.", + "Supported values: meta, global-index, bucket-index, data, file-index.", name, ) return result From 7d8a7d00ece352ad6cf55a484d00e7a3c918c9f5 Mon Sep 17 00:00:00 2001 From: "xiaohongbo.xhb" Date: Mon, 28 Sep 2026 02:26:38 -0700 Subject: [PATCH 09/10] [python] Route Avro reads through data cache --- .../read/reader/format_avro_reader.py | 8 +++- .../tests/parquet_local_cache_test.py | 39 +++++++++++++++++++ 2 files changed, 45 insertions(+), 2 deletions(-) diff --git a/paimon-python/pypaimon/read/reader/format_avro_reader.py b/paimon-python/pypaimon/read/reader/format_avro_reader.py index 57309fab1d56..ec20f9d8e305 100644 --- a/paimon-python/pypaimon/read/reader/format_avro_reader.py +++ b/paimon-python/pypaimon/read/reader/format_avro_reader.py @@ -36,8 +36,12 @@ class FormatAvroReader(RecordBatchReader): def __init__(self, file_io: FileIO, file_path: str, read_fields: List[str], full_fields: List[DataField], push_down_predicate: Any, batch_size: int = 1024, nested_name_paths: Optional[List[List[str]]] = None): - file_path_for_io = file_io.to_filesystem_path(file_path) - self._file = file_io.filesystem.open_input_file(file_path_for_io) + from pypaimon.filesystem.caching_file_io import CachingFileIO + if isinstance(file_io, CachingFileIO) and file_io._is_cacheable(file_path): + self._file = file_io.new_input_stream(file_path) + else: + file_path_for_io = file_io.to_filesystem_path(file_path) + self._file = file_io.filesystem.open_input_file(file_path_for_io) self._avro_reader = fastavro.reader(self._file) self._batch_size = batch_size self._push_down_predicate = push_down_predicate diff --git a/paimon-python/pypaimon/tests/parquet_local_cache_test.py b/paimon-python/pypaimon/tests/parquet_local_cache_test.py index b46b940cc4d6..9b397781d585 100644 --- a/paimon-python/pypaimon/tests/parquet_local_cache_test.py +++ b/paimon-python/pypaimon/tests/parquet_local_cache_test.py @@ -17,6 +17,7 @@ from unittest.mock import patch +import fastavro import pyarrow as pa import pyarrow.orc as orc import pyarrow.parquet as pq @@ -25,6 +26,7 @@ from pypaimon.common.options import Options from pypaimon.filesystem.caching_file_io import CachingFileIO from pypaimon.filesystem.local_file_io import LocalFileIO +from pypaimon.read.reader.format_avro_reader import FormatAvroReader from pypaimon.read.reader.format_pyarrow_reader import ( FormatPyArrowReader, _reset_file_format_dataset_cache, ) @@ -103,3 +105,40 @@ def test_known_parquet_size_avoids_source_stat(tmp_path): delegate, opts, CachingFileIO.create_cache_manager(opts)) with patch.object(delegate, 'get_file_size', side_effect=AssertionError('source stat')): assert _read(file_io, str(path), 'parquet', None, path.stat().st_size) == [1, 2, 3] + + +@pytest.mark.parametrize('disk', [False, True]) +def test_avro_reader_uses_block_cache(tmp_path, disk): + path = tmp_path / 'data.avro' + schema = {'type': 'record', 'name': 'row', + 'fields': [{'name': 'v', 'type': 'int'}]} + with path.open('wb') as output: + fastavro.writer(output, schema, [{'v': 1}, {'v': 2}, {'v': 3}]) + + opts = Options({'local-cache.enabled': 'true', 'local-cache.block-size': '64 b', + 'local-cache.whitelist': 'data', + 'local-cache.exclude-extensions': 'blob'}) + if disk: + opts.data['local-cache.dir'] = str(tmp_path / 'cache') + delegate = LocalFileIO(catalog_options=opts) + file_io = CachingFileIO.wrap_with_caching_if_needed( + delegate, opts, CachingFileIO.create_cache_manager(opts)) + + def read(): + reader = FormatAvroReader(file_io, path.as_uri(), ['v'], + [DataField(0, 'v', AtomicType('INT'))], None, batch_size=2) + values = [] + try: + while True: + batch = reader.read_arrow_batch() + if batch is None: + return values + values.extend(batch.column(0).to_pylist()) + finally: + reader.close() + + assert read() == [1, 2, 3] + path.unlink() + with patch.object(delegate, 'new_input_stream', side_effect=AssertionError('source read')): + with patch.object(delegate, 'get_file_size', side_effect=AssertionError('source stat')): + assert read() == [1, 2, 3] From ab77e38a992e27c82f5d0c4f02f2f2a613f52400 Mon Sep 17 00:00:00 2001 From: "xiaohongbo.xhb" Date: Mon, 28 Sep 2026 02:29:15 -0700 Subject: [PATCH 10/10] [python] Match cache exclusions only to file extensions --- paimon-python/pypaimon/filesystem/caching_file_io.py | 3 ++- paimon-python/pypaimon/tests/caching_file_io_test.py | 8 ++++++++ 2 files changed, 10 insertions(+), 1 deletion(-) diff --git a/paimon-python/pypaimon/filesystem/caching_file_io.py b/paimon-python/pypaimon/filesystem/caching_file_io.py index efc055c1e104..8cce717aaa54 100644 --- a/paimon-python/pypaimon/filesystem/caching_file_io.py +++ b/paimon-python/pypaimon/filesystem/caching_file_io.py @@ -419,7 +419,8 @@ def properties(self): def _is_cacheable(self, path: str): file_type = FileType.classify(path) - extension = os.path.basename(path).rsplit('.', 1)[-1].lower() + stem, _, extension = os.path.basename(path).rpartition('.') + extension = extension.lower() if stem else '' return (self._cache is not None and file_type in self._whitelist and extension not in self._excluded_extensions diff --git a/paimon-python/pypaimon/tests/caching_file_io_test.py b/paimon-python/pypaimon/tests/caching_file_io_test.py index 42beef5d064f..43e0308099cc 100644 --- a/paimon-python/pypaimon/tests/caching_file_io_test.py +++ b/paimon-python/pypaimon/tests/caching_file_io_test.py @@ -336,6 +336,14 @@ def test_custom_whitelist_meta_only(self): result = caching_io.new_input_stream("global-index-uuid.index") self.assertNotIsInstance(result, CachingInputStream) + def test_extension_exclusion_does_not_match_extensionless_filename(self): + from pypaimon.utils.file_type import FileType + delegate = self._make_delegate({'snapshot-1': b'snap'}) + cache = LocalDiskCacheManager(self.cache_dir, 1024, block_size=4) + caching_io = CachingFileIO(delegate, cache, {FileType.META}, 'snapshot-1') + with caching_io.new_input_stream('snapshot-1') as stream: + self.assertIsInstance(stream, CachingInputStream) + def test_data_cache_excludes_blob_extension(self): from pypaimon.filesystem.caching_file_io import LocalMemoryCacheManager from pypaimon.utils.file_type import FileType