diff --git a/docs/docs/concepts/spec/manifest.md b/docs/docs/concepts/spec/manifest.md index 7f8bf7a8a5e7..5b4271cc215e 100644 --- a/docs/docs/concepts/spec/manifest.md +++ b/docs/docs/concepts/spec/manifest.md @@ -88,6 +88,13 @@ Selected block bytes still share the manifest content cache without populating t whole-manifest entry cache with partial results. The low-level `build` method returns sidecar bytes without writing or publishing another file. +PyPaimon can read these sidecars and prune manifest blocks using partition, row-ID and bucket +filters. Its `manifest.sidecar.enabled` option inherits `manifest-sort.enabled` when unset. +Entry filters and ADD/DELETE reconciliation still apply after block selection. Missing or +unusable sidecars fall back to full manifest reads; scans without pruning filters and +explain/statistics scans do not perform sidecar I/O. The standalone codec can build sidecar +bytes, but automatic Python writer publication and cleanup are not integrated yet. + Callers decide whether to invoke `build` and `read`; these utilities have no read/write switches. `build` and `Builder` accept `rowIdEnabled` and `bucketEnabled` arguments for independent payload generation. Partition generation is always enabled, diff --git a/paimon-python/pypaimon/common/options/core_options.py b/paimon-python/pypaimon/common/options/core_options.py index 0ac461086425..ef9d45fd82bb 100644 --- a/paimon-python/pypaimon/common/options/core_options.py +++ b/paimon-python/pypaimon/common/options/core_options.py @@ -301,6 +301,20 @@ class CoreOptions: .with_description("The parallelism for scanning manifest files.") ) + MANIFEST_SIDECAR_ENABLED: ConfigOption[bool] = ( + ConfigOptions.key("manifest.sidecar.enabled") + .boolean_type() + .no_default_value() + .with_description("Enable sidecar pruning on reads. Defaults to manifest-sort.enabled when unset.") + ) + + MANIFEST_SORT_ENABLED: ConfigOption[bool] = ( + ConfigOptions.key("manifest-sort.enabled") + .boolean_type() + .default_value(False) + .with_description("Manifest sort setting. Also supplies the default for manifest sidecar reads.") + ) + MANIFEST_COMPRESSION: ConfigOption[str] = ( ConfigOptions.key("manifest.compression") .string_type() @@ -1266,6 +1280,13 @@ def manifest_target_size(self, default=None): default = MemorySize.of_bytes(default) if isinstance(default, int) else MemorySize.parse(default) return self.options.get(CoreOptions.MANIFEST_TARGET_FILE_SIZE, default).get_bytes() + def manifest_sidecar_enabled(self): + enabled = self.options.get(CoreOptions.MANIFEST_SIDECAR_ENABLED) + return self.manifest_sort_enabled() if enabled is None else enabled + + def manifest_sort_enabled(self): + return self.options.get(CoreOptions.MANIFEST_SORT_ENABLED) + def manifest_merge_skip_on_write_only(self, default=None): return self.options.get(CoreOptions.MANIFEST_MERGE_SKIP_ON_WRITE_ONLY, default) diff --git a/paimon-python/pypaimon/manifest/manifest_file_manager.py b/paimon-python/pypaimon/manifest/manifest_file_manager.py index 044d8985f9bc..1d833f3d3516 100644 --- a/paimon-python/pypaimon/manifest/manifest_file_manager.py +++ b/paimon-python/pypaimon/manifest/manifest_file_manager.py @@ -31,6 +31,9 @@ from datetime import datetime +from pypaimon.manifest.manifest_sidecar import ( + Query, read_sidecar, read_selected_bytes, +) from pypaimon.manifest.schema.data_file_meta import DataFileMeta from pypaimon.manifest.schema.manifest_entry import (MANIFEST_ENTRY_SCHEMA, ManifestEntry) @@ -128,13 +131,27 @@ def read_entries_parallel(self, manifest_files: List[ManifestFileMeta], manifest early_entry_filter: Optional[Callable[[int, int], bool]] = None, early_record_filter: Optional[Callable[[dict], bool]] = None, partition_filter=None, + row_ranges=None, ) -> List[ManifestEntry]: - def _process_single_manifest(manifest_file: ManifestFileMeta) -> List[ManifestEntry]: - return self.read(manifest_file.file_name, manifest_entry_filter, drop_stats, - early_entry_filter=early_entry_filter, - early_record_filter=early_record_filter, - partition_filter=partition_filter) + enabled = self.table.options.manifest_sidecar_enabled() + query = Query(row_ranges) if enabled and row_ranges is not None else None + + def _process_single_manifest(manifest_file: ManifestFileMeta): + path = f"{self.manifest_path}/{manifest_file.file_name}" + selected = None + if enabled and ( + query is not None or partition_filter is not None or early_entry_filter is not None): + selected = read_sidecar(self.file_io, path, manifest_file, query, + partition_filter, self.partition_keys_fields, early_entry_filter) + if selected is not None and not selected.blocks: + return [] + return self.read( + manifest_file.file_name, manifest_entry_filter, drop_stats, + early_entry_filter=early_entry_filter, + early_record_filter=early_record_filter, + partition_filter=partition_filter, + selected_blocks=selected) def _entry_identifier(e: ManifestEntry) -> tuple: return ( @@ -177,6 +194,7 @@ def read(self, manifest_file_name: str, manifest_entry_filter=None, drop_stats=T early_entry_filter: Optional[Callable[[int, int], bool]] = None, early_record_filter: Optional[Callable[[dict], bool]] = None, partition_filter=None, + selected_blocks=None, ) -> List[ManifestEntry]: """ early_entry_filter: ``(bucket, total_buckets) -> bool``, skip before deserializing _FILE. @@ -190,8 +208,11 @@ def read(self, manifest_file_name: str, manifest_entry_filter=None, drop_stats=T manifest_file_path = f"{self.manifest_path}/{manifest_file_name}" entries = [] - with self.file_io.new_input_stream(manifest_file_path) as input_stream: - avro_bytes = input_stream.read() + if selected_blocks is not None: + avro_bytes = read_selected_bytes(self.file_io, manifest_file_path, selected_blocks) + else: + with self.file_io.new_input_stream(manifest_file_path) as input_stream: + avro_bytes = input_stream.read() buffer = BytesIO(avro_bytes) records = _read_manifest_records( buffer, early_entry_filter, partition_filter, diff --git a/paimon-python/pypaimon/manifest/manifest_sidecar.py b/paimon-python/pypaimon/manifest/manifest_sidecar.py new file mode 100644 index 000000000000..8c1397333c21 --- /dev/null +++ b/paimon-python/pypaimon/manifest/manifest_sidecar.py @@ -0,0 +1,533 @@ +# 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. + +"""Independent partition, row-id and bucket coverage for each Avro block.""" + +import logging +import struct +import zlib +from bisect import bisect_left +from concurrent.futures import CancelledError +from dataclasses import dataclass +from io import BytesIO +from typing import Tuple + +from pyarrow import ArrowCancelled + +from pypaimon.utils.range import Range +from pypaimon.table.row.generic_row import GenericRowSerializer, GenericRowDeserializer + +LOG = logging.getLogger(__name__) +SUFFIX = '.avro.sidecar' +MAGIC = b'PMSC' +FORMAT_VERSION = 1 +MAX_ROW_ID = (1 << 63) - 1 +MAX_INT = (1 << 31) - 1 +READ_BUFFER_BYTES = 1024 * 1024 +LONG = struct.Struct('>q') +ROW_BOUNDS = struct.Struct('>qq') +_PROPAGATED_ERRORS = (InterruptedError, CancelledError, ArrowCancelled, MemoryError, RecursionError) + + +@dataclass(frozen=True) +class Settings: + + enabled: bool = False + row_id_enabled: bool = True + bucket_enabled: bool = True + + @classmethod + def from_options(cls, options): + return cls(options.manifest_sidecar_enabled(), + options.data_evolution_enabled(), options.bucket() != -1) + + +@dataclass(frozen=True) +class Block: + + offset: int + length: int + first_record: int + record_count: int + + +@dataclass(frozen=True) +class Selection: + + header: bytes + blocks: Tuple[Block, ...] + + +class Query: + + def __init__(self, ranges): + normalized = Range.sort_and_merge_overlap(list(ranges), True) + self.starts = [r.from_ for r in normalized] + self.ends = [r.to for r in normalized] + + def intersects(self, first, last): + candidate = bisect_left(self.ends, first) + return candidate < len(self.starts) and self.starts[candidate] <= last + + +def _require(condition): + if not condition: + raise ValueError('Invalid, unsupported or mismatched manifest sidecar') + + +def _varint(value, maximum=MAX_ROW_ID): + _require(0 <= value <= maximum) + result = bytearray() + while value >= 128: + result.append((value & 127) | 128) + value >>= 7 + result.append(value) + return result + + +def _deltas(values, base=0, signed=False): + result = _varint(len(values), MAX_INT) + for value in values: + _require(0 <= value <= (MAX_INT if signed else MAX_ROW_ID)) + delta = value - base + result.extend(_varint((delta << 1) ^ (delta >> 63) if signed else delta)) + base = value + return result + + +class _Buffer: + + """A bounded cursor over shared bytes; payloads need not allocate memoryview slices.""" + + __slots__ = ('data', 'position', 'limit') + + def __init__(self, data, position=0, limit=None): + self.data = data if isinstance(data, memoryview) else memoryview(data) + self.position = position + self.limit = len(self.data) if limit is None else limit + + @property + def remaining(self): + return self.limit - self.position + + def take(self, count): + start = self.position + end = start + count + if count < 0 or end > self.limit: + _require(False) + self.position = end + return self.data[start:end] + + def uint(self, maximum=MAX_ROW_ID): + # Do not create a memoryview slice or call take() for every encoded byte. + data, position = self.data, self.position + limit = self.limit + if position >= limit: + _require(False) + byte = data[position] + position += 1 + if byte < 128: + self.position = position + if byte > maximum: + _require(False) + return byte + value = byte & 127 + for shift in (7, 14, 21, 28, 35, 42, 49, 56): + if position >= limit: + self.position = position + _require(False) + byte = data[position] + position += 1 + value |= (byte & 127) << shift + if byte < 128: + self.position = position + if byte == 0 or value > maximum: + _require(False) + return value + self.position = position + raise ValueError('Invalid variable-length integer') + + def long(self): + position = self.position + if position + 8 > self.limit: + _require(False) + self.position = position + 8 + return LONG.unpack_from(self.data, position)[0] + + +def _delta_count(data): + count = data.uint(MAX_INT) + if count > data.limit - data.position: + _require(False) + return count + + +class _Deltas: + + __slots__ = ('data', 'value', 'maximum', 'signed', 'count', 'remaining') + + def __init__(self, data, base=0, maximum=MAX_ROW_ID, signed=False, count=None): + _require(0 <= base <= maximum and (not signed or maximum <= MAX_INT)) + self.data, self.value, self.maximum, self.signed = data, base, maximum, signed + # select() validates framing before filtering, but only needs a decoder + # for payloads whose enclosing block survives earlier filters. + if count is None: + count = _delta_count(data) + else: + _require(0 <= count <= MAX_INT and count <= data.remaining) + self.count = count + self.remaining = self.count + + def next(self): + _require(self.remaining > 0) + delta = self.data.uint(2 * MAX_INT if self.signed else MAX_ROW_ID) + if self.signed: + delta = (delta >> 1) ^ -(delta & 1) + _require(-self.value <= delta <= self.maximum - self.value) + self.value += delta + self.remaining -= 1 + return self.value + + +class Builder: + + def __init__(self, settings, header): + self.settings, self.header = settings, header + self.dictionary = {} + self.blocks = [] + self.ranges = [] + self.partition_ids = set() + self.bucket_pairs = set() + self.next_offset = len(header) + self.next_record = 0 + self.current = None + + def begin_block(self, offset, length, records): + _require(self.current is None and offset == self.next_offset and length > 0 and records > 0) + self.current = Block(offset, length, self.next_record, records) + self.entries_in_block = 0 + self.row_available = self.settings.row_id_enabled + self.partition_available = True + self.bucket_available = self.settings.bucket_enabled + self.ranges.clear() + self.partition_ids.clear() + self.bucket_pairs.clear() + + def add(self, first, count, partition=None, bucket=None, total_buckets=None): + _require(self.current is not None) + self.entries_in_block += 1 + if self.partition_available: + if partition is None: + self.partition_available = False + self.partition_ids.clear() + else: + partition = bytes(partition) + if partition not in self.dictionary: + self.dictionary[partition] = len(self.dictionary) + self.partition_ids.add(self.dictionary[partition]) + if self.bucket_available: + if bucket is None or total_buckets is None or not 0 <= bucket < total_buckets <= MAX_INT: + self.bucket_available = False + self.bucket_pairs.clear() + else: + self.bucket_pairs.add((bucket, total_buckets)) + if not self.row_available: + return + if first is None or not 0 <= first <= MAX_ROW_ID or count <= 0 or count - 1 > MAX_ROW_ID - first: + self.row_available = False + self.ranges.clear() + return + end = first + count - 1 + left = bisect_left(self.ranges, (first, -1)) + if left and self.ranges[left - 1][1] >= first - 1: + left -= 1 + right = left + while right < len(self.ranges) and self.ranges[right][0] <= end + 1: + first = min(first, self.ranges[right][0]) + end = max(end, self.ranges[right][1]) + right += 1 + self.ranges[left:right] = [(first, end)] + + def end_block(self): + block = self.current + _require(block is not None and self.entries_in_block == block.record_count) + partitions = _deltas(sorted(self.partition_ids)) if self.partition_available else b'' + rows = b'' + if self.row_available: + endpoints = [value for interval in self.ranges for value in interval] + rows = LONG.pack(endpoints[0]) + LONG.pack(endpoints[-1]) + rows += _deltas(endpoints[1:-1], endpoints[0]) + buckets = b'' + if self.bucket_available: + pairs = sorted(self.bucket_pairs) + buckets = _deltas([b for b, _ in pairs]) + _deltas([t for _, t in pairs], signed=True) + self.blocks.append((block, partitions, rows, buckets)) + self.next_offset = block.offset + block.length + self.next_record = block.first_record + block.record_count + self.current = None + + def serialize(self, file_size, entry_count): + _require(self.current is None and self.next_offset == file_size and self.next_record == entry_count) + data = bytearray(MAGIC) + _varint(FORMAT_VERSION) + _varint(len(self.header), MAX_INT) + self.header + data.extend(_varint(len(self.dictionary), MAX_INT)) + for partition in self.dictionary: + data.extend(_varint(len(partition), MAX_INT)) + data.extend(partition) + data.extend(_varint(len(self.blocks), MAX_INT)) + for block, *payloads in self.blocks: + for value in (block.offset, block.length, block.record_count): + data.extend(_varint(value)) + for payload in payloads: + data.append(1 if payload else 0) + if payload: + data.extend(_varint(len(payload), MAX_INT)) + data.extend(payload) + data.extend(struct.pack('>I', zlib.crc32(data))) + return bytes(data) + + +def build_from_entries(avro_bytes, entries, settings): + import fastavro + blocks = iter(fastavro.block_reader(BytesIO(avro_bytes))) + first = next(blocks, None) + builder = Builder(settings, avro_bytes[:first.offset] if first else avro_bytes) + position = 0 + block = first + while block is not None: + builder.begin_block(block.offset, block.size, block.num_records) + end = position + block.num_records + _require(end <= len(entries)) + for entry in entries[position:end]: + builder.add(entry.file.first_row_id, entry.file.row_count, + GenericRowSerializer.to_bytes(entry.partition), entry.bucket, entry.total_buckets) + builder.end_block() + position = end + block = next(blocks, None) + return builder.serialize(len(avro_bytes), len(entries)) + + +def _payload(data, payload): + position = data.position + if position >= data.limit: + _require(False) + encoding = data.data[position] + data.position = position + 1 + if encoding == 0: + return None + length = data.uint(MAX_INT) + start = data.position + end = start + length + if end > data.limit: + _require(False) + data.position = end + if encoding != 1: + return None + payload.position, payload.limit = start, end + return payload + + +def select(data, manifest, query, partition_filter=None, partition_fields=None, bucket_filter=None): + if query is not None and not isinstance(query, Query): + query = Query(query) + _require(len(data) >= 33) + view = memoryview(data) + _require(zlib.crc32(view[:-4]) == struct.unpack('>I', view[-4:])[0]) + stream = _Buffer(view[:-4]) + _require(bytes(stream.take(4)) == MAGIC and stream.uint(MAX_INT) == FORMAT_VERSION) + size = manifest.file_size + entries = manifest.num_added_files + manifest.num_deleted_files + _require(0 <= entries <= MAX_ROW_ID) + header_length = stream.uint(MAX_INT) + _require(21 <= header_length <= stream.remaining - 2 and header_length <= size) + header = bytes(stream.take(header_length)) + _require(header[:4] == b'Obj\x01') + partitions = stream.uint(MAX_INT) + _require(partitions <= stream.remaining // 13) + matches = [] if partition_filter is not None else None + unique = set() + for _ in range(partitions): + length = stream.uint(MAX_INT) + _require(length >= 12) + partition = bytes(stream.take(length)) + arity, = struct.unpack_from('>i', partition) + _require(arity >= 0 and 4 + ((arity + 71) // 64) * 8 + arity * 8 <= length) + _require(partition_fields is None or arity == len(partition_fields)) + _require(partition not in unique) + unique.add(partition) + if partition_filter is not None: + _require(partition_fields is not None) + matches.append(partition_filter.test(GenericRowDeserializer.from_bytes(partition, partition_fields))) + blocks = stream.uint(MAX_INT) + _require(blocks <= stream.remaining // 6) + next_offset, first_record = header_length, 0 + selected = [] + # Payloads are consumed within their block. Reuse bounded cursors within + # this invocation, not across files or concurrent queries. + partition_cursor, row_cursor, bucket_cursor = (_Buffer(stream.data) for _ in range(3)) + read_uint = stream.uint + intersects = None if query is None else query.intersects + for _ in range(blocks): + if stream.limit - stream.position < 6: + _require(False) + offset, length, count = read_uint(), read_uint(), read_uint() + if not (offset == next_offset and 0 < length <= size - offset + and 0 < count <= entries - first_record): + _require(False) + partitions_data = _payload(stream, partition_cursor) + row_data = _payload(stream, row_cursor) + bucket_data = _payload(stream, bucket_cursor) + if partitions_data is not None: + partition_count = _delta_count(partitions_data) + if not (0 < partition_count <= count and partition_count <= partitions): + _require(False) + if row_data is not None: + if row_data.limit - row_data.position < 17: + _require(False) + minimum, maximum = ROW_BOUNDS.unpack_from(row_data.data, row_data.position) + row_data.position += ROW_BOUNDS.size + if not 0 <= minimum <= maximum: + _require(False) + endpoint_count = _delta_count(row_data) + if (endpoint_count % 2 != 0 or endpoint_count // 2 >= count + or (endpoint_count == 0 and row_data.position != row_data.limit)): + _require(False) + if bucket_data is not None: + prefix = _Buffer(bucket_data.data, bucket_data.position, bucket_data.limit) + pairs = prefix.uint(MAX_INT) + _require(0 < pairs <= count and 2 * pairs + 1 <= prefix.remaining) + block_first_record = first_record + first_record += count + next_offset = offset + length + + if intersects is not None and row_data is not None: + if not intersects(minimum, maximum): + continue + ranges = endpoint_count // 2 + 1 + row_hit, start = ranges == 1, minimum + rows = None if row_hit else _Deltas(row_data, minimum, maximum, count=endpoint_count) + for i in range(ranges): + if row_hit: + break + end = maximum if i + 1 == ranges else rows.next() + _require(end >= start) + row_hit = intersects(start, end) + if not row_hit and i + 1 < ranges: + start = rows.next() + _require(start > end) + _require(rows.remaining or row_data.remaining == 0) + if not row_hit: + continue + + if partition_filter is not None and partitions_data is not None: + ids = _Deltas(partitions_data, maximum=partitions - 1, count=partition_count) + partition_hit, previous = False, -1 + while not partition_hit and ids.remaining: + id_ = ids.next() + _require(id_ > previous) + _require(ids.remaining or partitions_data.remaining == 0) + previous = id_ + partition_hit = matches[id_] + if not partition_hit: + continue + + if bucket_filter is not None and bucket_data is not None: + directory_data = _Buffer(bucket_data.data, bucket_data.position, bucket_data.limit) + directory = _Deltas(directory_data, maximum=MAX_INT) + while directory.remaining: + directory.next() + buckets = _Deltas(_Buffer(bucket_data.data, bucket_data.position, directory_data.position), maximum=MAX_INT) + totals_data = _Buffer(bucket_data.data, directory_data.position, bucket_data.limit) + totals = _Deltas(totals_data, maximum=MAX_INT, signed=True) + _require(totals.count == buckets.count) + previous, bucket_hit = (-1, -1), False + while not bucket_hit and buckets.remaining: + bucket, total = buckets.next(), totals.next() + _require(total > bucket and (bucket, total) > previous) + _require(buckets.remaining or (buckets.data.remaining == 0 and totals_data.remaining == 0)) + previous = (bucket, total) + bucket_hit = bucket_filter(bucket, total) + if not bucket_hit: + continue + selected.append(Block(offset, length, block_first_record, count)) + _require(stream.remaining == 0 and next_offset == size and first_record == entries) + return Selection(header, tuple(selected)) + + +def sidecar_file_name(manifest): + return next((name for name in manifest.extra_files or [] if name.endswith(SUFFIX)), None) + + +def read_sidecar(file_io, manifest_path, manifest, query, partition_filter=None, partition_fields=None, + bucket_filter=None): + name = sidecar_file_name(manifest) + if name is None: + return None + sidecar_path = manifest_path.rsplit('/', 1)[0] + '/' + name + try: + with file_io.new_input_stream(sidecar_path) as stream: + data = bytearray() + while True: + chunk = stream.read(READ_BUFFER_BYTES) + if not chunk: + break + data.extend(chunk) + return select(data, manifest, query, partition_filter, partition_fields, bucket_filter) + except _PROPAGATED_ERRORS: + raise + except Exception as error: + pending, visited = [error], set() + while pending: + cause = pending.pop() + if id(cause) in visited: + continue + visited.add(id(cause)) + if not isinstance(cause, Exception) or isinstance(cause, _PROPAGATED_ERRORS): + raise cause + if cause.__cause__ is not None: + pending.append(cause.__cause__) + if cause.__context__ is not None: + pending.append(cause.__context__) + LOG.debug('Cannot use manifest sidecar for %s; reading manifest: %s', manifest_path, error) + return None + + +def read_selected_bytes(file_io, manifest_path, selected): + """Read complete selected blocks with seek; adjacent blocks share one contiguous span. + + The concatenated original header and blocks form a valid Avro OCF. Partial entries + must not be stored in a cache keyed by the complete manifest. + """ + data = bytearray(selected.header) + with file_io.new_input_stream(manifest_path) as stream: + block_position = 0 + while block_position < len(selected.blocks): + block = selected.blocks[block_position] + block_position += 1 + end = block.offset + block.length + while (block_position < len(selected.blocks) + and selected.blocks[block_position].offset == end): + end += selected.blocks[block_position].length + block_position += 1 + stream.seek(block.offset) + remaining = end - block.offset + while remaining: + chunk = stream.read(min(remaining, READ_BUFFER_BYTES)) + if not chunk: + raise EOFError('Truncated manifest block') + data.extend(chunk) + remaining -= len(chunk) + return bytes(data) diff --git a/paimon-python/pypaimon/read/scanner/file_scanner.py b/paimon-python/pypaimon/read/scanner/file_scanner.py index c07bd7e2d3a8..bc89cf16d62c 100755 --- a/paimon-python/pypaimon/read/scanner/file_scanner.py +++ b/paimon-python/pypaimon/read/scanner/file_scanner.py @@ -599,7 +599,7 @@ def read_manifest_entries(self, manifest_files: List[ManifestFileMeta], self.scan_stats.manifest_files_after_partition += len(manifest_files) # Force single-threaded so we can mutate stats without locking. max_workers = 1 - # Disable both early filters in explain mode (scan_stats) so all entries + # Disable early entry filters and sidecar pruning in explain mode so all entries # flow through _filter_manifest_entry for accurate funnel counting. early_row_filter = None if self.scan_stats is not None \ else _build_early_row_range_filter(row_ranges) @@ -613,6 +613,7 @@ def read_manifest_entries(self, manifest_files: List[ManifestFileMeta], early_entry_filter=self._build_early_bucket_filter(), early_record_filter=early_row_filter, partition_filter=partition_filter, + row_ranges=row_ranges if self.scan_stats is None else None, ) def _build_early_bucket_filter(self): diff --git a/paimon-python/pypaimon/tests/file_type_test.py b/paimon-python/pypaimon/tests/file_type_test.py index b82604eaede1..0efabd1833ee 100644 --- a/paimon-python/pypaimon/tests/file_type_test.py +++ b/paimon-python/pypaimon/tests/file_type_test.py @@ -48,6 +48,9 @@ def test_manifest(self): self.assertEqual(FileType.META, FileType.classify("manifest-abc123")) self.assertEqual(FileType.META, FileType.classify("manifest-list-abc")) self.assertEqual(FileType.META, FileType.classify("index-manifest-abc")) + self.assertEqual(FileType.META, FileType.classify("manifest-abc123.avro.sidecar")) + self.assertEqual(FileType.META, FileType.classify("custom.avro.sidecar")) + self.assertEqual(FileType.META, FileType.classify("index-abc123.avro.sidecar")) def test_hint_files(self): self.assertEqual(FileType.META, FileType.classify("EARLIEST")) diff --git a/paimon-python/pypaimon/tests/manifest/manifest_block_index_test.py b/paimon-python/pypaimon/tests/manifest/manifest_block_index_test.py new file mode 100644 index 000000000000..b031e5f57f1f --- /dev/null +++ b/paimon-python/pypaimon/tests/manifest/manifest_block_index_test.py @@ -0,0 +1,542 @@ +# 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 zlib +import struct +import random +from pathlib import Path +from types import SimpleNamespace +from unittest.mock import Mock, patch + +import pytest + +from pypaimon.common.options.core_options import CoreOptions +from pypaimon.common.options.options import Options +from pypaimon.manifest import manifest_sidecar +from pypaimon.manifest.manifest_sidecar import Builder, Settings, select +from pypaimon.schema.data_types import AtomicType, DataField +from pypaimon.table.row.generic_row import GenericRow, GenericRowSerializer +from pypaimon.tests.manifest.manifest_sidecar_test import avro_header, golden, golden_meta, make_sidecar +from pypaimon.utils.range import Range + +FIELDS = [DataField(0, 'p', AtomicType('INT')), DataField(1, 'q', AtomicType('STRING'))] + + +def partition(p, q): + return GenericRowSerializer.to_bytes(GenericRow([p, q], FIELDS)) + + +def part(p): + return SimpleNamespace(test=lambda row: row.values[0] == p) + + +def fixture(key): + if key == 'avroHeader': + return avro_header() + partitions = (partition(7, 'left'), partition(9, None)) + if key == 'partitionA': + return partitions[0] + if key == 'partitionB': + return partitions[1] + return make_sidecar(partitions if key != 'index' else None, key == 'indexWithBuckets') + + +def meta(name, size, count): + return SimpleNamespace(file_name=name, file_size=size, num_added_files=count, num_deleted_files=0) + + +def test_v1_golden_file_matches_java(): + path = Path(__file__).resolve().parents[4] / 'paimon-core/src/test/resources/compatibility/manifest-sidecar-v1' + data = path.read_bytes() + header = avro_header() + block_length = 1024 * 1024 + 17 + builder = Builder(Settings(), header) + for block in range(10): + builder.begin_block(len(header) + block * block_length, block_length, 16) + step = 16400 if block % 4 == 0 else 16 + for entry in range(16): + p = (block * 7 + entry * 3) % 16 + builder.add((block << 40) + entry * step, 3, + partition(p, None if p % 7 == 0 else 'partition-' + str(p)), + block + entry * 13, 512 if entry % 2 == 0 else 1024) + builder.end_block() + file_size = len(header) + 10 * block_length + assert data == builder.serialize(file_size, 160) + metadata = meta('arbitrary-name', file_size, 160) + assert len(select(data, metadata, None).blocks) == 10 + for block in range(10): + step = 16400 if block % 4 == 0 else 16 + for entry in range(16): + row_id = (block << 40) + entry * step + p = (block * 7 + entry * 3) % 16 + pair = (block + entry * 13, 512 if entry % 2 == 0 else 1024) + selected = select(data, metadata, [Range(row_id, row_id + 2)], part(p), FIELDS, + lambda b, t: (b, t) == pair) + assert selected.blocks == (manifest_sidecar.Block( + len(header) + block * block_length, block_length, block * 16, 16),) + assert not select(data, metadata, [Range(row_id + 3, row_id + 3)]).blocks + assert not select(data, metadata, None, part(99), FIELDS).blocks + assert not select(data, metadata, None, bucket_filter=lambda b, t: t == 4096).blocks + for length in (0, 32, len(data) // 2, len(data) - 1): + with pytest.raises(ValueError): + select(data[:length], metadata, None) + + +@pytest.mark.parametrize('partition_enabled', [False, True]) +@pytest.mark.parametrize('row_id_enabled', [False, True]) +@pytest.mark.parametrize('bucket_count', [-2, -1, 4]) +def test_settings_derive_payload_generation_from_table_metadata(partition_enabled, row_id_enabled, bucket_count): + header = avro_header() + bucket_enabled = bucket_count != -1 + settings = Settings.from_options(CoreOptions(Options({ + 'data-evolution.enabled': row_id_enabled, 'bucket': bucket_count}))) + builder = Builder(settings, header) + fields = FIELDS if partition_enabled else [] + for block in range(2): + builder.begin_block(len(header) + block * 100, 100, 1) + p = partition(7 + block, 'p') if partition_enabled else GenericRowSerializer.to_bytes(GenericRow([], [])) + builder.add(100 + block * 100, 10, p, 1, 4) + builder.end_block() + data = builder.serialize(len(header) + 200, 2) + for position in positions(data): + assert data[position[1]] == 1 + assert data[position[2]] == int(row_id_enabled) + assert data[position[3]] == int(bucket_enabled) + metadata = meta('m', len(header) + 200, 2) + assert len(select(data, metadata, [Range(999, 999)]).blocks) == (0 if row_id_enabled else 2) + assert not select(data, metadata, None, SimpleNamespace(test=lambda row: False), fields).blocks + assert len(select(data, metadata, None, bucket_filter=lambda b, t: b == 99).blocks) == (0 if bucket_enabled else 2) + + +def test_partition_dictionary_golden_tuples_nulls_and_derived_ordinals(): + a, b, header = partition(7, 'left'), partition(9, None), avro_header() + assert a == fixture('partitionA') + assert b == fixture('partitionB') + builder = Builder(Settings(), header) + values = [(0, 100, [(0, 10, a), (5, 5, a), (20, 5, b)]), + (100, 200, [((1 << 32) - 2, 5, b), (8254058425445, 1, a)]), + (300, 100, [(20, 5, a), ((1 << 63) - 1, 1, b)])] + for offset, length, entries in values: + builder.begin_block(len(header) + offset, length, len(entries)) + for first, count, p in entries: + builder.add(first, count, p) + builder.end_block() + data = builder.serialize(len(header) + 400, 7) + assert data == fixture('indexWithPartitions') + predicate = part(7) + with patch.object(predicate, 'test', wraps=predicate.test) as evaluated: + selected = select(data, golden_meta(), [Range(20, 20)], predicate, FIELDS) + assert evaluated.call_count == 2 + assert [b.first_record for b in selected.blocks] == [0, 5] + nulls = SimpleNamespace(test=lambda row: row.values[1] is None) + assert len(select(data, golden_meta(), None, nulls, FIELDS).blocks) == 3 + assert not select(data, golden_meta(), None, part(99), FIELDS).blocks + assert len(select(golden(), golden_meta(), None, part(99), FIELDS).blocks) == 3 + + +def test_unavailable_dimensions_are_independent_and_dictionary_misses_keep_unknown_blocks(): + settings, header = Settings(), avro_header() + builder = Builder(settings, header) + for i, (first, p) in enumerate([(None, partition(7, 'left')), (200, partition(9, 'x' * 600)), + (300, partition(7, 'left'))]): + builder.begin_block(len(header) + 100 * i, 100, 1) + builder.add(first, 10, p) + builder.end_block() + data = builder.serialize(len(header) + 300, 3) + metadata = meta('m', len(header) + 300, 3) + assert [b.first_record for b in select(data, metadata, None, part(9), FIELDS).blocks] == [1] + selected = select(data, metadata, [Range(999, 999)], part(7), FIELDS) + assert [b.first_record for b in selected.blocks] == [0] + selected = select(data, metadata, [Range(200, 200)], part(9), FIELDS) + assert [b.first_record for b in selected.blocks] == [1] + + +@pytest.mark.parametrize('unknown', [False, True]) +def test_exact_row_ranges_keep_gaps_and_detect_late_unknowns(unknown): + settings, header = Settings(), avro_header() + builder = Builder(settings, header) + builder.begin_block(len(header), 100, 66) + for i in range(64): + builder.add(100 + i * 1000, 10, partition(7, 'left')) + for first, count in [(10, 10), (None if unknown else (1 << 63) - 1, 1)]: + builder.add(first, count, partition(7, 'left')) + builder.end_block() + data = builder.serialize(len(header) + 100, 66) + metadata = meta('m', len(header) + 100, 66) + for point in [10, 100, (1 << 63) - 1]: + assert len(select(data, metadata, [Range(point, point)]).blocks) == 1 + assert len(select(data, metadata, [Range(0, 0)]).blocks) == int(unknown) + assert len(select(data, metadata, [Range(200, 200)]).blocks) == int(unknown) + assert not select(data, metadata, None, part(9), FIELDS).blocks + + +def positions(data): + reader = manifest_sidecar._Buffer(data[:-4]) + reader.take(4) + reader.uint() + reader.take(reader.uint()) + for _ in range(reader.uint()): + reader.take(reader.uint()) + result = [] + for _ in range(reader.uint()): + block = reader.position + for _ in range(3): + reader.uint() + payloads = [] + for _ in range(3): + payloads.append(reader.position) + if reader.take(1)[0]: + reader.take(reader.uint()) + result.append((block, *payloads)) + return result + + +def checksum(data): + data[-4:] = struct.pack('>I', zlib.crc32(data[:-4])) + return data + + +def replace_payload(data, block, dimension, payload, encoding=1): + start = positions(data)[block][dimension] + reader = manifest_sidecar._Buffer(data) + reader.position = start + if reader.take(1)[0]: + reader.take(reader.uint()) + framed = bytes([encoding]) + if encoding: + framed += manifest_sidecar._varint(len(payload)) + payload + return checksum(bytearray(data[:start] + framed + data[reader.position:])) + + +def vint(value): + return bytes(manifest_sidecar._varint(value)) + + +def row_payload(minimum, maximum, deltas): + return struct.pack('>qq', minimum, maximum) + vint(len(deltas)) + b''.join(vint(d) for d in deltas) + + +def test_row_miss_skips_partition_and_bucket_payloads(): + data = replace_payload(fixture('indexWithBuckets'), 0, 1, b'\2' + vint(999) + b'\1') + data = replace_payload(data, 0, 3, b'\2\0\0\2\0\0') + buckets = Mock(return_value=True) + assert not select(data, golden_meta(), [Range(15, 15)], part(7), FIELDS, buckets).blocks + buckets.assert_not_called() + + +@pytest.mark.parametrize('row_ranges', [None, [Range(0, 0)]]) +def test_partition_miss_skips_bucket_matching(row_ranges): + data = replace_payload(fixture('indexWithBuckets'), 0, 3, b'\2\0\0\2\0\0') + buckets = Mock(return_value=True) + assert not select(data, golden_meta(), row_ranges, part(99), FIELDS, buckets).blocks + buckets.assert_not_called() + + +def test_absent_partition_filter_keeps_row_and_bucket_matching(): + data = replace_payload(fixture('indexWithBuckets'), 0, 1, b'\2' + vint(999) + b'\1') + with patch('pypaimon.manifest.manifest_sidecar.GenericRowDeserializer.from_bytes') as decode: + selected = select(data, golden_meta(), [Range(0, 0)], None, FIELDS, lambda b, t: b == 1) + assert [b.first_record for b in selected.blocks] == [0] + decode.assert_not_called() + + +def test_absent_row_or_bucket_filters_keep_remaining_dimensions(): + data = fixture('indexWithBuckets') + selected = select(data, golden_meta(), None, part(7), FIELDS, lambda b, t: b == 1) + assert [b.first_record for b in selected.blocks] == [0] + assert [b.first_record for b in select(data, golden_meta(), [Range(20, 20)], part(7), FIELDS).blocks] == [0, 5] + assert len(select(data, golden_meta(), None).blocks) == 3 + + +def test_partition_and_bucket_matches_skip_unused_payload_elements(): + data = replace_payload(fixture('indexWithBuckets'), 0, 1, b'\2\0' + vint(999)) + assert [b.first_record for b in select(data, golden_meta(), [Range(0, 0)], part(7), FIELDS).blocks] == [0] + with pytest.raises(ValueError): + select(data, golden_meta(), [Range(0, 0)], part(99), FIELDS) + data = replace_payload(fixture('indexWithBuckets'), 0, 3, b'\2\1\0\2\10\0') + selected = select(data, golden_meta(), [Range(0, 0)], None, FIELDS, lambda b, t: b == 1) + assert [b.first_record for b in selected.blocks] == [0] + with pytest.raises(ValueError): + select(data, golden_meta(), [Range(0, 0)], None, FIELDS, lambda b, t: False) + + +def test_skipped_payloads_still_require_valid_framing_and_directory(): + for dimension in (1, 2, 3): + for payload in (b'', b'\x80', vint(1 << 31)): + data = replace_payload(fixture('indexWithBuckets'), 0, dimension, payload) + with pytest.raises(ValueError): + select(data, golden_meta(), [Range(999, 999)]) + data = bytearray(fixture('indexWithBuckets')) + reader = manifest_sidecar._Buffer(data) + reader.position = positions(data)[0][0] + reader.uint() + reader.uint() + data[reader.position] = 2 + with pytest.raises(ValueError): + select(checksum(data), golden_meta(), [Range(999, 999)]) + + +def test_unsigned_unknown_encodings_skip_only_one_payload_and_validate_lengths(): + for dimension in (1, 2, 3): + data = replace_payload(fixture('indexWithBuckets'), 0, dimension, b'\x80', 202) + point = 999 if dimension == 2 else 0 + selected = select(data, golden_meta(), [Range(point, point)], + part(99 if dimension == 1 else 7), FIELDS, + lambda b, t: b == (99 if dimension == 3 else 1)) + assert [b.first_record for b in selected.blocks] == [0] + start = positions(data)[0][dimension] + data[start + 1] = 127 + with pytest.raises(ValueError): + select(checksum(data), golden_meta(), None) + + +def test_complete_directory_and_payloads_are_not_dropped(): + header = avro_header() + builder = Builder(Settings(), header) + for i in range(3): + builder.begin_block(len(header) + 100 * i, 100, 1) + builder.add(100 * i, 10, partition(7, 'left')) + builder.end_block() + data = builder.serialize(len(header) + 300, 3) + assert not select(data, meta('m', len(header) + 300, 3), [Range(999, 999)], part(99), FIELDS).blocks + assert len(select(data, meta('m', len(header) + 300, 3), None).blocks) == 3 + assert builder.serialize(len(header) + 300, 3) == data + builder.begin_block(len(header) + 300, 100, 1) + builder.add(300, 1, partition(7, 'left')) + builder.end_block() + data = builder.serialize(len(header) + 400, 4) + assert len(select(data, meta('m', len(header) + 400, 4), None).blocks) == 4 + + +def test_randomized_exact_coverage_has_no_false_negatives(): + rng = random.Random(9743) + header = avro_header() + for _ in range(60): + settings = Settings() + builder = Builder(settings, header) + blocks = [] + for i in range(5): + values = [(rng.choice([None, rng.randrange(100)]), rng.randrange(1, 10), rng.randrange(5)) + for _ in range(6)] + blocks.append(values) + builder.begin_block(len(header) + 100 * i, 100, len(values)) + for first, count, p in values: + builder.add(first, count, partition(p, None)) + builder.end_block() + data = builder.serialize(len(header) + 500, 30) + assert data is not None + metadata = meta('m', len(header) + 500, 30) + for point in range(0, 110, 11): + for p in range(5): + selected = select(data, metadata, [Range(point, point)], part(p), FIELDS) + ordinals = {b.first_record for b in selected.blocks} + expected = set() + for i, values in enumerate(blocks): + if (any(row_p == p for _, _, row_p in values) + and any(first is None or first <= point < first + count for first, count, _ in values)): + expected.add(i * 6) + assert ordinals == expected + + +def test_bucket_payload_golden_rescale_and_unavailable_payloads(): + a, b, header = partition(7, 'left'), partition(9, None), avro_header() + builder = Builder(Settings(), header) + for offset, length, values in [ + (0, 100, [(0, 10, a, 1, 4), (5, 5, a, 1, 4), (20, 5, b, 1, 8)]), + (100, 200, [((1 << 32) - 2, 5, b, 2, 4), (8254058425445, 1, a, 2, 8)]), + (300, 100, [(20, 5, a, 0, 1), ((1 << 63) - 1, 1, b, 3, 4)])]: + builder.begin_block(len(header) + offset, length, len(values)) + for value in values: + builder.add(*value) + builder.end_block() + data = builder.serialize(len(header) + 400, 7) + assert data == fixture('indexWithBuckets') + selected = select(data, golden_meta(), None, bucket_filter=lambda bucket, total: bucket == 1) + assert [b.first_record for b in selected.blocks] == [0] + selected = select(data, golden_meta(), None, bucket_filter=lambda bucket, total: (bucket, total) == (2, 8)) + assert [b.first_record for b in selected.blocks] == [3] + assert not select(data, golden_meta(), [Range(0, 0)], part(7), FIELDS, + bucket_filter=lambda bucket, total: bucket == 2).blocks + for data in [golden(), fixture('indexWithPartitions')]: + assert len(select(data, golden_meta(), None, bucket_filter=lambda bucket, total: False).blocks) == 3 + + +@pytest.mark.parametrize('pair', [(None, None), (-1, 4), (4, 4), (0, 0), (2, 8)]) +def test_unknown_pairs_disable_only_bucket_payload(pair): + header = avro_header() + builder = Builder(Settings(), header) + builder.begin_block(len(header), 100, 2) + builder.add(100, 10, partition(7, 'left'), 1, 4) + builder.add(200, 10, partition(7, 'left'), *pair) + builder.end_block() + data = builder.serialize(len(header) + 100, 2) + metadata = meta('m', len(header) + 100, 2) + valid = pair == (2, 8) + assert data[positions(data)[0][3]] == int(valid) + assert len(select(data, metadata, None, bucket_filter=lambda b, t: False).blocks) == int(not valid) + assert not select(data, metadata, [Range(999, 999)]).blocks + + +def test_large_payloads_keep_exact_coverage(): + header, settings = avro_header(), Settings() + builder = Builder(settings, header) + blocks, entries_per_block, block_bytes = 33, 4097, 1024 * 1024 + entries = blocks * entries_per_block + for block in range(blocks): + builder.begin_block(len(header) + block * block_bytes, block_bytes, entries_per_block) + for i in range(entries_per_block): + entry = block * entries_per_block + i + builder.add(entry * 2, 1, partition(entry, None), i, entries_per_block + 1) + builder.end_block() + file_size = len(header) + blocks * block_bytes + data = builder.serialize(file_size, entries) + assert data is not None + metadata = meta('m', file_size, entries) + last = (entries - 1) * 2 + selected = select(data, metadata, [Range(last, last)]) + assert [b.first_record for b in selected.blocks] == [(blocks - 1) * entries_per_block] + assert not select(data, metadata, [Range(last - 1, last - 1)]).blocks + assert not select(data, metadata, None, part(entries), FIELDS).blocks + assert not select(data, metadata, None, bucket_filter=lambda bucket, total: bucket == entries_per_block).blocks + + +def test_absent_payloads_omit_length_fields_for_every_dimension_combination(): + header, settings = avro_header(), Settings() + builder = Builder(settings, header) + for mask in range(8): + builder.begin_block(len(header) + mask * 100, 100, 1) + builder.add(100 + mask if mask & 2 else None, 1, + partition(7, 'left') if mask & 1 else None, + 0 if mask & 4 else None, 1 if mask & 4 else None) + builder.end_block() + data = builder.serialize(len(header) + 800, 8) + locations = positions(data) + for mask in range(8): + for dimension, present_size in enumerate((4, 19, 6)): + start = locations[mask][dimension + 1] + end = (locations[mask][dimension + 2] if dimension < 2 + else locations[mask + 1][0] if mask < 7 else len(data) - 4) + present = bool(mask & (1 << dimension)) + assert data[start] == int(present) + assert end - start == (present_size if present else 1) + metadata = meta('m', len(header) + 800, 8) + selected = select(data, metadata, None, part(99), FIELDS) + assert [b.first_record for b in selected.blocks] == [0, 2, 4, 6] + selected = select(data, metadata, [Range(999, 999)]) + assert [b.first_record for b in selected.blocks] == [0, 1, 4, 5] + buckets = lambda bucket, total: False + selected = select(data, metadata, None, None, FIELDS, buckets) + assert [b.first_record for b in selected.blocks] == [0, 1, 2, 3] + selected = select(data, metadata, [Range(999, 999)], part(99), FIELDS, buckets) + assert [b.first_record for b in selected.blocks] == [0] + + +def test_malformed_bucket_payload_invalidates_the_container(): + for payload in (b'\0', b'\x80', b'\1\1\1\1', b'\2\1\0\1\10\10', + b'\2\1\0\2\10\0', b'\1\1\1' + vint(2 * ((1 << 31) - 1) + 1)): + data = replace_payload(fixture('indexWithBuckets'), 0, 3, payload) + with pytest.raises(ValueError): + select(data, golden_meta(), None, bucket_filter=lambda b, t: False) + + +def test_varint_cursor_boundaries_and_single_byte_fast_path(): + values = {0, 1, 127, 128, manifest_sidecar.MAX_INT, manifest_sidecar.MAX_ROW_ID} + for bits in range(7, 63, 7): + values.update((2 ** bits - 1, 2 ** bits, 2 ** bits + 1)) + rng = random.Random(9908) + values.update(rng.getrandbits(63) for _ in range(1000)) + for value in sorted(values): + encoded = vint(value) + for source in (b'\xff' + encoded + b'\x55', bytearray(b'\xff' + encoded + b'\x55')): + reader = manifest_sidecar._Buffer(source, 1, 1 + len(encoded)) + assert reader.uint(value) == value + assert reader.position == reader.limit and reader.remaining == 0 + with pytest.raises(ValueError): + reader.take(1) + with pytest.raises(ValueError): + reader.uint() + reader = manifest_sidecar._Buffer(memoryview(b'prefix' + encoded + b'\x55')[6:]) + assert reader.uint() == value + assert bytes(reader.take(1)) == b'\x55' + + +@pytest.mark.parametrize('encoded', [b'', b'\x80', b'\x80\x00', b'\x81\x00', b'\xff\x00', + b'\x80' * 8, b'\xff' * 9, b'\xff' * 9 + b'\x01', + b'\xff' * 8 + b'\x00']) +def test_varint_rejects_truncation_noncanonical_and_overflow_without_reading_next_payload(encoded): + # The following byte could terminate a truncated varint, but belongs to a different payload. + reader = manifest_sidecar._Buffer(b'\x55' + encoded + b'\x01', 1, 1 + len(encoded)) + with pytest.raises(ValueError): + reader.uint() + assert reader.position <= reader.limit + + +@pytest.mark.parametrize('value,maximum', [ + (1, 0), (127, 126), (128, 127), + (manifest_sidecar.MAX_INT + 1, manifest_sidecar.MAX_INT), + (manifest_sidecar.MAX_ROW_ID, manifest_sidecar.MAX_INT)]) +def test_varint_fast_path_and_multibyte_path_both_enforce_maximum(value, maximum): + with pytest.raises(ValueError): + manifest_sidecar._Buffer(vint(value)).uint(maximum) + + +def test_fixed_long_and_take_respect_payload_window(): + for value in (-(1 << 63), -1, 0, manifest_sidecar.MAX_ROW_ID): + encoded = b'prefix' + struct.pack('>q', value) + b'suffix' + reader = manifest_sidecar._Buffer(encoded, 6, 14) + assert reader.long() == value and reader.remaining == 0 + assert bytes(reader.take(0)) == b'' + for size in (-1, 1): + with pytest.raises(ValueError): + reader.take(size) + with pytest.raises(ValueError): + reader.long() + # Underlying bytes contain a full long, but this payload is truncated. + with pytest.raises(ValueError): + manifest_sidecar._Buffer(encoded, 6, 13).long() + + +def test_delta_reader_is_not_created_for_rejected_or_single_interval_blocks(): + header = avro_header() + builder = Builder(Settings(), header) + builder.begin_block(len(header), 100, 1) + builder.add(100, 10, partition(7, 'left'), 1, 4) + builder.end_block() + data = builder.serialize(len(header) + 100, 1) + metadata = meta('m', len(header) + 100, 1) + with patch.object(manifest_sidecar._Deltas, '__init__', side_effect=AssertionError('Unneeded decoder')): + assert not select(data, metadata, [Range(999, 999)], part(7), FIELDS).blocks + assert len(select(data, metadata, [Range(105, 105)]).blocks) == 1 + # Even a guaranteed row miss must validate the other payloads' framing and counts. + for dimension in (1, 3): + with pytest.raises(ValueError): + select(replace_payload(data, 0, dimension, b'\0'), metadata, [Range(999, 999)]) + with pytest.raises(ValueError): + select(replace_payload(data, 0, 2, row_payload(200, 100, [])), metadata, [Range(999, 999)]) + + +def test_payload_cursor_is_shared_but_cannot_cross_its_declared_limit(): + stream = manifest_sidecar._Buffer(b'\1\1\x80\1\1\0\0') + cursor = manifest_sidecar._Buffer(stream.data) + first = manifest_sidecar._payload(stream, cursor) + assert first is cursor and first.data is stream.data + with pytest.raises(ValueError): + first.uint() + second = manifest_sidecar._payload(stream, cursor) + assert second is cursor and second.uint() == 0 + assert manifest_sidecar._payload(stream, cursor) is None + assert stream.remaining == 0 diff --git a/paimon-python/pypaimon/tests/manifest/manifest_sidecar_test.py b/paimon-python/pypaimon/tests/manifest/manifest_sidecar_test.py new file mode 100644 index 000000000000..405c4c92b14a --- /dev/null +++ b/paimon-python/pypaimon/tests/manifest/manifest_sidecar_test.py @@ -0,0 +1,801 @@ +# 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 +import unittest +from concurrent.futures import CancelledError +from copy import deepcopy +from io import BytesIO +from itertools import product + +import fastavro +from pyarrow import ArrowCancelled +from dataclasses import replace +from pathlib import Path +from tempfile import TemporaryDirectory +from types import SimpleNamespace +from unittest.mock import Mock, call, patch + +from pypaimon.common.options.core_options import CoreOptions +from pypaimon.common.options.options import Options +from pypaimon.filesystem.caching_file_io import CachingFileIO +from pypaimon.globalindex.global_index_result import GlobalIndexResult +from pypaimon.manifest import manifest_sidecar +from pypaimon.manifest.manifest_sidecar import ( + Block, Builder, Selection, Settings, SUFFIX, MAX_ROW_ID, Query, select, read_sidecar, + read_selected_bytes, sidecar_file_name, +) +from pypaimon.manifest.schema.manifest_entry import ManifestEntry +from pypaimon.manifest.manifest_list_manager import ManifestListManager +from pypaimon.manifest.manifest_file_manager import ManifestFileManager +from pypaimon.manifest.schema.manifest_file_meta import MANIFEST_FILE_META_SCHEMA +from pypaimon.read.scanner.file_scanner import FileScanner +from pypaimon.read.scan_stats import ScanStats +from pypaimon.tests.manifest import manifest_entry_identifier_test as existing +from pypaimon.schema.schema import Schema +from pypaimon.table.row.generic_row import GenericRow +from pypaimon.utils.range import Range + + +def golden(): + return make_sidecar() + + +def avro_header(): + return b'Obj\x01\x04\x14avro.codec\x08null\x16avro.schema\x0c"long"\x00' + bytes(16) + + +def make_sidecar(partitions=None, buckets=False): + header = avro_header() + builder = Builder(Settings(), header) + blocks = [(0, 100, [(0, 10, 0, 1, 4), (5, 5, 0, 1, 4), (20, 5, 1, 1, 8)]), + (100, 200, [((1 << 32) - 2, 5, 1, 2, 4), (8254058425445, 1, 0, 2, 8)]), + (300, 100, [(20, 5, 0, 0, 1), (MAX_ROW_ID, 1, 1, 3, 4)])] + for offset, length, entries in blocks: + builder.begin_block(len(header) + offset, length, len(entries)) + for first, count, p, bucket, total in entries: + builder.add(first, count, partitions[p] if partitions is not None else None, + bucket if buckets else None, total if buckets else None) + builder.end_block() + return builder.serialize(len(header) + 400, 7) + + +def golden_meta(): + return SimpleNamespace(file_name='manifest-golden', file_size=len(avro_header()) + 400, + num_added_files=7, num_deleted_files=0) + + +def intersects(data, meta, ranges, settings): + return bool(select(data, meta, ranges).blocks) + + +class CountingInput(BytesIO): + + def __init__(self, data, max_read=None): + super().__init__(data) + self.max_read = max_read + self.reads = [] + self.requests = [] + self.seeks = [] + + def read(self, size=-1): + if size < 0: + raise AssertionError('Unbounded read') + self.requests.append(size) + position = self.tell() + data = super().read(size if self.max_read is None else min(size, self.max_read)) + if data: + self.reads.append((position, len(data))) + return data + + def seek(self, offset, whence=0): + self.seeks.append(offset) + return super().seek(offset, whence) + + +class FailingIndexInput(BytesIO): + + def __init__(self, data, failure, phase, close_failure=None): + super().__init__(data) + self.failure = failure + self.phase = phase + self.close_failure = close_failure + + def read(self, size=-1): + if self.phase == 'read': + raise self.failure + return super().read(size) + + def close(self): + was_closed = self.closed + super().close() + if not was_closed and self.close_failure is not None: + raise self.close_failure + if self.phase == 'close' and not was_closed: + raise self.failure + + +class ManifestSidecarReadTest(unittest.TestCase): + + def test_local_cache_shares_sidecar_bytes_across_queries_and_readers(self): + data, meta = golden(), golden_meta() + meta.extra_files = ['custom' + SUFFIX] + for disk in (False, True): + with self.subTest(disk=disk), TemporaryDirectory() as directory: + options = Options({'local-cache.enabled': True, 'local-cache.max-size': '1 mb', + 'local-cache.block-size': '128 bytes'}) + if disk: + options.set(CoreOptions.LOCAL_CACHE_DIR, directory) + cache = CachingFileIO.create_cache_manager(options) + delegate = SimpleNamespace(new_input_stream=Mock(side_effect=lambda path: BytesIO(data)), + get_file_size=Mock(return_value=len(data))) + for parent in ('/table-a', '/table-b'): + path = parent + '/' + meta.file_name + for point, expected in ((20, [0, 5]), (0, [0]), (16, [])): + file_io = CachingFileIO.wrap_with_caching_if_needed(delegate, options, cache) + selected = read_sidecar(file_io, path, meta, [Range(point, point)]) + self.assertEqual([b.first_record for b in selected.blocks], expected) + expected_calls = [call('/table-a/custom' + SUFFIX), call('/table-b/custom' + SUFFIX)] + self.assertEqual(delegate.new_input_stream.call_args_list, expected_calls) + self.assertEqual(delegate.get_file_size.call_args_list, expected_calls) + + def test_local_cache_respects_disable_whitelist_and_byte_budget(self): + data, meta = golden(), golden_meta() + meta.extra_files = ['custom' + SUFFIX] + for overrides in ({'local-cache.enabled': False}, {'local-cache.whitelist': 'global-index'}, + {'local-cache.max-size': '128 bytes'}): + with self.subTest(overrides=overrides): + options = Options({'local-cache.enabled': True, 'local-cache.max-size': '1 mb', + 'local-cache.block-size': '128 bytes', **overrides}) + cache = CachingFileIO.create_cache_manager(options) + delegate = SimpleNamespace(new_input_stream=Mock(side_effect=lambda path: BytesIO(data)), + get_file_size=Mock(return_value=len(data))) + file_io = CachingFileIO.wrap_with_caching_if_needed(delegate, options, cache) + for _ in range(2): + selected = read_sidecar(file_io, '/table/' + meta.file_name, meta, + [Range(20, 20)]) + self.assertEqual([b.first_record for b in selected.blocks], [0, 5]) + self.assertEqual(delegate.new_input_stream.call_count, 2) + + def test_index_reads_use_bounded_bulk_requests(self): + header = avro_header() + for block_count in (5000, 25000, 131073): + with self.subTest(block_count=block_count): + builder = Builder(Settings(), header) + for block_number in range(block_count): + builder.begin_block(len(header) + block_number * 100, 100, 1) + builder.add(block_number, 1) + builder.end_block() + size = len(header) + block_count * 100 + data = builder.serialize(size, block_count) + meta = SimpleNamespace(file_name='manifest-large', file_size=size, + num_added_files=block_count, num_deleted_files=0, + extra_files=['manifest-large' + SUFFIX]) + stream = CountingInput(data) + file_io = SimpleNamespace(new_input_stream=lambda path: stream) + actual = read_sidecar(file_io, '/manifest/manifest-large', meta, + [Range(0, 0)]) + self.assertEqual(actual, select(data, meta, [Range(0, 0)])) + self.assertEqual(len(stream.reads), (len(data) + (1 << 20) - 1) // (1 << 20)) + self.assertLessEqual(max(stream.requests), 1 << 20) + self.assertTrue(stream.closed) + + def test_index_short_reads_and_exact_budget(self): + data, meta = golden(), golden_meta() + meta.extra_files = [meta.file_name + SUFFIX] + for max_read in (None, 7): + with self.subTest(max_read=max_read): + stream = CountingInput(data, max_read) + file_io = SimpleNamespace(new_input_stream=lambda path: stream) + actual = read_sidecar(file_io, '/manifest/manifest-golden', meta, + [Range(20, 20)]) + self.assertEqual(actual, select(data, meta, [Range(20, 20)])) + self.assertTrue(stream.closed) + + def test_adjacent_blocks_share_reads_without_reading_gaps(self): + header = avro_header() + body = bytes(range(200)) * 2 + for points, spans in [([0, 8254058425445], [(0, 300)]), + ([20], [(0, 100), (300, 100)]), ([16], [])]: + with self.subTest(points=points): + selected = select(golden(), golden_meta(), + [Range(point, point) for point in points]) + stream = CountingInput(header + body) + file_io = SimpleNamespace(new_input_stream=lambda path: stream) + actual = read_selected_bytes(file_io, '/manifest/manifest-golden', selected) + expected = header + b''.join(body[start:start + size] for start, size in spans) + self.assertEqual(actual, expected) + self.assertEqual(stream.reads, [(len(header) + start, size) for start, size in spans]) + self.assertEqual(stream.seeks, [len(header) + start for start, _ in spans]) + self.assertTrue(stream.closed) + + def test_large_block_spans_use_bounded_reads(self): + header = avro_header() + block_size = 512 * 1024 + body = bytes(2 * block_size + 257) + selected = Selection(header, (Block(len(header), block_size, 0, 1), + Block(len(header) + block_size, block_size, 1, 1), + Block(len(header) + 2 * block_size, 257, 2, 1))) + stream = CountingInput(header + body) + file_io = SimpleNamespace(new_input_stream=lambda path: stream) + self.assertEqual(read_selected_bytes(file_io, '/manifest/manifest-large', selected), header + body) + self.assertEqual(stream.reads, [(len(header), 1 << 20), (len(header) + (1 << 20), 257)]) + self.assertEqual(stream.seeks, [len(header)]) + self.assertTrue(stream.closed) + + def test_block_short_reads_and_truncation(self): + header = avro_header() + body = bytes(range(200)) * 2 + selected = select(golden(), golden_meta(), [Range(0, MAX_ROW_ID)]) + stream = CountingInput(header + body, 7) + file_io = SimpleNamespace(new_input_stream=lambda path: stream) + self.assertEqual(read_selected_bytes(file_io, '/manifest/manifest-golden', selected), header + body) + self.assertTrue(stream.closed) + stream = CountingInput(header + body[:-1], 7) + with self.assertRaises(EOFError): + read_selected_bytes(file_io, '/manifest/manifest-golden', selected) + self.assertTrue(stream.closed) + + +class ManifestSidecarFormatTest(unittest.TestCase): + + def test_settings_default_to_manifest_sort_with_explicit_override(self): + self.assertIsNone(CoreOptions.MANIFEST_SIDECAR_ENABLED.default_value()) + for sort, enabled in product((None, False, True), repeat=2): + with self.subTest(sort=sort, enabled=enabled): + values = {'manifest-sort.enabled': sort, 'manifest.sidecar.enabled': enabled} + options = CoreOptions(Options({key: value for key, value in values.items() if value is not None})) + settings = Settings.from_options(options) + self.assertEqual(settings.enabled, bool(sort) if enabled is None else enabled) + + def test_large_avro_headers(self): + stream = BytesIO() + fastavro.writer(stream, 'long', [42], metadata={'large-metadata': 'x' * (1 << 20)}) + avro_bytes = stream.getvalue() + block = next(fastavro.block_reader(BytesIO(avro_bytes))) + header = avro_bytes[:block.offset] + self.assertGreater(len(header), 1 << 20) + builder = Builder(Settings(), header) + builder.begin_block(block.offset, block.size, block.num_records) + builder.add(42, 1) + builder.end_block() + meta = SimpleNamespace(file_name='large-header.avro', file_size=len(avro_bytes), + num_added_files=1, num_deleted_files=0) + data = builder.serialize(meta.file_size, 1) + self.assertIsNotNone(data) + selected = select(data, meta, [Range(42, 42)]) + self.assertEqual(selected.header, header) + file_io = SimpleNamespace(new_input_stream=lambda path: BytesIO(avro_bytes)) + restored = read_selected_bytes(file_io, meta.file_name, selected) + self.assertEqual(restored, avro_bytes) + self.assertEqual(list(fastavro.reader(BytesIO(restored))), [42]) + + def test_disabled_option_is_independent_of_manifest_target_size(self): + for target in ('1 bytes', '1 gb'): + with self.subTest(target=target): + settings = Settings.from_options(CoreOptions(Options({ + 'manifest.target-file-size': target, 'manifest-sort.enabled': True, + 'manifest.sidecar.enabled': False}))) + self.assertFalse(settings.enabled) + + def test_row_id_coverage_and_block_ordinals(self): + data, meta, header = golden(), golden_meta(), avro_header() + for point in (0, 9, 20, 24, (1 << 32) - 2, 1 << 32, (1 << 32) + 2, + 8254058425445, MAX_ROW_ID): + self.assertTrue(intersects(data, meta, [Range(point, point)], Settings())) + for point in (10, 19, 25, (1 << 32) - 3, (1 << 32) + 3, 8254058425444, MAX_ROW_ID - 1): + self.assertFalse(intersects(data, meta, [Range(point, point)], Settings())) + selected = select(data, meta, [Range(20, 20)]) + self.assertEqual([b.first_record for b in selected.blocks], [0, 5]) + self.assertEqual([b.offset for b in selected.blocks], [len(header), len(header) + 300]) + self.assertEqual([b.length for b in selected.blocks], [100, 100]) + + gap = select(data, meta, [Range(16, 16)]) + + self.assertFalse(gap.blocks) + ranges = [Range(10, 19), Range(25, 40)] + self.assertFalse(intersects(data, meta, ranges, Settings())) + self.assertEqual(ranges, [Range(10, 19), Range(25, 40)]) + b = Builder(Settings(), header) + for offset, length, values in [ + (len(header), 100, [(0, 10), (5, 5), (20, 5)]), + (len(header) + 100, 200, [((1 << 32) - 2, 5), (8254058425445, 1)]), + (len(header) + 300, 100, [(20, 5), (MAX_ROW_ID, 1)])]: + b.begin_block(offset, length, len(values)) + for first, count in values: + b.add(first, count) + b.end_block() + self.assertEqual(b.serialize(meta.file_size, 7), + golden()) + + def test_minmax_skips_exact_checks_and_one_interval_is_already_exact(self): + header = avro_header() + builder = Builder(Settings(), header) + for offset, values in [(0, [(0, 10), (20, 10)]), + (100, [(100, 10), (200, 10)]), + (200, [(1 << 32, 10)])]: + builder.begin_block(len(header) + offset, 100, len(values)) + for first, count in values: + builder.add(first, count) + builder.end_block() + data = builder.serialize(len(header) + 300, 5) + meta = SimpleNamespace(file_name='m', file_size=len(header) + 300, + num_added_files=5, num_deleted_files=0) + for point, expected in [(50, 0), ((1 << 32) + 9, 1)]: + query = Query([Range(point, point)]) + with patch.object(query, 'intersects', wraps=query.intersects) as check: + selected = select(data, meta, query) + + self.assertEqual(len(selected.blocks), expected) + self.assertEqual(check.call_count, 3) + check.assert_any_call(0, 29) + check.assert_any_call(100, 209) + check.assert_any_call(1 << 32, (1 << 32) + 9) + + def test_single_interval_needs_no_delta_decoding(self): + header = avro_header() + for first, count in [(0, 1), (42, 10), (MAX_ROW_ID, 1)]: + builder = Builder(Settings(), header) + builder.begin_block(len(header), 100, 1) + builder.add(first, count) + builder.end_block() + data = builder.serialize(len(header) + 100, 1) + meta = SimpleNamespace(file_name='m', file_size=len(header) + 100, + num_added_files=1, num_deleted_files=0) + last = first + count - 1 + missing = first - 1 if first > 0 else last + 1 + for ranges, expected in [(None, 1), ([], 0), ([Range(first, first)], 1), + ([Range(last, last)], 1), ([Range(missing, missing)], 0)]: + with patch.object(manifest_sidecar._Deltas, 'next', side_effect=AssertionError('No deltas')) as decode: + self.assertEqual(len(select(data, meta, ranges).blocks), expected) + decode.assert_not_called() + + def test_consumed_intervals_still_require_valid_contents(self): + from pypaimon.tests.manifest.manifest_block_index_test import replace_payload, row_payload + for minimum, maximum, deltas in [(-1, 24, [10, 11]), (30, 24, [0, 1]), (0, 24, [9, 0]), + (0, 24, [9, 100])]: + data = replace_payload(golden(), 0, 2, row_payload(minimum, maximum, deltas)) + with self.assertRaises(ValueError): + select(data, golden_meta(), [Range(15, 15)]) + + def test_row_bounds_and_matches_skip_unused_intervals(self): + from pypaimon.tests.manifest.manifest_block_index_test import replace_payload, row_payload + data = replace_payload(golden(), 0, 2, row_payload(0, 24, [9, 0])) + self.assertEqual(len(select(data, golden_meta(), [Range(0, 0)]).blocks), 1) + self.assertFalse(select(data, golden_meta(), [Range(100, 100)]).blocks) + self.assertFalse(select(data, golden_meta(), []).blocks) + with self.assertRaises(ValueError): + select(data, golden_meta(), [Range(15, 15)]) + + def test_unknown_coverage_and_exact_gaps(self): + header = avro_header() + for first, count in [(None, 1), (-1, 1), (10, 0), (10, -1), (MAX_ROW_ID, 2)]: + b = Builder(Settings(), header) + b.begin_block(len(header), 100, 1) + b.add(first, count) + b.end_block() + data = b.serialize(len(header) + 100, 1) + meta = SimpleNamespace(file_name='m', file_size=len(header) + 100, + num_added_files=1, num_deleted_files=0) + self.assertEqual(len(select(data, meta, [Range(100, 100)]).blocks), 1) + b = Builder(Settings(), header) + b.begin_block(len(header), 100, 2) + b.add(0, MAX_ROW_ID) + b.add(MAX_ROW_ID, 1) + b.end_block() + self.assertLess(len(b.serialize(len(header) + 100, 2)), 512) + b = Builder(Settings(), header) + b.begin_block(len(header), 100, 64) + for i in range(64): + b.add(1 + i * 100, 1) + b.end_block() + data = b.serialize(len(header) + 100, 64) + meta = SimpleNamespace(file_name='m', file_size=len(header) + 100, + num_added_files=64, num_deleted_files=0) + self.assertEqual(len(select(data, meta, [Range(10, 10)]).blocks), 0) + + def test_invalid_envelopes(self): + from pypaimon.tests.manifest.manifest_block_index_test import checksum + meta, data = golden_meta(), golden() + for index in (0, 9, 11, 15, 16, 55, 63, 67, 75, len(data) - 1): + bad = bytearray(data) + bad[index] ^= 2 + with self.assertRaises(ValueError): + select(bad, meta, [Range(10, 10)]) + for version in (0, 2, 99): + bad = bytearray(data) + bad[4] = version + with self.assertRaises(ValueError): + select(checksum(bad), meta, [Range(10, 10)]) + for length in range(len(data)): + with self.assertRaises(ValueError): + select(data[:length], meta, [Range(10, 10)]) + meta.file_name = 'renamed' + self.assertEqual(len(select(data, meta, [Range(20, 20)]).blocks), 2) + meta.file_size += 1 + with self.assertRaises(ValueError): + select(data, meta, [Range(10, 10)]) + + +class ManifestSidecarScanTest(existing.ManifestEntryIdentifierTest): + + def setUp(self): + super().setUp() + self.table.options.options.set(CoreOptions.MANIFEST_SORT_ENABLED, True) + self.table.options.options.set(CoreOptions.DATA_EVOLUTION_ENABLED, True) + + def entry(self, name, first, count=10, kind=0): + return ManifestEntry(kind, self._create_file_meta('unused').min_key, 0, 1, + replace(self._create_file_meta(name), first_row_id=first, row_count=count)) + + def test_enabling_sidecar_reads_does_not_change_python_writes(self): + manager = self.manifest_file_manager + entries = [self.entry('data.parquet', 100)] + self.assertTrue(self.table.options.manifest_sidecar_enabled()) + self.assertIsNone(manager.write('plain-writer', entries)) + self.assertFalse(Path(manager.manifest_path, 'plain-writer' + SUFFIX).exists()) + for meta in manager.rolling_write(entries, 1, 'rolling-writer'): + self.assertIsNone(meta.extra_files) + self.assertFalse(Path(manager.manifest_path, meta.file_name + SUFFIX).exists()) + + def write_meta(self, name, entries): + # Emulate an externally published sidecar; production Python writes are unchanged. + manager = self.manifest_file_manager + manager.write(name, entries) + path = Path(manager.manifest_path, name) + avro_bytes = path.read_bytes() + meta = manager._build_meta(name, entries, len(avro_bytes)) + settings = Settings.from_options(self.table.options) + if settings.enabled: + data = manifest_sidecar.build_from_entries(avro_bytes, entries, settings) + sidecar_path = path.with_name(name + SUFFIX) + sidecar_path.write_bytes(data) + meta = replace(meta, extra_files=[sidecar_path.name]) + return meta + + def test_bucket_point_lookup_with_rescale_and_delete_entries(self): + self.table.options.options.set(CoreOptions.DATA_EVOLUTION_ENABLED, False) + self.table.options.options.set(CoreOptions.BUCKET, 4) + from pypaimon.common.predicate import Predicate + from pypaimon.read.scanner.bucket_select_converter import create_bucket_selector + selector = create_bucket_selector(Predicate('equal', 0, 'id', [7]), self.table.fields[:1]) + self.assertIsNotNone(selector) + entries = [replace(self.entry('%s-%s-%s.parquet' % (total, bucket, i), None), + bucket=bucket, total_buckets=total) + for total in (4, 8) for bucket in range(total) for i in range(300)] + metadata = self.write_meta('buckets', entries) + results = [] + for enabled in (False, True): + self.table.options.options.set(CoreOptions.MANIFEST_SIDECAR_ENABLED, enabled) + scanner = FileScanner(self.table, lambda: ([metadata], None)) + scanner._bucket_selector = selector + with patch('pypaimon.manifest.manifest_file_manager.read_selected_bytes', + wraps=read_selected_bytes) as read_blocks: + actual = scanner.read_manifest_entries([metadata]) + results.append([e.file.file_name for e in actual]) + if enabled: + read_blocks.assert_called_once() + selected = read_blocks.call_args[0][2] + self.assertLess(sum(b.record_count for b in selected.blocks), 1000) + else: + read_blocks.assert_not_called() + self.assertEqual(results[0], results[1]) + self.assertEqual(len(results[1]), 600) + chosen = next(bucket for bucket in range(4) if selector(bucket, 4)) + added = [replace(self.entry('point.' + suffix, None), bucket=chosen, total_buckets=4) + for suffix in ('parquet', 'blob')] + metas = [self.write_meta('point-add', added), + self.write_meta('point-delete', [replace(e, kind=1) for e in added])] + self.assertEqual(scanner.read_manifest_entries(metas), []) + + def test_partition_only_and_conjunctive_planning_keep_entry_and_delete_filters(self): + import pyarrow as pa + schema = Schema.from_pyarrow_schema( + pa.schema([('p', pa.int32()), ('q', pa.string()), ('value', pa.string())]), + partition_keys=['p', 'q'], + options={'manifest.sidecar.enabled': 'true'}) + self.catalog.create_table('default.partition_block_index', schema, False) + self.table = self.catalog.get_table('default.partition_block_index') + self.manifest_file_manager = ManifestFileManager(self.table) + fields = self.table.partition_keys_fields + + def entry(name, first, p, kind=0): + return replace(self.entry(name, first, kind=kind), partition=GenericRow([p, None], fields)) + + # Build many blocks without any row tracking; partition-only planning must use the index. + entries = [entry('part-%d.parquet' % i, None, i // 1000) for i in range(4000)] + manifest = self.write_meta('partitioned', entries) + from pypaimon.common.predicate import Predicate + predicate = Predicate('equal', 0, 'p', [1]) + results = [] + for enabled in (False, True): + self.table.options.options.set(CoreOptions.MANIFEST_SIDECAR_ENABLED, enabled) + scanner = FileScanner(self.table, lambda: ([manifest], None), partition_predicate=predicate) + with patch('pypaimon.manifest.manifest_file_manager.read_selected_bytes', + wraps=read_selected_bytes) as read_blocks: + actual = scanner.read_manifest_entries([manifest]) + results.append([e.file.file_name for e in actual]) + if enabled: + read_blocks.assert_called_once() + selected = read_blocks.call_args[0][2] + self.assertLess(sum(b.record_count for b in selected.blocks), 1500) + else: + read_blocks.assert_not_called() + self.assertEqual(results[0], results[1]) + self.assertEqual(len(results[1]), 1000) + + # A block can match the two dimensions through different entries; keep entry filtering. + mixed = self.write_meta('mixed', [entry('a.parquet', 100, 1), entry('b.parquet', 5, 2)]) + scanner = FileScanner(self.table, lambda: ([mixed], None), partition_predicate=predicate) + self.assertEqual(scanner.read_manifest_entries([mixed], row_ranges=[Range(5, 5)]), []) + # Both column groups and DELETE blocks must survive the same partition + row-id filter. + add = [entry('data.parquet', 100, 1), entry('data.blob', 100, 1)] + metas = [self.write_meta('adds', add), self.write_meta('deletes', [replace(e, kind=1) for e in add])] + self.assertEqual(scanner.read_manifest_entries(metas, row_ranges=[Range(105, 105)]), []) + + def test_explain_keeps_complete_entry_counts_with_sidecar_enabled(self): + import pyarrow as pa + from pypaimon.common.predicate import Predicate + + schema = Schema.from_pyarrow_schema( + pa.schema([('p', pa.int32()), ('q', pa.string()), ('value', pa.string())]), + partition_keys=['p', 'q'], + options={'manifest.sidecar.enabled': 'true', + 'data-evolution.enabled': 'true'}) + self.catalog.create_table('default.explain_sidecar', schema, False) + self.table = self.catalog.get_table('default.explain_sidecar') + self.manifest_file_manager = ManifestFileManager(self.table) + fields = self.table.partition_keys_fields + entries = [replace(self.entry('part-%d.parquet' % i, i * 1000), + partition=GenericRow([i // 1000, None], fields)) for i in range(4000)] + manifest = self.write_meta('explain-manifest', entries) + predicate = Predicate('equal', 0, 'p', [1]) + ranges = [Range(1000000, 1000000)] + for partition_filter, row_ranges in [(predicate, None), (None, ranges), (predicate, ranges)]: + results = [] + for enabled in (False, True): + with self.subTest(partition=partition_filter, row_ranges=row_ranges, sidecar=enabled): + self.table.options.options.set(CoreOptions.MANIFEST_SIDECAR_ENABLED, enabled) + scanner = FileScanner(self.table, lambda: ([manifest], None), + partition_predicate=partition_filter) + scanner.scan_stats = ScanStats() + with patch('pypaimon.manifest.manifest_file_manager.read_sidecar', + wraps=read_sidecar) as read_metadata: + actual = scanner.read_manifest_entries([manifest], row_ranges=row_ranges) + read_metadata.assert_not_called() + stats = scanner.scan_stats + self.assertEqual(stats.entries_potential_total, 4000) + self.assertEqual(stats.entries_total, 4000) + self.assertEqual(stats.entries_after_partition, 1000 if partition_filter else 4000) + self.assertEqual(stats.partition_keys_before, {(p, None) for p in range(4)}) + results.append([entry.file.file_name for entry in actual]) + self.assertEqual(results[0], results[1]) + + def test_explicit_reference_and_null_does_not_probe(self): + manager = self.manifest_file_manager + written = self.write_meta('explicit', [self.entry('data.parquet', 100)]) + self.assertEqual(sidecar_file_name(written), written.file_name + SUFFIX) + index_path = Path(manager.manifest_path, sidecar_file_name(written)) + explicit_path = index_path.with_name('independent-index' + SUFFIX) + index_path.rename(explicit_path) + other_path = index_path.with_name('other-partition-index') + other_path.write_bytes(b'not a manifest sidecar') + indexed = replace(written, extra_files=[other_path.name, explicit_path.name]) + with patch.object(self.table.file_io, 'new_input_stream', + wraps=self.table.file_io.new_input_stream) as opened: + self.assertEqual(manager.read_entries_parallel([indexed], row_ranges=[Range(0, 0)]), []) + self.assertEqual([call[0][0] for call in opened.call_args_list], [str(explicit_path)]) + + for extra_files in (None, [], [other_path.name]): + with self.subTest(extra_files=extra_files): + unindexed = replace(indexed, extra_files=extra_files) + with patch.object(self.table.file_io, 'new_input_stream', + wraps=self.table.file_io.new_input_stream) as opened: + actual = manager.read_entries_parallel([unindexed], row_ranges=[Range(0, 0)]) + self.assertEqual(len(actual), 1) + self.assertEqual([call[0][0] for call in opened.call_args_list], + [str(Path(manager.manifest_path, written.file_name))]) + + def test_manifest_list_index_reference_compatibility(self): + indexed = self.write_meta('indexed', [self.entry('data.parquet', 100)]) + indexed = replace(indexed, extra_files=['other-index'] + indexed.extra_files) + self.table.options.options.set(CoreOptions.MANIFEST_SIDECAR_ENABLED, False) + unindexed = self.write_meta('legacy-entry', [self.entry('old.parquet', None)]) + self.assertIsNone(sidecar_file_name(unindexed)) + lists = ManifestListManager(self.table) + lists.write('references', [indexed, unindexed]) + actual = lists.read('references') + self.assertEqual([meta.extra_files for meta in actual], [indexed.extra_files, None]) + self.assertEqual([sidecar_file_name(meta) for meta in actual], [sidecar_file_name(indexed), None]) + self.assertEqual([meta.file_name for meta in actual], [indexed.file_name, unindexed.file_name]) + + data = Path(lists.manifest_path, 'references').read_bytes() + legacy_schema = deepcopy(MANIFEST_FILE_META_SCHEMA) + legacy_schema['fields'] = [field for field in legacy_schema['fields'] + if field['name'] != '_EXTRA_FILES'] + legacy_records = list(fastavro.reader(BytesIO(data), reader_schema=legacy_schema)) + self.assertTrue(all('_EXTRA_FILES' not in record for record in legacy_records)) + self.assertEqual([record['_FILE_NAME'] for record in legacy_records], + [indexed.file_name, unindexed.file_name]) + with self.table.file_io.new_output_stream(str(Path(lists.manifest_path, 'old-list'))) as stream: + fastavro.writer(stream, legacy_schema, legacy_records) + self.assertTrue(all(sidecar_file_name(meta) is None for meta in lists.read('old-list'))) + + def test_skips_blocks_inside_a_matching_manifest(self): + entries = [self.entry('file-%d.parquet' % i, i * 1000) for i in range(4000)] + meta = self.write_meta('many-blocks', entries) + outputs = [] + reader = fastavro.reader + for enabled in (False, True): + self.table.options.options.set(CoreOptions.MANIFEST_SIDECAR_ENABLED, enabled) + scanner = FileScanner(self.table, lambda: ([meta], None)) + scanner.with_global_index_result(GlobalIndexResult.from_ranges([Range(2000005, 2000005)])) + decoded = [] + + def observed_reader(stream): + for record in reader(stream): + decoded.append(record) + yield record + + with patch('pypaimon.manifest.manifest_file_manager.fastavro.reader', side_effect=observed_reader), \ + patch('pypaimon.manifest.manifest_file_manager.read_selected_bytes', + wraps=read_selected_bytes) as selected_read: + actual, _ = scanner._create_data_evolution_split_generator() + outputs.append([e.file.file_name for e in actual]) + if enabled: + self.assertEqual(selected_read.call_count, 1) + selected = selected_read.call_args[0][2] + self.assertEqual(len(selected.blocks), 1) + self.assertLess(sum(block.length for block in selected.blocks), meta.file_size // 10) + self.assertLess(len(decoded), 200) + self.assertEqual(len(decoded), selected.blocks[0].record_count) + else: + self.assertEqual(selected_read.call_count, 0) + self.assertEqual(len(decoded), 4000) + self.assertEqual(outputs, [['file-2000.parquet']] * 2) + + def test_actual_global_index_scanner_72_to_2(self): + metas = [] + for i in range(72): + entries = [self.entry('a%d.parquet' % i, 0), self.entry('b%d.parquet' % i, 100)] + if i < 2: + entries.append(self.entry('hit.' + ('parquet' if i == 0 else 'blob'), 45)) + metas.append(self.write_meta('manifest-%d' % i, entries)) + results = [] + for enabled in (False, True): + self.table.options.options.set(CoreOptions.MANIFEST_SIDECAR_ENABLED, enabled) + scanner = FileScanner(self.table, lambda: (metas, None)) + scanner.with_global_index_result(GlobalIndexResult.from_ranges([Range(50, 50)])) + manager = scanner.manifest_file_manager + with patch.object(manager, 'read', wraps=manager.read) as read_body, \ + patch('pypaimon.manifest.manifest_file_manager.read_sidecar', wraps=read_sidecar) as read_metadata: + entries, _ = scanner._create_data_evolution_split_generator() + results.append(sorted(e.file.file_name for e in entries)) + self.assertEqual(len(read_body.call_args_list), 2 if enabled else 72) + self.assertEqual(len(read_metadata.call_args_list), 72 if enabled else 0) + self.assertEqual(results, [['hit.blob', 'hit.parquet']] * 2) + + def test_delete_union_no_resurrection_and_no_query_no_index_io(self): + add = self.entry('data.parquet', 45) + blob = self.entry('data.blob', 45) + metas = [self.write_meta('add', [add, blob]), + self.write_meta('delete', [replace(add, kind=1), replace(blob, kind=1)]), + self.write_meta('gap', [self.entry('lo', 0), self.entry('hi', 100)])] + manager = self.manifest_file_manager + for enabled in (False, True): + self.table.options.options.set(CoreOptions.MANIFEST_SIDECAR_ENABLED, enabled) + with patch.object(manager, 'read', wraps=manager.read) as read_body: + entries = manager.read_entries_parallel(metas[:2], row_ranges=[Range(50, 50)]) + self.assertEqual(entries, []) + self.assertEqual(len(read_body.call_args_list), 2) + with patch.object(self.table.file_io, 'new_input_stream', + wraps=self.table.file_io.new_input_stream) as opened: + manager.read_entries_parallel(metas) + self.assertTrue(all(not call[0][0].endswith(SUFFIX) for call in opened.call_args_list)) + # Missing and corrupt objects retain their manifests and read the full body. + path = manager.manifest_path + '/gap' + SUFFIX + for bad in (None, b'partial'): + if bad is None: + os.unlink(path) + else: + Path(path).write_bytes(bad) + with patch.object(manager, 'read', wraps=manager.read) as read_body: + entries = manager.read_entries_parallel(metas[2:], row_ranges=[Range(50, 50)]) + self.assertEqual(len(entries), 2) + self.assertEqual(read_body.call_count, 1) + self.assertIsNone(read_body.call_args[1]['selected_blocks']) + with patch.object(self.table.file_io, 'new_input_stream', side_effect=InterruptedError('stop')): + with self.assertRaises(InterruptedError): + read_sidecar(self.table.file_io, path, metas[0], [Range(0, 0)]) + + def test_sidecar_cancellation_during_open(self): + self._check_sidecar_cancellation('open') + + def test_sidecar_cancellation_during_read(self): + self._check_sidecar_cancellation('read') + + def test_sidecar_cancellation_during_close(self): + self._check_sidecar_cancellation('close') + + def _check_sidecar_cancellation(self, phase): + for failure_type in (ArrowCancelled, CancelledError, InterruptedError): + with self.subTest(phase=phase, failure_type=failure_type): + self._check_sidecar_io_failure(phase, failure_type('cancelled'), cancelled=True) + + def test_sidecar_io_failures_fall_back_to_manifest(self): + for phase in ('open', 'read', 'close'): + for failure_type in (FileNotFoundError, TimeoutError, OSError): + with self.subTest(phase=phase, failure_type=failure_type): + self._check_sidecar_io_failure(phase, failure_type('unavailable'), cancelled=False) + + def test_sidecar_cancellation_survives_close_failure(self): + for failure_type in (ArrowCancelled, CancelledError, InterruptedError): + with self.subTest(failure_type=failure_type): + self._check_sidecar_io_failure('read', failure_type('cancelled'), cancelled=True, + close_failure=OSError('close failed')) + + def test_sidecar_wrapped_cancellation_propagates(self): + for failure_type in (ArrowCancelled, CancelledError, InterruptedError): + with self.subTest(failure_type=failure_type): + cancellation = failure_type('cancelled') + wrapped = OSError('wrapped failure') + wrapped.__cause__ = cancellation + self._check_sidecar_io_failure('open', wrapped, cancelled=True, + expected_failure=cancellation) + + def test_sidecar_exception_cycle_falls_back(self): + first = OSError('first') + second = OSError('second') + first.__cause__ = second + second.__cause__ = first + self._check_sidecar_io_failure('open', first, cancelled=False) + + def _check_sidecar_io_failure(self, phase, failure, cancelled, close_failure=None, + expected_failure=None): + manager = self.manifest_file_manager + meta = self.write_meta('failure-' + phase + '-' + type(failure).__name__, + [self.entry('data.parquet', 100)]) + index_path = str(Path(manager.manifest_path, sidecar_file_name(meta))) + body_path = str(Path(manager.manifest_path, meta.file_name)) + stream = (FailingIndexInput(Path(index_path).read_bytes(), failure, phase, close_failure) + if phase != 'open' else None) + original_open = self.table.file_io.new_input_stream + + def open_stream(path): + if path == index_path: + if phase == 'open': + raise failure + return stream + return original_open(path) + + with patch.object(self.table.file_io, 'new_input_stream', side_effect=open_stream) as opened, \ + patch.object(manager, 'read', wraps=manager.read) as read_body: + if cancelled: + expected = failure if expected_failure is None else expected_failure + with self.assertRaises(type(expected)) as raised: + manager.read_entries_parallel([meta], row_ranges=[Range(100, 100)]) + self.assertIs(raised.exception, expected) + read_body.assert_not_called() + self.assertEqual([call[0][0] for call in opened.call_args_list], [index_path]) + else: + entries = manager.read_entries_parallel([meta], row_ranges=[Range(100, 100)]) + self.assertEqual([entry.file.file_name for entry in entries], ['data.parquet']) + read_body.assert_called_once() + self.assertIsNone(read_body.call_args[1]['selected_blocks']) + self.assertEqual([call[0][0] for call in opened.call_args_list], [index_path, body_path]) + if stream is not None: + self.assertTrue(stream.closed) diff --git a/paimon-python/pypaimon/tests/reader_append_only_test.py b/paimon-python/pypaimon/tests/reader_append_only_test.py index 480cd4ebbe4b..1c783bfe0611 100644 --- a/paimon-python/pypaimon/tests/reader_append_only_test.py +++ b/paimon-python/pypaimon/tests/reader_append_only_test.py @@ -1090,7 +1090,8 @@ def test_is_in_with_partitions(self): def counting_read(self_mgr, manifest_file_name, manifest_entry_filter=None, drop_stats=True, early_entry_filter=None, - early_record_filter=None, partition_filter=None): + early_record_filter=None, partition_filter=None, + selected_blocks=None): # avro_total = every entry in the manifest (no manifest-file pruning # here: single file, is_in spans its partition stats). path = f"{self_mgr.manifest_path}/{manifest_file_name}" @@ -1100,7 +1101,8 @@ def counting_read(self_mgr, manifest_file_name, return original_read( self_mgr, manifest_file_name, manifest_entry_filter, drop_stats, - early_entry_filter, early_record_filter, partition_filter) + early_entry_filter, early_record_filter, partition_filter, + selected_blocks=selected_blocks) def counting_dfm_init(self_dfm, *args, **kwargs): entry_counts['constructed'] += 1 diff --git a/paimon-python/pypaimon/utils/file_type.py b/paimon-python/pypaimon/utils/file_type.py index 92819f3618bc..7dadc93af631 100644 --- a/paimon-python/pypaimon/utils/file_type.py +++ b/paimon-python/pypaimon/utils/file_type.py @@ -28,7 +28,7 @@ class FileType(Enum): """Classification of Paimon files. - - META: snapshot, schema, manifest, statistics, tag, changelog metadata, + - META: snapshot, schema, manifest, manifest sidecar, statistics, tag, changelog metadata, hint files, _SUCCESS, consumer, service files - DATA: data files and any unrecognized files (default) - BUCKET_INDEX: bucket level index files (Hash, DV) @@ -67,7 +67,7 @@ def classify(file_path: str) -> 'FileType': return FileType.GLOBAL_INDEX return FileType.FILE_INDEX - if "manifest" in name: + if "manifest" in name or name.endswith(".avro.sidecar"): return FileType.META if name.startswith("index-"):