diff --git a/docs/docs/program-api/file-cache.mdx b/docs/docs/program-api/file-cache.mdx index 8c58a398bcab..e6f194230d0a 100644 --- a/docs/docs/program-api/file-cache.mdx +++ b/docs/docs/program-api/file-cache.mdx @@ -127,11 +127,15 @@ The Java snippet is a method-body fragment; declare or handle the catalog except | `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.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. +Use `local-cache.whitelist=meta,global-index,data` with `local-cache.exclude-extensions=blob` +to cache data files except BLOB bodies. + ## How It Works 1. The reader requests bytes from an eligible file. The cache maps the byte range to fixed-size 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 3fad543d8d8c..8cce717aaa54 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,17 +410,26 @@ 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) - if self._cache is None or file_type not in self._whitelist or FileType.is_mutable(path): + 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 + 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_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/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 940a4f4ab29a..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,50 @@ 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 + + for disk in (False, True): + 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), ' .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) + 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_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/parquet_local_cache_test.py b/paimon-python/pypaimon/tests/parquet_local_cache_test.py new file mode 100644 index 000000000000..9b397781d585 --- /dev/null +++ b/paimon-python/pypaimon/tests/parquet_local_cache_test.py @@ -0,0 +1,144 @@ +# 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 fastavro +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_avro_reader import FormatAvroReader +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] + + +@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]