This is an automated email from the ASF dual-hosted git repository.
JingsongLi pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/paimon.git
The following commit(s) were added to refs/heads/master by this push:
new 935867d123 [python] Avoid redundant file metadata requests in Parquet
reads (#10152)
935867d123 is described below
commit 935867d1230144e0ed9479adbe2455ec0bc7f3de
Author: XiaoHongbo <[email protected]>
AuthorDate: Thu Sep 24 10:48:02 2026 +0800
[python] Avoid redundant file metadata requests in Parquet reads (#10152)
---
.../filesystem/jindo_file_system_handler.py | 43 +++++++++++++++++---
.../pypaimon/read/reader/format_pyarrow_reader.py | 35 +++++++++++-----
paimon-python/pypaimon/read/split_read.py | 3 +-
.../pypaimon/tests/jindo_file_system_test.py | 47 +++++++++++++++++++++-
.../pypaimon/tests/parquet_metadata_cache_test.py | 34 ++++++++++++++++
5 files changed, 146 insertions(+), 16 deletions(-)
diff --git a/paimon-python/pypaimon/filesystem/jindo_file_system_handler.py
b/paimon-python/pypaimon/filesystem/jindo_file_system_handler.py
index 73ad09077f..eaad4a89be 100644
--- a/paimon-python/pypaimon/filesystem/jindo_file_system_handler.py
+++ b/paimon-python/pypaimon/filesystem/jindo_file_system_handler.py
@@ -16,6 +16,9 @@
# under the License.
import logging
+import os
+import threading
+from collections import OrderedDict
import pyarrow as pa
from pyarrow import PythonFile
@@ -50,6 +53,7 @@ _CASE_SENSITIVE_JINDO_CONFIG_KEYS = {
OssOptions.OSS_ACCESS_KEY_SECRET.key().lower():
OssOptions.OSS_ACCESS_KEY_SECRET.key(),
OssOptions.OSS_SECURITY_TOKEN.key().lower():
OssOptions.OSS_SECURITY_TOKEN.key(),
}
+_KNOWN_FILE_SIZE_CACHE_MAX_ENTRIES = 4096
def _jindo_config_value(value) -> str:
@@ -118,8 +122,9 @@ def create_jindo_oss_filesystem(root_uri: str,
catalog_options: Options):
class JindoInputFile:
- def __init__(self, jindo_stream):
+ def __init__(self, jindo_stream, file_size=None):
self._stream = jindo_stream
+ self._file_size = file_size
self._closed = False
@property
@@ -132,13 +137,18 @@ class JindoInputFile:
if self.closed:
raise ValueError("I/O operation on closed file")
if nbytes is None or nbytes < 0:
- return self._stream.read()
+ if self._file_size is None:
+ return self._stream.read()
+ nbytes = max(0, self._file_size - self.tell())
return self._stream.read(nbytes)
def seek(self, position: int, whence: int = 0):
if self.closed:
raise ValueError("I/O operation on closed file")
- self._stream.seek(position, whence)
+ if whence == os.SEEK_END and self._file_size is not None:
+ position = self._file_size + position
+ whence = os.SEEK_SET
+ return self._stream.seek(position, whence)
def tell(self) -> int:
if self.closed:
@@ -210,10 +220,29 @@ class JindoFileSystemHandler(FileSystemHandler):
self.logger = logging.getLogger(__name__)
self.root_path = root_path
self.properties = catalog_options
+ self._known_file_sizes = OrderedDict()
+ self._known_file_sizes_lock = threading.Lock()
config = build_jindo_config(catalog_options)
self._jindo_fs = jfs.connect(self.root_path, "root", config)
+ def register_file_size(self, path: str, file_size: int):
+ """Register an immutable file's size supplied by Paimon metadata."""
+ normalized = self._normalize_path(path)
+ with self._known_file_sizes_lock:
+ self._known_file_sizes[normalized] = file_size
+ self._known_file_sizes.move_to_end(normalized)
+ while (len(self._known_file_sizes)
+ > _KNOWN_FILE_SIZE_CACHE_MAX_ENTRIES):
+ self._known_file_sizes.popitem(last=False)
+
+ def _known_file_size(self, path: str):
+ with self._known_file_sizes_lock:
+ file_size = self._known_file_sizes.get(path)
+ if file_size is not None:
+ self._known_file_sizes.move_to_end(path)
+ return file_size
+
def __eq__(self, other):
if isinstance(other, JindoFileSystemHandler):
return self.root_path == other.root_path
@@ -325,12 +354,16 @@ class JindoFileSystemHandler(FileSystemHandler):
def open_input_stream(self, path: str):
normalized = self._normalize_path(path)
jindo_stream = self._jindo_fs.open(normalized, "rb")
- return PythonFile(JindoInputFile(jindo_stream), mode="r")
+ return PythonFile(
+ JindoInputFile(jindo_stream, self._known_file_size(normalized)),
+ mode="r")
def open_input_file(self, path: str):
normalized = self._normalize_path(path)
jindo_stream = self._jindo_fs.open(normalized, "rb")
- return PythonFile(JindoInputFile(jindo_stream), mode="r")
+ return PythonFile(
+ JindoInputFile(jindo_stream, self._known_file_size(normalized)),
+ mode="r")
def open_output_stream(self, path: str, metadata):
normalized = self._normalize_path(path)
diff --git a/paimon-python/pypaimon/read/reader/format_pyarrow_reader.py
b/paimon-python/pypaimon/read/reader/format_pyarrow_reader.py
index ed6dd37215..572f90a54b 100644
--- a/paimon-python/pypaimon/read/reader/format_pyarrow_reader.py
+++ b/paimon-python/pypaimon/read/reader/format_pyarrow_reader.py
@@ -32,6 +32,7 @@ from pyarrow import RecordBatch
from pypaimon.common.file_io import FileIO
from pypaimon.common.options.config import CatalogOptions
from pypaimon.common.options.core_options import CoreOptions
+from pypaimon.filesystem.pyarrow_file_io import _pyarrow_lt_7
from pypaimon.data.map_shared_shredding import (
assemble_normal_map_selected_keys,
assemble_shared_shredding_map,
@@ -241,7 +242,7 @@ def _estimate_file_format_dataset_size(dataset,
file_format: str) -> Optional[in
def _estimate_file_format_cache_entry_size(
- key: Tuple[Any, str, str],
+ key: Tuple[Any, ...],
dataset,
file_format: str) -> Optional[int]:
metadata_size = _estimate_file_format_dataset_size(dataset, file_format)
@@ -255,9 +256,7 @@ def _estimate_file_format_cache_entry_size(
visible_size = (
metadata_size
+ sys.getsizeof(key)
- + sys.getsizeof(key[0])
- + sys.getsizeof(key[1])
- + sys.getsizeof(key[2])
+ + sum(sys.getsizeof(part) for part in key)
+ sys.getsizeof(dataset)
+ sys.getsizeof((dataset, metadata_size))
+ _FILE_FORMAT_METADATA_CACHE_CONTAINER_OVERHEAD
@@ -274,22 +273,39 @@ def _file_format_metadata_cache_max_size(file_io: FileIO)
-> int:
def _file_format_dataset(file_io: FileIO, file_format: str, file_path: str,
- cache_max_size: int):
+ cache_max_size: int,
+ file_size: Optional[int] = None):
file_path_for_pyarrow = file_io.to_filesystem_path(file_path)
filesystem = file_io.filesystem
+ known_size = file_size if file_size is not None and file_size > 0 else None
+
+ # PyArrow's Python FileSystemHandler API opens files by path, so its
+ # FileInfo overload cannot forward file_size to handlers such as Jindo.
+ # Let capable handlers consume the size from Paimon's immutable file
+ # metadata before PyArrow opens the file.
+ handler = getattr(filesystem, "handler", None)
+ register_file_size = getattr(handler, "register_file_size", None)
+ if known_size is not None and register_file_size is not None:
+ register_file_size(file_path_for_pyarrow, known_size)
def load():
if file_format == 'parquet':
parquet_format = ds.ParquetFileFormat()
+ fragment_options = {}
+ if known_size is not None and not _pyarrow_lt_7():
+ fragment_options["file_size"] = known_size
fragment = parquet_format.make_fragment(
- file_path_for_pyarrow, filesystem=filesystem)
+ file_path_for_pyarrow, filesystem=filesystem,
+ **fragment_options)
# Reuse this fragment's footer for schema discovery and scanning.
return ds.FileSystemDataset(
[fragment], fragment.physical_schema, parquet_format,
filesystem)
return ds.dataset(
file_path_for_pyarrow, format=file_format, filesystem=filesystem)
- key = (_FilesystemIdentity(filesystem), file_format, file_path_for_pyarrow)
+ key = (
+ _FilesystemIdentity(filesystem), file_format, file_path_for_pyarrow,
+ known_size)
if cache_max_size <= 0:
_reset_file_format_dataset_cache()
return load()
@@ -347,7 +363,8 @@ class FormatPyArrowReader(RecordBatchReader):
predicate_field_names: Optional[Set[str]] = None,
row_indices: Optional[List[int]] = None,
row_ranges: Optional[List[Tuple[int, int]]] = None,
- row_group_cache: Optional[_DecodedRowGroupCache] = None):
+ row_group_cache: Optional[_DecodedRowGroupCache] = None,
+ file_size: Optional[int] = None):
from pypaimon.filesystem.resolving_file_io import ResolvingFileIO
if isinstance(file_io, ResolvingFileIO):
file_io = file_io._get_fileio(file_path)
@@ -359,7 +376,7 @@ class FormatPyArrowReader(RecordBatchReader):
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(
- file_io, file_format, file_path, cache_max_size)
+ file_io, file_format, file_path, cache_max_size, file_size)
self._range_slicer = None
self._selected_parquet_row_groups = None
self._exhausted = False
diff --git a/paimon-python/pypaimon/read/split_read.py
b/paimon-python/pypaimon/read/split_read.py
index 1932852563..2a044c4045 100644
--- a/paimon-python/pypaimon/read/split_read.py
+++ b/paimon-python/pypaimon/read/split_read.py
@@ -456,7 +456,8 @@ class SplitRead(ABC):
nested_name_paths=ordered_nested_paths,
predicate_field_names=predicate_fields,
row_ranges=parquet_row_ranges,
- row_group_cache=self._parquet_row_group_cache)
+ row_group_cache=self._parquet_row_group_cache,
+ file_size=file.file_size)
elif file_format == CoreOptions.FILE_FORMAT_ROW:
if has_nested:
raise NotImplementedError(
diff --git a/paimon-python/pypaimon/tests/jindo_file_system_test.py
b/paimon-python/pypaimon/tests/jindo_file_system_test.py
index b828fd8b83..8c62113ca9 100644
--- a/paimon-python/pypaimon/tests/jindo_file_system_test.py
+++ b/paimon-python/pypaimon/tests/jindo_file_system_test.py
@@ -27,7 +27,11 @@ from pyarrow.fs import PyFileSystem
from pypaimon.common.options import Options
from pypaimon.common.options.config import OssOptions
from pypaimon.filesystem import jindo_file_system_handler as jindo_module
-from pypaimon.filesystem.jindo_file_system_handler import
JindoFileSystemHandler, JINDO_AVAILABLE
+from pypaimon.filesystem.jindo_file_system_handler import (
+ JindoFileSystemHandler,
+ JindoInputFile,
+ JINDO_AVAILABLE,
+)
class _RecordingConfig:
@@ -38,6 +42,47 @@ class _RecordingConfig:
self.values[key] = value
+class _RecordingInputStream:
+ def __init__(self, position=0):
+ self.position = position
+ self.seeks = []
+ self.reads = []
+ self.closed = False
+
+ def seek(self, position, whence=0):
+ self.seeks.append((position, whence))
+ self.position = position
+ return position
+
+ def tell(self):
+ return self.position
+
+ def read(self, size=None):
+ self.reads.append(size)
+ self.position += size or 0
+ return b"x" * (size or 0)
+
+ def close(self):
+ self.closed = True
+
+
+class JindoInputFileTest(unittest.TestCase):
+
+ def test_known_size_avoids_backend_seek_from_end(self):
+ stream = _RecordingInputStream()
+ input_file = JindoInputFile(stream, file_size=100)
+
+ self.assertEqual(95, input_file.seek(-5, os.SEEK_END))
+ self.assertEqual([(95, os.SEEK_SET)], stream.seeks)
+
+ def test_known_size_bounds_unlimited_read(self):
+ stream = _RecordingInputStream(position=40)
+ input_file = JindoInputFile(stream, file_size=100)
+
+ self.assertEqual(b"x" * 60, input_file.read())
+ self.assertEqual([60], stream.reads)
+
+
class JindoConfigTest(unittest.TestCase):
def test_forwards_native_options_to_connect(self):
diff --git a/paimon-python/pypaimon/tests/parquet_metadata_cache_test.py
b/paimon-python/pypaimon/tests/parquet_metadata_cache_test.py
index 12cead96fa..711e73191e 100644
--- a/paimon-python/pypaimon/tests/parquet_metadata_cache_test.py
+++ b/paimon-python/pypaimon/tests/parquet_metadata_cache_test.py
@@ -95,6 +95,9 @@ class _CountingFileSystemHandler(pafs.FSSpecHandler):
self.calls.append("open_input_file")
return super().open_input_file(path)
+ def register_file_size(self, path, file_size):
+ self.calls.append(("register_file_size", path, file_size))
+
class FileFormatMetadataCacheTest(unittest.TestCase):
def setUp(self):
@@ -232,6 +235,37 @@ class FileFormatMetadataCacheTest(unittest.TestCase):
self.assertEqual({"value": list(range(10))}, results[0])
self.assertEqual(results[0], results[1])
+ def test_forwards_known_file_size(self):
+ parquet_format = unittest.mock.Mock()
+ fragment = unittest.mock.Mock(physical_schema=pa.schema([]))
+ parquet_format.make_fragment.return_value = fragment
+ with patch.object(
+ reader_module.ds, "ParquetFileFormat",
+ return_value=parquet_format), patch.object(
+ reader_module.ds, "FileSystemDataset",
+ return_value=unittest.mock.sentinel.dataset):
+ dataset = reader_module._file_format_dataset(
+ self.file_io, "parquet", self.paths[0], 0, 123)
+ self.assertIs(unittest.mock.sentinel.dataset, dataset)
+ expected_options = (
+ {} if reader_module._pyarrow_lt_7()
+ else {"file_size": 123})
+ parquet_format.make_fragment.assert_called_once_with(
+ self.paths[0], filesystem=self.file_io.filesystem,
+ **expected_options)
+
+ handler = _CountingFileSystemHandler()
+ self.file_io.filesystem = pafs.PyFileSystem(handler)
+ file_size = os.path.getsize(self.paths[0])
+
+ dataset = reader_module._file_format_dataset(
+ self.file_io, "parquet", self.paths[0], 0, file_size)
+
+ self.assertEqual({"value": list(range(10))},
+ dataset.to_table().to_pydict())
+ self.assertIn(
+ ("register_file_size", self.paths[0], file_size), handler.calls)
+
def test_fragment_metadata_is_reused_without_io(self):
handler = _CountingFileSystemHandler()
self.file_io.filesystem = pafs.PyFileSystem(handler)