| # 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) |