blob: 2f6f09665f9b8dca7c6a837f2d946b8380153a25 [file]
# 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 os
from abc import ABC, abstractmethod
from functools import partial
from typing import Callable, Dict, List, Optional, Tuple
from pypaimon.common.merge_engine_dispatch import build_merge_function
from pypaimon.common.options.core_options import CoreOptions, MergeEngine
from pypaimon.common.predicate import Predicate
from pypaimon.deletionvectors import (
ApplyDeletionVectorReader,
PositionMappedDeletionVector,
)
from pypaimon.deletionvectors.deletion_vector import DeletionVector
from pypaimon.globalindex import Range
from pypaimon.manifest.schema.data_file_meta import DataFileMeta
from pypaimon.read.interval_partition import IntervalPartition, SortedRun
from pypaimon.read.partition_info import PartitionInfo
from pypaimon.read.push_down_utils import (
predicate_field_names,
rewrite_predicate_indices,
trim_predicate_by_fields,
)
from pypaimon.read.reader.concat_batch_reader import (
BlobFallbackBatchReader, ConcatBatchReader,
MergeAllBatchReader, DataEvolutionMergeReader)
from pypaimon.read.reader.concat_record_reader import ConcatRecordReader
from pypaimon.read.reader.auth_masking_reader import AuthFilterReader
from pypaimon.read.reader.data_file_batch_reader import DataFileBatchReader
from pypaimon.read.reader.deferred_blob_resolve_reader import \
DeferredBlobResolveReader
from pypaimon.read.reader.drop_delete_reader import DropDeleteRecordReader
from pypaimon.read.reader.empty_record_reader import EmptyFileRecordReader
from pypaimon.read.reader.field_bunch import BlobBunch, DataBunch, FieldBunch, VectorBunch
from pypaimon.read.reader.field_indices import (
blob_field_indices, vector_field_indices)
from pypaimon.read.reader.filter_record_reader import FilterRecordReader
from pypaimon.read.reader.format_avro_reader import FormatAvroReader
from pypaimon.read.reader.blob_descriptor_convert_reader import BlobInlineConvertReader
from pypaimon.read.reader.filter_record_batch_reader import FilterRecordBatchReader
from pypaimon.read.reader.limited_record_reader import LimitedRecordBatchReader, LimitedRecordReader
from pypaimon.read.reader.row_range_filter_record_reader import RowIdFilterRecordBatchReader
from pypaimon.read.reader.format_blob_reader import FormatBlobReader
from pypaimon.read.reader.format_lance_reader import FormatLanceReader
from pypaimon.read.reader.format_pyarrow_reader import FormatPyArrowReader
from pypaimon.read.reader.format_row_reader import FormatRowReader
from pypaimon.read.reader.format_mosaic_reader import FormatMosaicReader
from pypaimon.read.reader.format_vortex_reader import FormatVortexReader
from pypaimon.read.reader.iface.record_batch_reader import (RecordBatchReader,
RowPositionReader, EmptyRecordBatchReader)
from pypaimon.read.reader.iface.record_reader import RecordReader
from pypaimon.read.reader.key_value_unwrap_reader import \
KeyValueUnwrapRecordReader
from pypaimon.read.reader.key_value_wrap_reader import KeyValueWrapReader
from pypaimon.read.reader.shard_batch_reader import ShardBatchReader
from pypaimon.read.reader.aggregation_merge_function import (
AggregateMergeFunction, build_field_aggregators)
from pypaimon.read.reader.sort_merge_reader import (SortMergeReaderWithMinHeap,
builtin_seq_comparator)
from pypaimon.read.split import Split
from pypaimon.read.sliced_split import SlicedSplit
from pypaimon.schema.data_types import DataField, PyarrowFieldParser
from pypaimon.table.special_fields import SpecialFields
from pypaimon.globalindex.indexed_split import IndexedSplit
from pypaimon.utils.data_evolution_utils import retrieve_anchor_file
KEY_PREFIX = "_KEY_"
KEY_FIELD_ID_START = 1000000
NULL_FIELD_INDEX = -1
def deferred_blob_field_names(table, read_fields: List[DataField],
predicate: Optional[Predicate],
limit: Optional[int],
has_post_filter: bool = False) -> set:
# An auth filter also selects rows; defer past it too, like a predicate/limit.
if ((predicate is None and limit is None and not has_post_filter)
or CoreOptions.blob_as_descriptor(table.options)):
return set()
inline_fields = (
CoreOptions.blob_descriptor_fields(table.options)
| CoreOptions.blob_view_fields(table.options)
)
predicate_fields = (
predicate_field_names(predicate) if predicate is not None else set()
)
return {
read_fields[index].name
for index in blob_field_indices(read_fields)
if read_fields[index].name not in inline_fields
and read_fields[index].name not in predicate_fields
}
ROW_SIDECAR_FORMAT = CoreOptions.FILE_FORMAT_ROW
_COMPRESS_EXTENSIONS = frozenset(['gz', 'bz2', 'deflate', 'snappy', 'lz4', 'zst'])
def format_identifier(file_name):
idx = file_name.rfind('.')
assert idx != -1, "%s is not a legal file name." % file_name
ext = file_name[idx + 1:]
if ext.lower() in _COMPRESS_EXTENSIONS:
second_idx = file_name.rfind('.', 0, idx)
assert second_idx != -1, "%s is not a legal file name." % file_name
return file_name[second_idx + 1:idx]
return ext
class SplitRead(ABC):
"""Abstract base class for split reading operations."""
def __init__(
self,
table,
predicate: Optional[Predicate],
read_type: List[DataField],
split: Split,
row_tracking_enabled: bool,
nested_name_paths: Optional[List[List[str]]] = None,
limit: Optional[int] = None):
from pypaimon.table.file_store_table import FileStoreTable
self.table: FileStoreTable = table
self.predicate = predicate
self.push_down_predicate = self._push_down_predicate()
self.split = split
self.row_tracking_enabled = row_tracking_enabled
self.value_arity = len(read_type)
self.nested_name_paths = nested_name_paths
self.limit = limit
self._blob_parallelism = 1
# Snapshot the raw value-side schema before _create_key_value_fields
# wraps it, so MergeFileSplitRead can hand per-value-field nullable
# flags to merge functions that enforce NOT-NULL on every add().
self.value_fields = list(read_type)
self.trimmed_primary_key = self.table.trimmed_primary_keys
self.read_fields = read_type
if isinstance(self, MergeFileSplitRead):
self.read_fields = self._create_key_value_fields(read_type)
self._cached_nested_path_by_name = self._compute_nested_path_by_name()
self.schema_id_2_fields = {}
self.deletion_file_readers = {}
# Apply filter only when all predicate columns are read by this scan,
# AND remap predicate leaf indices into the row layout the reader sees.
# Predicate leaves carry an `index` baked in by PredicateBuilder against
# the *original* table schema; if `read_type` is narrower or reordered,
# that index no longer matches the OffsetRow handed to
# FilterRecordReader (which would otherwise raise IndexError).
# We use `read_type` here, not `self.read_fields`: MergeFileSplitRead
# augments `read_fields` with _KEY_*/_SEQ/_KIND prefixes, but
# KeyValueUnwrapRecordReader returns kv.value whose arity equals
# len(read_type) and whose coordinate space is read_type — that is
# the space FilterRecordReader actually evaluates against.
read_type_names = {f.name for f in read_type}
if (
self.predicate is not None
and predicate_field_names(self.predicate).issubset(read_type_names)
):
self.predicate_for_reader = rewrite_predicate_indices(
self.predicate, read_type
)
else:
self.predicate_for_reader = None
def _compute_nested_path_by_name(self) -> Optional[Dict[str, List[str]]]:
if not self.nested_name_paths:
return None
if not any(len(p) > 1 for p in self.nested_name_paths):
return None
out: Dict[str, List[str]] = {}
for f, path in zip(self.read_fields[:self.value_arity],
self.nested_name_paths):
out[f.name] = path
return out
def _nested_path_by_name(self) -> Optional[Dict[str, List[str]]]:
return self._cached_nested_path_by_name
def _resolve_schema(self, schema_id: int):
"""Resolve schema, short-circuiting current table schema id to avoid
filesystem access (REST catalog would get 403).
"""
if schema_id == self.table.table_schema.id:
return self.table.table_schema
return self.table.schema_manager.get_schema(schema_id)
def _push_down_predicate(self) -> Optional[Predicate]:
if self.predicate is None:
return None
elif self.table.is_primary_key_table:
pk_predicate = trim_predicate_by_fields(self.predicate, self.table.primary_keys)
if not pk_predicate:
return None
return pk_predicate
else:
return self.predicate
@abstractmethod
def create_reader(self) -> RecordReader:
"""Create a record reader for the given split."""
# row_ranges: from IndexedSplit (ANN vector search), a list of discrete global row ID ranges.
# shard_range: from SlicedSplit (parallel shard scan), a contiguous [start, end) row range within the file.
def file_reader_supplier(self, file: DataFileMeta, for_merge_read: bool,
read_fields: List[str], row_tracking_enabled: bool,
row_ranges: Optional[List[Range]] = None,
shard_range: Optional[Tuple[int, int]] = None) -> RecordBatchReader:
(
read_file_fields,
read_arrow_predicate,
read_paimon_predicate,
) = self._get_fields_and_predicate(file.schema_id, read_fields)
# Use external_path if available, otherwise use file_path
file_path = file.external_path if file.external_path else file.file_path
file_format = format_identifier(os.path.basename(file_path))
batch_size = self.table.options.read_batch_size()
effective_row_ranges = None
if row_ranges is not None:
effective_row_ranges = Range.and_(row_ranges, [file.row_id_range()])
if len(effective_row_ranges) == 0:
return EmptyRecordBatchReader()
row_sidecar_file = self._row_sidecar_file_name(file)
if row_sidecar_file is not None and self._should_read_row_sidecar(
file,
effective_row_ranges,
row_sidecar_file,
self.table.options.data_evolution_row_sidecar_max_selected_rows(),
self.table.options.data_evolution_row_sidecar_max_selection_ratio()):
file_path = self._aligned_extra_file_path(file, row_sidecar_file)
file_format = ROW_SIDECAR_FORMAT
# Prepare file-local native row selection. Existing native formats
# consume row indices; Parquet keeps compact ranges to avoid expanding
# large selections into millions of Python integers.
row_indices = None
parquet_row_ranges = None
if effective_row_ranges is not None:
row_index_formats = (CoreOptions.FILE_FORMAT_BLOB,
CoreOptions.FILE_FORMAT_VORTEX,
CoreOptions.FILE_FORMAT_LANCE,
CoreOptions.FILE_FORMAT_ROW)
if file_format in row_index_formats:
row_indices = [
row_id - file.first_row_id
for row_range in effective_row_ranges
for row_id in range(row_range.from_, row_range.to + 1)
]
elif (file_format == CoreOptions.FILE_FORMAT_PARQUET
and read_arrow_predicate is None):
parquet_row_ranges = []
merged_ranges = Range.sort_and_merge_overlap(
effective_row_ranges, True)
for r in merged_ranges:
start = max(0, r.from_ - file.first_row_id)
end = min(
file.row_count - 1,
r.to - file.first_row_id,
)
if end >= start:
parquet_row_ranges.append((start, end))
if not parquet_row_ranges:
return EmptyRecordBatchReader()
# Map nested paths into the order the format reader will see.
nested_path_by_name = self._nested_path_by_name()
has_nested = nested_path_by_name is not None
# Field-id based per-file read (non-nested): select the file's OWN
# physical fields by field id and read them under the file's original
# names/types. A normalize step in DataFileBatchReader then aligns the
# batch to the latest read schema by field id (not by name), so a
# rename follows the id and a dropped-then-readded name cannot revive
# stale data. Nested-projection reads stay on the legacy name path.
file_read_fields = None if has_nested else self._file_read_fields(file)
target_fields = None if has_nested else self._target_read_fields()
if file_read_fields is not None:
read_file_fields = [f.name for f in file_read_fields]
name_to_field: Dict[str, DataField] = {
f.name: f for f in file_read_fields}
else:
# Cover both the merge-internal aliases (``_KEY_id``) and the
# bare user-facing PK name (``id``) the file actually stores.
name_to_field = {f.name: f for f in self.read_fields}
_, _trimmed_lookup_fields = self._get_trimmed_fields(
self._get_read_data_fields(), self._get_all_data_fields()
)
for f in _trimmed_lookup_fields:
name_to_field.setdefault(f.name, f)
format_reader: RecordBatchReader
if file_format == CoreOptions.FILE_FORMAT_AVRO:
avro_nested_paths = (
[nested_path_by_name[name] for name in read_file_fields]
if has_nested else None
)
# Pass the alias-safe union so FormatAvroReader can resolve
# the bare PK name (e.g. ``id``) requested by read_file_fields,
# even when value projection drops it from self.read_fields.
format_reader = FormatAvroReader(
self.table.file_io, file_path, read_file_fields,
list(name_to_field.values()),
read_arrow_predicate, batch_size=batch_size,
nested_name_paths=avro_nested_paths)
elif file_format == CoreOptions.FILE_FORMAT_BLOB:
if has_nested:
raise NotImplementedError(
"Nested-field projection is not supported on BLOB files")
blob_as_descriptor = self._read_blob_as_descriptor(read_file_fields)
blob_parallelism = self._blob_parallelism
format_reader = FormatBlobReader(self.table.file_io, file_path, read_file_fields,
self.read_fields, read_arrow_predicate, blob_as_descriptor,
batch_size=batch_size,
row_indices=row_indices,
blob_parallelism=blob_parallelism)
elif file_format == CoreOptions.FILE_FORMAT_LANCE:
if has_nested:
raise NotImplementedError(
"Nested-field projection is not supported on Lance files")
ordered_read_fields = [name_to_field[n] for n in read_file_fields if n in name_to_field]
format_reader = FormatLanceReader(self.table.file_io, file_path, ordered_read_fields,
read_arrow_predicate, batch_size=batch_size,
row_indices=row_indices,
shard_range=shard_range)
elif file_format == CoreOptions.FILE_FORMAT_VORTEX:
if has_nested:
raise NotImplementedError(
"Nested-field projection is not supported on Vortex files")
ordered_read_fields = [name_to_field[n] for n in read_file_fields if n in name_to_field]
predicate_fields = (
predicate_field_names(self.push_down_predicate)
if self.push_down_predicate else set())
format_reader = FormatVortexReader(self.table.file_io, file_path, ordered_read_fields,
read_arrow_predicate, batch_size=batch_size,
row_indices=row_indices,
shard_range=shard_range,
predicate_fields=predicate_fields)
elif file_format == CoreOptions.FILE_FORMAT_MOSAIC:
if has_nested:
raise NotImplementedError(
"Nested-field projection is not supported on Mosaic files")
ordered_read_fields = [name_to_field[n] for n in read_file_fields if n in name_to_field]
row_group_predicate = (
read_paimon_predicate
if file.schema_id == self.table.table_schema.id else None
)
format_reader = FormatMosaicReader(self.table.file_io, file_path, ordered_read_fields,
read_arrow_predicate, batch_size=batch_size,
row_group_predicate=row_group_predicate)
elif file_format == CoreOptions.FILE_FORMAT_PARQUET or file_format == CoreOptions.FILE_FORMAT_ORC:
ordered_read_fields = [name_to_field[n] for n in read_file_fields if n in name_to_field]
ordered_nested_paths = (
[nested_path_by_name[f.name] for f in ordered_read_fields]
if has_nested else None
)
predicate_fields = (
predicate_field_names(self.push_down_predicate)
if self.push_down_predicate else set())
format_reader = FormatPyArrowReader(
self.table.file_io, file_format, file_path,
ordered_read_fields, read_arrow_predicate, batch_size=batch_size,
options=self.table.options,
nested_name_paths=ordered_nested_paths,
predicate_field_names=predicate_fields,
row_ranges=parquet_row_ranges)
elif file_format == CoreOptions.FILE_FORMAT_ROW:
if has_nested:
raise NotImplementedError(
"Nested-field projection is not supported on ROW files")
file_schema = self._resolve_schema(file.schema_id)
if file.write_cols:
field_map = {f.name: f for f in file_schema.fields}
row_full_fields = [field_map[n] for n in file.write_cols
if n in field_map]
elif self.table.is_primary_key_table:
row_full_fields = self._create_key_value_fields(
file_schema.fields)
else:
row_full_fields = file_schema.fields
format_reader = FormatRowReader(
self.table.file_io, file_path, read_file_fields,
row_full_fields,
read_arrow_predicate, batch_size=batch_size,
row_indices=row_indices)
elif file_format in ('json', 'csv'):
raise NotImplementedError(
f"Reading '{file_format}' format is not yet supported in Python SDK. "
f"Supported formats: parquet, orc, avro, lance, vortex, mosaic, blob, row.")
else:
raise ValueError(f"Unexpected file format: {file_format}")
index_mapping = self.create_index_mapping()
partition_info = self._create_partition_info()
system_fields = SpecialFields.find_system_fields(self.read_fields)
table_schema_fields = (
SpecialFields.row_type_with_row_tracking(self.table.table_schema.fields)
if row_tracking_enabled else self.table.table_schema.fields
)
# When native shard pushdown is used, the format reader only returns rows
# starting from shard_range[0], so _ROW_ID must be offset accordingly.
effective_first_row_id = file.first_row_id
if (shard_range is not None and file.first_row_id is not None
and file_format in (
CoreOptions.FILE_FORMAT_VORTEX, CoreOptions.FILE_FORMAT_LANCE)):
effective_first_row_id = file.first_row_id + shard_range[0]
if for_merge_read:
reader = DataFileBatchReader(
format_reader,
index_mapping,
partition_info,
self.trimmed_primary_key,
table_schema_fields,
file.max_sequence_number,
effective_first_row_id,
row_tracking_enabled,
system_fields,
file_io=self.table.file_io,
row_id_offsets=row_indices,
row_id_offset_ranges=parquet_row_ranges,
file_data_fields=file_read_fields,
target_data_fields=target_fields)
else:
reader = DataFileBatchReader(
format_reader,
index_mapping,
partition_info,
None,
table_schema_fields,
file.max_sequence_number,
effective_first_row_id,
row_tracking_enabled,
system_fields,
file_io=self.table.file_io,
row_id_offsets=row_indices,
row_id_offset_ranges=parquet_row_ranges,
file_data_fields=file_read_fields,
target_data_fields=target_fields)
# For non-Vortex formats, wrap with RowIdFilterRecordBatchReader
if (row_ranges is not None
and row_indices is None
and parquet_row_ranges is None):
reader = RowIdFilterRecordBatchReader(reader, file.first_row_id, effective_row_ranges)
# For formats without native shard support, wrap with ShardBatchReader
if shard_range is not None and file_format not in (
CoreOptions.FILE_FORMAT_VORTEX, CoreOptions.FILE_FORMAT_LANCE):
reader = ShardBatchReader(reader, shard_range[0], shard_range[1])
return reader
def _read_blob_as_descriptor(self, field_names: List[str]) -> bool:
if CoreOptions.blob_as_descriptor(self.table.options):
return True
deferred_fields = getattr(self, '_deferred_blob_fields', set())
return any(field_name in deferred_fields for field_name in field_names)
@staticmethod
def _row_sidecar_file_name(file: DataFileMeta) -> Optional[str]:
row_files = [
extra_file for extra_file in file.extra_files
if SplitRead._is_row_sidecar_file(extra_file)
]
return row_files[0] if len(row_files) == 1 else None
@staticmethod
def _should_read_row_sidecar(file: DataFileMeta,
effective_row_ranges: Optional[List[Range]],
row_sidecar_file: Optional[str],
max_selected_rows: int,
max_selection_ratio: float) -> bool:
if (not effective_row_ranges
or file.row_count <= 0
or DataFileMeta.is_blob_file(file.file_name)
or DataFileMeta.is_vector_file(file.file_name)
or row_sidecar_file is None):
return False
selected_row_count = sum(
r.count()
for r in Range.sort_and_merge_overlap(effective_row_ranges, True)
)
if selected_row_count <= 0 or selected_row_count >= file.row_count:
return False
selection_ratio = float(selected_row_count) / float(file.row_count)
return (selected_row_count <= max_selected_rows
and selection_ratio <= max_selection_ratio)
@staticmethod
def _is_row_sidecar_file(file_name: str) -> bool:
try:
return format_identifier(file_name) == ROW_SIDECAR_FORMAT
except Exception:
return False
@staticmethod
def _aligned_extra_file_path(file: DataFileMeta, extra_file: str) -> str:
if "://" in extra_file or extra_file.startswith("/"):
return extra_file
file_path = file.external_path if file.external_path else file.file_path
if not file_path or "/" not in file_path:
return extra_file
return f"{file_path.rsplit('/', 1)[0]}/{extra_file}"
def _get_fields_and_predicate(self, schema_id: int, read_fields):
key = (schema_id, tuple(read_fields))
if key not in self.schema_id_2_fields:
nested_path_by_name = self._nested_path_by_name()
schema = self._resolve_schema(schema_id)
schema_fields = (
SpecialFields.row_type_with_row_tracking(schema.fields)
if self.row_tracking_enabled else schema.fields
)
schema_field_names = set(field.name for field in schema_fields)
if self.table.is_primary_key_table:
schema_field_names.add('_SEQUENCE_NUMBER')
schema_field_names.add('_VALUE_KIND')
def _is_reachable(name: str) -> bool:
if name in schema_field_names:
return True
if nested_path_by_name is not None:
path = nested_path_by_name.get(name)
if path:
return path[0] in schema_field_names
return False
read_file_fields = [
read_field for read_field in read_fields
if _is_reachable(read_field)
]
read_predicate = trim_predicate_by_fields(self.push_down_predicate, read_file_fields)
read_arrow_predicate = read_predicate.to_arrow() if read_predicate else None
self.schema_id_2_fields[key] = (
read_file_fields,
read_arrow_predicate,
read_predicate,
)
return self.schema_id_2_fields[key]
@abstractmethod
def _all_data_fields_from(self, fields: List[DataField]) -> List[DataField]:
"""Apply this split-read's data-field shaping (row-tracking / kv
wrapping) to the given base ``fields``. Called both for the latest
table schema and for an older file schema."""
def _get_all_data_fields(self):
return self._all_data_fields_from(self.table.fields)
def _get_read_data_fields(self):
return self._read_data_fields_from(self._get_all_data_fields())
def _read_data_fields_from(self, all_data_fields):
read_field_ids = {field.id for field in self.read_fields}
return [f for f in all_data_fields if f.id in read_field_ids]
def _final_data_fields_from(self, all_data_fields: List[DataField]) -> List[DataField]:
"""The per-position target fields a batch must end up as: trimmed for
kv (``_KEY_*``) duplicates and stripped of partition columns. The
DataField analogue of ``_get_final_read_data_fields()``. Called with
the latest-schema fields it yields the read target (latest names +
types); called with a file's fields it yields what to physically read
from that file (the file's own names + types)."""
_, trimmed = self._get_trimmed_fields(
self._read_data_fields_from(all_data_fields), all_data_fields)
partition_keys = self.table.partition_keys
if not partition_keys:
return list(trimmed)
return [f for f in trimmed if f.name not in partition_keys]
def _target_read_fields(self) -> Optional[List[DataField]]:
"""Latest-schema target fields (names + types) that a normalized batch
must align to, in order. None for nested-projection reads (kept on the
legacy name-based path)."""
if self._nested_path_by_name() is not None:
return None
return self._final_data_fields_from(self._get_all_data_fields())
def _file_read_fields(self, file: DataFileMeta) -> Optional[List[DataField]]:
"""The fields to physically read from ``file``, in the file's own
names/types, selected by field id against the read set. None for
nested-projection reads."""
if self._nested_path_by_name() is not None:
return None
file_schema = self._resolve_schema(file.schema_id)
if file_schema is None:
return None
return self._final_data_fields_from(
self._all_data_fields_from(file_schema.fields))
def _create_key_value_fields(self, value_field: List[DataField]):
all_fields: List[DataField] = self.table.fields
all_data_fields = []
for field in all_fields:
if field.name in self.trimmed_primary_key:
key_field_name = f"{KEY_PREFIX}{field.name}"
key_field_id = field.id + KEY_FIELD_ID_START
key_field = DataField(key_field_id, key_field_name, field.type)
all_data_fields.append(key_field)
all_data_fields.append(SpecialFields.SEQUENCE_NUMBER)
all_data_fields.append(SpecialFields.VALUE_KIND)
for field in value_field:
all_data_fields.append(field)
return all_data_fields
def create_index_mapping(self):
if self._nested_path_by_name() is not None:
return None
base_index_mapping = self._create_base_index_mapping(self.read_fields, self._get_read_data_fields())
trimmed_key_mapping, _ = self._get_trimmed_fields(self._get_read_data_fields(), self._get_all_data_fields())
if base_index_mapping is None:
mapping = trimmed_key_mapping
elif trimmed_key_mapping is None:
mapping = base_index_mapping
else:
combined = [0] * len(base_index_mapping)
for i in range(len(base_index_mapping)):
if base_index_mapping[i] < 0:
combined[i] = base_index_mapping[i]
else:
combined[i] = trimmed_key_mapping[base_index_mapping[i]]
mapping = combined
if mapping is not None:
for i in range(len(mapping)):
if mapping[i] != i:
return mapping
return None
def _create_base_index_mapping(self, table_fields: List[DataField], data_fields: List[DataField]):
index_mapping = [0] * len(table_fields)
field_id_to_index = {field.id: i for i, field in enumerate(data_fields)}
for i, table_field in enumerate(table_fields):
field_id = table_field.id
data_field_index = field_id_to_index.get(field_id)
if data_field_index is not None:
index_mapping[i] = data_field_index
else:
index_mapping[i] = NULL_FIELD_INDEX
for i in range(len(index_mapping)):
if index_mapping[i] != i:
return index_mapping
return None
def _get_final_read_data_fields(self) -> List[str]:
if self._nested_path_by_name() is not None:
return self._remove_partition_fields(list(self.read_fields))
_, trimmed_fields = self._get_trimmed_fields(
self._get_read_data_fields(), self._get_all_data_fields()
)
return self._remove_partition_fields(trimmed_fields)
def _remove_partition_fields(self, fields: List[DataField]) -> List[str]:
partition_keys = self.table.partition_keys
if not partition_keys:
return [field.name for field in fields]
fields_without_partition = []
for field in fields:
if field.name not in partition_keys:
fields_without_partition.append(field)
return [field.name for field in fields_without_partition]
def _get_trimmed_fields(self, read_data_fields: List[DataField],
all_data_fields: List[DataField]) -> Tuple[List[int], List[DataField]]:
trimmed_mapping = [0] * len(read_data_fields)
trimmed_fields = []
field_id_to_field = {field.id: field for field in all_data_fields}
position_map = {}
for i, field in enumerate(read_data_fields):
is_key_field = field.name.startswith(KEY_PREFIX)
if is_key_field:
original_id = field.id - KEY_FIELD_ID_START
else:
original_id = field.id
original_field = field_id_to_field.get(original_id)
if original_id in position_map:
trimmed_mapping[i] = position_map[original_id]
else:
position = len(trimmed_fields)
position_map[original_id] = position
trimmed_mapping[i] = position
if is_key_field:
trimmed_fields.append(original_field)
else:
trimmed_fields.append(field)
return trimmed_mapping, trimmed_fields
def _create_partition_info(self):
if not self.table.partition_keys:
return None
partition_mapping = self._construct_partition_mapping()
if not partition_mapping:
return None
return PartitionInfo(partition_mapping, self.split.partition)
def _construct_partition_mapping(self) -> List[int]:
if self._nested_path_by_name() is not None:
partition_names = self.table.partition_keys
mapping = [0] * (len(self.read_fields) + 1)
p_count = 0
for i, field in enumerate(self.read_fields):
if field.name in partition_names:
partition_index = partition_names.index(field.name)
mapping[i] = -(partition_index + 1)
p_count += 1
else:
mapping[i] = (i - p_count) + 1
return mapping
_, trimmed_fields = self._get_trimmed_fields(
self._get_read_data_fields(), self._get_all_data_fields()
)
partition_names = self.table.partition_keys
mapping = [0] * (len(trimmed_fields) + 1)
p_count = 0
for i, field in enumerate(trimmed_fields):
if field.name in partition_names:
partition_index = partition_names.index(field.name)
mapping[i] = -(partition_index + 1)
p_count += 1
else:
mapping[i] = (i - p_count) + 1
return mapping
def _genarate_deletion_file_readers(self):
self.deletion_file_readers = {}
if self.split.data_deletion_files:
for data_file, deletion_file in zip(self.split.files, self.split.data_deletion_files):
if deletion_file is not None:
# Create a callable method to read the deletion vector
self.deletion_file_readers[data_file.file_name] = lambda df=deletion_file: DeletionVector.read(
self.table.file_io, df)
class RawFileSplitRead(SplitRead):
def __init__(
self,
table,
predicate: Optional[Predicate],
read_type: List[DataField],
split: Split,
row_tracking_enabled: bool,
outer_extract_name_paths: Optional[List[List[str]]] = None,
outer_flat_read_type: Optional[List[DataField]] = None,
limit: Optional[int] = None):
# Nested-leaf projection is NOT pushed down by name: a leaf path is
# only valid against the latest schema, while each data file stores
# its own (possibly renamed / retyped) sub-fields. Instead the read
# widens to the full top-level columns, which the per-file field-id
# normalization aligns to the latest schema, and the requested leaf
# paths are extracted afterwards (``outer_extract_name_paths``).
super().__init__(
table=table,
predicate=predicate,
read_type=read_type,
split=split,
row_tracking_enabled=row_tracking_enabled,
nested_name_paths=None,
limit=limit)
self.outer_extract_name_paths = outer_extract_name_paths
self.outer_flat_read_type = outer_flat_read_type
def raw_reader_supplier(self, file: DataFileMeta, dv_factory: Optional[Callable] = None) -> Optional[RecordReader]:
read_fields = self._get_final_read_data_fields()
# Check if this is a SlicedSplit to get shard_file_idx_map
shard_file_idx_map = (
self.split.shard_file_idx_map() if isinstance(self.split, SlicedSplit) else {}
)
if file.file_name in shard_file_idx_map:
(start_pos, end_pos) = shard_file_idx_map[file.file_name]
if (start_pos, end_pos) == (-1, -1):
return None
file_batch_reader = self.file_reader_supplier(
file=file,
for_merge_read=False,
read_fields=read_fields,
row_tracking_enabled=True,
shard_range=(start_pos, end_pos))
else:
file_batch_reader = self.file_reader_supplier(
file=file,
for_merge_read=False,
read_fields=read_fields,
row_tracking_enabled=True)
dv = dv_factory() if dv_factory else None
if dv:
if file.file_name in shard_file_idx_map:
dv = PositionMappedDeletionVector(
dv,
file_offset=start_pos,
)
return ApplyDeletionVectorReader(RowPositionReader(file_batch_reader), dv)
else:
return file_batch_reader
def create_reader(self) -> RecordReader:
self._genarate_deletion_file_readers()
data_readers = []
for file in self.split.files:
supplier = partial(
self.raw_reader_supplier,
file=file,
dv_factory=self.deletion_file_readers.get(file.file_name, None)
)
data_readers.append(supplier)
if not data_readers:
return EmptyFileRecordReader()
concat_reader = ConcatBatchReader(
data_readers, file_io=self.table.file_io,
blob_field_indices=blob_field_indices(self.read_fields),
vector_field_indices=vector_field_indices(self.read_fields))
reader = concat_reader
if self.table.is_primary_key_table and self.predicate_for_reader:
reader = FilterRecordBatchReader(
reader,
self.predicate_for_reader,
field_names=[f.name for f in self.read_fields],
schema_fields=self.read_fields,
)
if self.outer_extract_name_paths:
from pypaimon.read.reader.nested_leaf_batch_reader import \
NestedLeafBatchReader
reader = NestedLeafBatchReader(
reader, self.outer_extract_name_paths,
self.outer_flat_read_type)
# A predicate on a projected nested leaf cannot be pushed down:
# its leaf path is absent from the widened top-level read fields,
# so SplitRead.__init__ dropped it (predicate_for_reader is None).
# Without re-applying it the filter is silently lost and every row
# is returned. Re-evaluate it on the extracted flat batches, whose
# column names match the predicate fields.
if self.predicate is not None and self.predicate_for_reader is None:
flat_names = [f.name for f in self.outer_flat_read_type]
trimmed = trim_predicate_by_fields(self.predicate, flat_names)
if trimmed is not None:
reader = FilterRecordBatchReader(reader, trimmed)
if self.limit is not None:
reader = LimitedRecordBatchReader(reader, self.limit)
return reader
def _all_data_fields_from(self, fields):
if self.row_tracking_enabled:
return SpecialFields.row_type_with_row_tracking(fields)
return fields
class MergeFileSplitRead(SplitRead):
def __init__(
self,
table,
predicate: Optional[Predicate],
read_type: List[DataField],
split: Split,
row_tracking_enabled: bool,
outer_extract_name_paths: Optional[List[List[str]]] = None,
outer_flat_read_type: Optional[List[DataField]] = None,
limit: Optional[int] = None):
self.row_ranges = None
if isinstance(split, IndexedSplit):
self.row_ranges = split.row_ranges()
split = split.data_split()
# Merge functions need full ROW sub-structures, so nested paths
# are not pushed down here; sub-path extraction happens above
# the merge via OuterProjectionRecordReader.
super().__init__(
table=table,
predicate=predicate,
read_type=read_type,
split=split,
row_tracking_enabled=row_tracking_enabled,
nested_name_paths=None,
limit=limit,
)
self.outer_extract_name_paths = outer_extract_name_paths
self.outer_flat_read_type = outer_flat_read_type
# Built once per split-read (value_fields and options are constant
# for the object's life), not per section. ``None`` when
# ``sequence.field`` is unset, in which case the heap falls back to
# the file-level sequence number.
self.seq_comparator = builtin_seq_comparator(
self.value_fields,
self.table.options.sequence_field(),
self.table.options.sequence_field_sort_order_is_ascending(),
)
def kv_reader_supplier(self, file: DataFileMeta, dv_factory: Optional[Callable] = None) -> RecordReader:
file_batch_reader = self.file_reader_supplier(file, True, self._get_final_read_data_fields(), False)
selected_positions = None
if self.row_ranges is not None:
selected_positions = [
position
for row_range in self.row_ranges
for position in range(row_range.from_, row_range.to + 1)
]
file_batch_reader = RowIdFilterRecordBatchReader(
file_batch_reader, 0, self.row_ranges)
dv = dv_factory() if dv_factory else None
if dv:
if selected_positions is not None:
dv = PositionMappedDeletionVector(
dv, row_positions=selected_positions)
return ApplyDeletionVectorReader(
KeyValueWrapReader(file_batch_reader,
len(self.trimmed_primary_key), self.value_arity), dv)
else:
return KeyValueWrapReader(file_batch_reader, len(self.trimmed_primary_key), self.value_arity)
def section_reader_supplier(self, section: List[SortedRun]) -> RecordReader:
readers = []
for sorter_run in section:
data_readers = []
for file in sorter_run.files:
supplier = partial(self.kv_reader_supplier, file, self.deletion_file_readers.get(file.file_name, None))
data_readers.append(supplier)
readers.append(ConcatRecordReader(data_readers))
merge_function = self._build_merge_function()
return SortMergeReaderWithMinHeap(
readers, self.table.table_schema, merge_function=merge_function,
seq_comparator=self.seq_comparator)
def _build_merge_function(self):
"""Pick the MergeFunction for the table's ``merge-engine`` option.
Delegates to the shared dispatch in
``pypaimon.common.merge_engine_dispatch`` so the read path and
the in-memory merge buffer on the write path cannot drift.
``AGGREGATE`` is special-cased here because building the per-
field aggregators needs the full ``DataField`` objects, the
full primary-key list and the parsed ``CoreOptions`` -- which
sit outside the dispatch's raw-options contract. The writer-
side merge buffer falls back to dedupe for aggregation anyway
(see :meth:`FileStoreWrite._build_pk_merge_function`), so the
two sides only need to share the simple engines.
"""
engine = self.table.options.merge_engine()
if engine == MergeEngine.AGGREGATE:
# Use the full primary-key list, not ``trimmed_primary_key``:
# ``value_fields`` still carries partition columns, so any PK
# column that is also a partition column must be recognised
# as PK here. Otherwise a table with
# ``fields.default-aggregate-function`` would apply the
# default aggregator to that partition-PK column.
field_aggregators = build_field_aggregators(
self.value_fields,
self.table.primary_keys,
self.table.options,
)
return AggregateMergeFunction(
key_arity=len(self.trimmed_primary_key),
value_arity=self.value_arity,
field_aggregators=field_aggregators,
)
return build_merge_function(
engine=engine,
raw_options=self.table.options.options.to_map(),
key_arity=len(self.trimmed_primary_key),
value_arity=self.value_arity,
value_field_nullables=[f.type.nullable for f in self.value_fields],
value_field_names=[f.name for f in self.value_fields],
)
def create_reader(self) -> RecordReader:
# Create a dict mapping data file name to deletion file reader method
self._genarate_deletion_file_readers()
section_readers = []
sections = IntervalPartition(self.split.files).partition()
for section in sections:
supplier = partial(self.section_reader_supplier, section)
section_readers.append(supplier)
concat_reader = ConcatRecordReader(section_readers)
kv_unwrap_reader = KeyValueUnwrapRecordReader(DropDeleteRecordReader(concat_reader))
if self.predicate_for_reader:
reader = FilterRecordReader(kv_unwrap_reader, self.predicate_for_reader)
else:
reader = kv_unwrap_reader
if self.outer_extract_name_paths:
from pypaimon.read.reader.outer_projection_record_reader import \
OuterProjectionRecordReader
inner_value_fields = self.read_fields[-self.value_arity:]
reader = OuterProjectionRecordReader(
reader, [f.name for f in inner_value_fields],
self.outer_extract_name_paths,
file_io=self.table.file_io,
blob_field_indices=blob_field_indices(inner_value_fields),
vector_field_indices=vector_field_indices(inner_value_fields))
# A predicate on a projected nested leaf is not pushed down (its leaf
# path is absent from the widened-to-full-ROW read fields, so it was
# dropped in __init__). Without re-applying it after extraction the
# filter is silently lost. Evaluate it on the extracted flat rows,
# whose fields are outer_flat_read_type; trim to the projected
# columns and rewrite indices into that flat row.
if (self.predicate is not None and self.predicate_for_reader is None
and self.outer_flat_read_type is not None):
flat_names = [f.name for f in self.outer_flat_read_type]
trimmed = trim_predicate_by_fields(self.predicate, flat_names)
if trimmed is not None:
reader = FilterRecordReader(
reader,
rewrite_predicate_indices(
trimmed, self.outer_flat_read_type))
if self.limit is not None:
reader = LimitedRecordReader(reader, self.limit)
return reader
def _all_data_fields_from(self, fields):
return self._create_key_value_fields(fields)
class DataEvolutionSplitRead(SplitRead):
def __init__(
self,
table,
predicate: Optional[Predicate],
read_type: List[DataField],
split: Split,
row_tracking_enabled: bool,
nested_name_paths: Optional[List[List[str]]] = None,
limit: Optional[int] = None,
outer_extract_name_paths: Optional[List[List[str]]] = None,
outer_flat_read_type: Optional[List[DataField]] = None,
post_merge_filter=None,
eager_blob_fields=None,
post_filter_after_inline=False):
self.row_ranges = None
actual_split = split
if isinstance(split, IndexedSplit):
self.row_ranges = split.row_ranges()
actual_split = split.data_split()
super().__init__(
table, predicate, read_type, actual_split, row_tracking_enabled,
nested_name_paths=nested_name_paths,
limit=limit,
)
self.outer_extract_name_paths = outer_extract_name_paths
self.outer_flat_read_type = outer_flat_read_type
self._post_merge_filter = post_merge_filter
# Apply the auth filter after inline BLOB resolution, so scalar BLOBs still defer.
self._post_filter_after_inline = post_filter_after_inline
self._eager_blob_fields = set(eager_blob_fields or [])
self._deferred_blob_fields = self._deferred_blob_field_names()
def _push_down_predicate(self) -> Optional[Predicate]:
# Data evolution: files may have different schemas, so we don't push predicate
# to file readers; filtering is done in FilterRecordBatchReader after merge.
return None
def create_reader(self) -> RecordReader:
reader = self._create_raw_reader()
if ((CoreOptions.blob_view_fields(self.table.options) and CoreOptions.blob_view_resolve_enabled(
self.table.options))
or (not CoreOptions.blob_as_descriptor(self.table.options)
and CoreOptions.blob_descriptor_fields(self.table.options))):
blob_parallelism = self._blob_parallelism
reader = BlobInlineConvertReader(
reader, self.table,
prescan_reader_factory=lambda names: self._create_prescan_reader(names),
blob_parallelism=blob_parallelism)
if self._post_filter_after_inline:
if self._post_merge_filter is not None:
reader = AuthFilterReader(reader, self._post_merge_filter)
if self.limit is not None:
reader = LimitedRecordBatchReader(reader, self.limit)
if self._deferred_blob_fields:
blob_names = [
field.name for field in self.read_fields
if field.name in self._deferred_blob_fields
]
reader = DeferredBlobResolveReader(
reader,
self.table.file_io,
blob_names,
blob_parallelism=self._blob_parallelism,
)
return reader
def _deferred_blob_field_names(self) -> set:
return deferred_blob_field_names(
self.table,
self.read_fields,
self.predicate_for_reader,
self.limit,
has_post_filter=self._post_merge_filter is not None,
) - self._eager_blob_fields
def _create_raw_reader(self) -> RecordReader:
"""Core read logic: split_by_row_id -> suppliers -> ConcatBatchReader -> filter."""
files = self.split.files
suppliers = []
self._genarate_deletion_file_readers()
# Split files by row ID
split_by_row_id = self._split_by_row_id(files)
for need_merge_files in split_by_row_id:
deletion_vector = self._read_deletion_vector(need_merge_files)
if len(need_merge_files) == 1 or not self.read_fields:
# No need to merge fields, just create a single file reader
suppliers.append(
lambda f=need_merge_files[0], dv=deletion_vector: self._create_file_reader(
f, self._get_final_read_data_fields(), dv)
)
else:
suppliers.append(
lambda files=need_merge_files, dv=deletion_vector: self._create_union_reader(files, dv)
)
merge_reader = ConcatBatchReader(
suppliers, file_io=self.table.file_io,
blob_field_indices=blob_field_indices(self.read_fields),
vector_field_indices=vector_field_indices(self.read_fields))
if self.predicate_for_reader is not None:
reader = FilterRecordBatchReader(
merge_reader,
self.predicate_for_reader,
field_names=[f.name for f in self.read_fields],
schema_fields=self.read_fields,
)
else:
reader = merge_reader
if self._post_merge_filter is not None and not self._post_filter_after_inline:
reader = AuthFilterReader(reader, self._post_merge_filter)
if self.outer_extract_name_paths:
if self.outer_flat_read_type is None:
raise ValueError(
"outer_flat_read_type is required when outer_extract_name_paths "
"is set")
from pypaimon.read.reader.nested_leaf_batch_reader import \
NestedLeafBatchReader
reader = NestedLeafBatchReader(
reader, self.outer_extract_name_paths, self.outer_flat_read_type)
if self.limit is not None and not self._post_filter_after_inline:
reader = LimitedRecordBatchReader(reader, self.limit)
return reader
def _read_deletion_vector(self, need_merge_files: List[DataFileMeta]):
if not getattr(self, "deletion_file_readers", None):
return None
anchor = retrieve_anchor_file(need_merge_files)
dv_factory = self.deletion_file_readers.get(anchor.file_name)
if dv_factory is None:
return None
deletion_vector = dv_factory()
if deletion_vector is None:
return None
return anchor.row_id_range(), deletion_vector
def _apply_deletion_vector(self, reader, reader_range: Range, deletion_vector):
if reader is None or deletion_vector is None:
return reader
dv_range, dv = deletion_vector
if dv.is_empty():
return reader
if dv_range.from_ > reader_range.from_ or dv_range.to < reader_range.to:
raise ValueError(
f"Deletion vector range {dv_range} should contain reader range {reader_range}."
)
mapped_dv = PositionMappedDeletionVector(
dv,
reader_range.from_ - dv_range.from_,
self._selected_local_positions(reader_range),
)
return ApplyDeletionVectorReader(reader, mapped_dv)
def _selected_local_positions(self, reader_range: Range) -> Optional[List[int]]:
if self.row_ranges is None:
return None
selected = Range.and_([reader_range], self.row_ranges)
return [
row_id - reader_range.from_
for row_range in selected
for row_id in range(row_range.from_, row_range.to + 1)
]
def _create_prescan_reader(self, field_names):
"""Create a prescan reader by constructing a new DataEvolutionSplitRead
instance that only projects the specified field names.
Align with Java's configureBlobViewPrescanRead: pass limit to prescan reader
to avoid scanning entire split when there's a LIMIT clause.
"""
from pypaimon.read.reader.iface.record_batch_reader import EmptyRecordBatchReader
prescan_fields = [f for f in self.read_fields if f.name in field_names]
if not prescan_fields:
return EmptyRecordBatchReader()
# Skip limit push-down when the outer reader also selects rows (predicate or auth
# filter): prescan's first-N rows would differ from the outer set. TODO: push down.
skip_limit = self.predicate is not None or self._post_merge_filter is not None
prescan_read = DataEvolutionSplitRead(
table=self.table,
predicate=self.predicate,
read_type=prescan_fields,
split=self.split,
row_tracking_enabled=False,
limit=None if skip_limit else self.limit,
)
prescan_read.row_ranges = self.row_ranges
return prescan_read._create_raw_reader()
def _split_by_row_id(self, files: List[DataFileMeta]) -> List[List[DataFileMeta]]:
"""Split files by firstRowId for data evolution."""
# Sort files by firstRowId and then by maxSequenceNumber
def sort_key(file: DataFileMeta) -> tuple:
first_row_id = file.first_row_id if file.first_row_id is not None else float('-inf')
is_special = 1 if (DataFileMeta.is_blob_file(file.file_name)
or DataFileMeta.is_vector_file(file.file_name)) else 0
max_seq = file.max_sequence_number
return (first_row_id, is_special, -max_seq)
sorted_files = sorted(files, key=sort_key)
# Split files by firstRowId
split_by_row_id = []
last_row_id = -1
check_row_id_start = 0
current_split = []
for file in sorted_files:
first_row_id = file.first_row_id
if first_row_id is None:
split_by_row_id.append([file])
continue
if (not DataFileMeta.is_blob_file(file.file_name)
and not DataFileMeta.is_vector_file(file.file_name)
and first_row_id != last_row_id):
if current_split:
split_by_row_id.append(current_split)
if first_row_id < check_row_id_start:
raise ValueError(
f"There are overlapping files in the split: {files}, "
f"the wrong file is: {file}"
)
current_split = []
last_row_id = first_row_id
check_row_id_start = first_row_id + file.row_count
current_split.append(file)
if current_split:
split_by_row_id.append(current_split)
return split_by_row_id
def _create_union_reader(self, need_merge_files: List[DataFileMeta], deletion_vector=None) -> RecordReader:
"""Create a DataEvolutionFileReader for merging multiple files."""
# Split field bunches
fields_files = self._split_field_bunches(need_merge_files)
# Validate row counts and first row IDs (skip when row ranges are pushed down)
row_count = fields_files[0].row_count()
first_row_id = fields_files[0].files()[0].first_row_id
if self.row_ranges is None:
for bunch in fields_files:
if bunch.row_count() != row_count:
raise ValueError(
"All files in a field merge split should have the same row count.")
if bunch.files()[0].first_row_id != first_row_id:
raise ValueError(
"All files in a field merge split should have the same "
"first row id and could not be null."
)
# Create the union reader
all_read_fields = self.read_fields
file_record_readers = [None] * len(fields_files)
read_field_index = [field.id for field in all_read_fields]
# Initialize offsets
row_offsets = [-1] * len(all_read_fields)
field_offsets = [-1] * len(all_read_fields)
schema_pos = {f.id: p for p, f in enumerate(self.table.fields)}
for i, bunch in enumerate(fields_files):
first_file = bunch.files()[0]
# Get field IDs for this bunch
if DataFileMeta.is_blob_file(first_file.file_name):
field_ids = [self._get_field_id_from_write_cols(first_file)]
elif DataFileMeta.is_vector_file(first_file.file_name):
field_ids = self._get_field_ids_from_write_cols(first_file.write_cols)
elif first_file.write_cols:
field_ids = self._get_field_ids_from_write_cols(first_file.write_cols)
else:
# For regular files without write_cols, derive field IDs from
# the file's schema version, not the current table schema.
# The file only contains columns from when it was written.
file_schema = self._resolve_schema(first_file.schema_id)
field_ids = [field.id for field in file_schema.fields]
field_ids.append(SpecialFields.ROW_ID.id)
field_ids.append(SpecialFields.SEQUENCE_NUMBER.id)
read_fields = []
for j, read_field_id in enumerate(read_field_index):
for field_id in field_ids:
if read_field_id == field_id:
if row_offsets[j] == -1:
row_offsets[j] = i
read_fields.append(all_read_fields[j])
break
if not read_fields:
file_record_readers[i] = None
else:
read_fields.sort(key=lambda f: schema_pos.get(f.id, float('inf')))
id_to_pos = {f.id: p for p, f in enumerate(read_fields)}
for j in range(len(read_field_index)):
if row_offsets[j] == i:
field_offsets[j] = id_to_pos[read_field_index[j]]
read_field_names = self._remove_partition_fields(read_fields)
table_fields = self.read_fields
self.read_fields = read_fields # create reader based on read_fields
batch_size = self.table.options.read_batch_size()
# DataEvolutionMergeReader aligns fields by row ordinal, so every
# non-empty bunch reader created below must return the same row-id
# sequence. Keep row_ranges and the group-level deletion vector
# applied uniformly across normal, blob, and vector bunches.
if len(bunch.files()) == 1:
suppliers = [lambda r=self._create_file_reader(
bunch.files()[0], read_field_names, deletion_vector
): r]
file_record_readers[i] = MergeAllBatchReader(suppliers, batch_size=batch_size)
elif DataFileMeta.is_blob_file(first_file.file_name):
file_reader_suppliers = [
(
file,
partial(
self._create_raw_blob_file_reader,
file=file,
read_fields=read_field_names,
),
)
for file in bunch.files()
]
file_record_readers[i] = BlobFallbackBatchReader(
file_reader_suppliers,
read_fields[0].name,
PyarrowFieldParser.from_paimon_schema(
[read_fields[0]]
).field(0).type,
self.row_ranges,
self._read_blob_as_descriptor([read_fields[0].name]),
deletion_vector=deletion_vector,
batch_size=batch_size,
blob_parallelism=self._blob_parallelism,
)
else:
# Create concatenated reader for multiple files
suppliers = [
partial(self._create_file_reader, file=file,
read_fields=read_field_names,
deletion_vector=deletion_vector) for file in bunch.files()
]
file_record_readers[i] = MergeAllBatchReader(suppliers, batch_size=batch_size)
self.read_fields = table_fields
# Validate that all required fields are found
for i, field in enumerate(all_read_fields):
if row_offsets[i] == -1:
if not field.type.nullable:
raise ValueError(f"Field {field} is not null but can't find any file contains it.")
output_schema = PyarrowFieldParser.from_paimon_schema(all_read_fields)
return DataEvolutionMergeReader(row_offsets, field_offsets, file_record_readers, schema=output_schema)
def _create_file_reader(
self, file: DataFileMeta, read_fields: [str], deletion_vector=None) -> Optional[RecordReader]:
"""Create a file reader for a single file."""
reader = self.file_reader_supplier(
file=file,
for_merge_read=False,
read_fields=read_fields,
row_tracking_enabled=True,
row_ranges=self.row_ranges)
return self._apply_deletion_vector(reader, file.row_id_range(), deletion_vector)
def _create_raw_blob_file_reader(
self, file: DataFileMeta, read_fields: [str]) -> Optional[FormatBlobReader]:
row_indices = None
if self.row_ranges is not None:
row_indices = [
row_id - file.first_row_id
for row_range in Range.and_([file.row_id_range()], self.row_ranges)
for row_id in range(row_range.from_, row_range.to + 1)
]
if not row_indices:
return None
file_path = file.external_path if file.external_path else file.file_path
blob_parallelism = self._blob_parallelism
return FormatBlobReader(
self.table.file_io,
file_path,
read_fields,
self.read_fields,
None,
self._read_blob_as_descriptor(read_fields),
batch_size=self.table.options.read_batch_size(),
row_indices=row_indices,
blob_parallelism=blob_parallelism,
)
def _split_field_bunches(self, need_merge_files: List[DataFileMeta]) -> List[FieldBunch]:
"""Split files into field bunches."""
fields_files = []
blob_bunch_map = {}
vector_bunch_map = {}
row_count = -1
row_id_push_down = self.row_ranges is not None
for file in need_merge_files:
if DataFileMeta.is_blob_file(file.file_name):
field_id = self._get_field_id_from_write_cols(file)
if field_id not in blob_bunch_map:
blob_bunch_map[field_id] = BlobBunch(row_count, row_id_push_down)
blob_bunch_map[field_id].add(file)
elif DataFileMeta.is_vector_file(file.file_name):
field_id = self._get_field_id_from_write_cols(file)
if field_id not in vector_bunch_map:
vector_bunch_map[field_id] = VectorBunch(row_count, row_id_push_down)
vector_bunch_map[field_id].add(file)
else:
fields_files.append(DataBunch(file))
row_count = file.row_count
for bunch in blob_bunch_map.values():
bunch.finish()
fields_files.append(bunch)
fields_files.extend(vector_bunch_map.values())
return fields_files
def _get_field_id_from_write_cols(self, file: DataFileMeta) -> int:
"""Get field ID from write columns for blob/vector files."""
if not file.write_cols or len(file.write_cols) == 0:
raise ValueError("Blob/vector file must have write columns")
# Find the field by name in the table schema
field_name = file.write_cols[0]
for field in self.table.fields:
if field.name == field_name:
return field.id
raise ValueError(f"Field {field_name} not found in table schema")
def _get_field_ids_from_write_cols(self, write_cols: List[str]) -> List[int]:
field_ids = []
for field_name in write_cols:
for field in self.table.fields:
if field.name == field_name:
field_ids.append(field.id)
field_ids.append(SpecialFields.ROW_ID.id)
field_ids.append(SpecialFields.SEQUENCE_NUMBER.id)
return field_ids
def _all_data_fields_from(self, fields):
return SpecialFields.row_type_with_row_tracking(fields)