Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 4 additions & 0 deletions docs/docs/program-api/file-cache.mdx
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
10 changes: 10 additions & 0 deletions paimon-python/pypaimon/common/options/core_options.py
Original file line number Diff line number Diff line change
Expand Up @@ -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()
Expand Down Expand Up @@ -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)

Expand Down
45 changes: 37 additions & 8 deletions paimon-python/pypaimon/filesystem/caching_file_io.py
Original file line number Diff line number Diff line change
Expand Up @@ -26,6 +26,7 @@
"""

import hashlib
import io
import os
import threading
from collections import OrderedDict
Expand Down Expand Up @@ -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)
Expand Down Expand Up @@ -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
Expand All @@ -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:
Expand Down Expand Up @@ -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)
Expand Down
78 changes: 78 additions & 0 deletions paimon-python/pypaimon/filesystem/caching_file_system.py
Original file line number Diff line number Diff line change
@@ -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)
8 changes: 6 additions & 2 deletions paimon-python/pypaimon/read/reader/format_avro_reader.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
30 changes: 22 additions & 8 deletions paimon-python/pypaimon/read/reader/format_pyarrow_reader.py
Original file line number Diff line number Diff line change
Expand Up @@ -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):
Expand All @@ -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()
Expand All @@ -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()
Expand All @@ -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")
Expand Down Expand Up @@ -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(
Expand Down Expand Up @@ -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
Expand All @@ -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(())
Expand Down
44 changes: 44 additions & 0 deletions paimon-python/pypaimon/tests/caching_file_io_test.py
Original file line number Diff line number Diff line change
Expand Up @@ -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})
Expand Down
Loading
Loading