| // 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. |
| |
| //! TableScan for full table scan. |
| //! |
| //! Reference: [pypaimon.read.table_scan.TableScan](https://github.com/apache/paimon/blob/release-1.3/paimon-python/pypaimon/read/table_scan.py) |
| //! and [FullStartingScanner](https://github.com/apache/paimon/blob/release-1.3/paimon-python/pypaimon/read/scanner/full_starting_scanner.py). |
| |
| use super::bucket_filter::compute_target_buckets; |
| use super::partition_filter::PartitionFilter; |
| use super::stats_filter::{ |
| data_evolution_group_matches_predicates, data_file_matches_predicates, |
| data_file_matches_predicates_for_table, group_by_overlapping_row_id, FileStatsRows, |
| ResolvedStatsSchema, |
| }; |
| use super::Table; |
| use crate::io::FileIO; |
| use crate::spec::{ |
| avro::SharedSchemaCache, bucket_dir_name, BinaryRow, BucketFunctionType, CoreOptions, |
| DataField, DataFileMeta, FileKind, IndexManifest, ManifestEntry, PartitionComputer, Predicate, |
| Snapshot, |
| }; |
| use crate::table::bin_pack::split_for_batch; |
| use crate::table::merge_tree_split_generator::{ |
| merge_tree_split_for_batch, KeyComparator, SplitGroup, |
| }; |
| use crate::table::source::{ |
| any_range_overlaps_file, intersect_ranges_with_file, merge_row_ranges, DataSplit, |
| DataSplitBuilder, DeletionFile, PartitionBucket, Plan, RowRange, |
| }; |
| use crate::table::SnapshotManager; |
| use futures::{StreamExt, TryStreamExt}; |
| use std::collections::{HashMap, HashSet}; |
| use std::sync::Arc; |
| |
| /// Path segment for manifest directory under table. |
| const MANIFEST_DIR: &str = "manifest"; |
| /// Path segment for index directory under table. |
| const INDEX_DIR: &str = "index"; |
| |
| /// Reads a manifest list file (Avro) and returns manifest file metas. |
| async fn read_manifest_list( |
| file_io: &FileIO, |
| table_path: &str, |
| list_name: &str, |
| ) -> crate::Result<Vec<crate::spec::ManifestFileMeta>> { |
| if list_name.is_empty() { |
| return Ok(Vec::new()); |
| } |
| let path = format!( |
| "{}/{}/{}", |
| table_path.trim_end_matches('/'), |
| MANIFEST_DIR, |
| list_name |
| ); |
| let input = file_io.new_input(&path)?; |
| let bytes = input.read().await?; |
| crate::spec::avro::from_avro_bytes_fast::<crate::spec::ManifestFileMeta>(&bytes) |
| } |
| |
| /// Reads all manifest entries for a snapshot (base + delta manifest lists, then each manifest file). |
| /// Applies filters during concurrent manifest reading to reduce entries early: |
| /// - Manifest-file-level partition stats pruning (skip entire manifest files) |
| /// - Level-0 filtering per entry (DV mode or FirstRow engine) |
| /// - Partition predicate filtering per entry |
| /// - Data-level stats pruning per entry (current schema only, cross-schema fail-open) |
| #[allow(clippy::too_many_arguments)] |
| async fn read_all_manifest_entries( |
| file_io: &FileIO, |
| table_path: &str, |
| snapshot: &Snapshot, |
| skip_level_zero: bool, |
| scan_all_files: bool, |
| has_primary_keys: bool, |
| partition_filter: Option<&PartitionFilter>, |
| partition_fields: &[DataField], |
| data_predicates: &[Predicate], |
| current_schema_id: i64, |
| schema_fields: &[DataField], |
| bucket_predicate: Option<&Predicate>, |
| bucket_key_fields: &[DataField], |
| bucket_function_type: BucketFunctionType, |
| ) -> crate::Result<Vec<ManifestEntry>> { |
| let (mut manifest_files, delta) = futures::try_join!( |
| read_manifest_list(file_io, table_path, snapshot.base_manifest_list()), |
| read_manifest_list(file_io, table_path, snapshot.delta_manifest_list()), |
| )?; |
| manifest_files.extend(delta); |
| |
| // Manifest-file-level partition stats pruning: skip entire manifest files |
| // whose partition range doesn't overlap the partition predicate. |
| if let Some(pf) = partition_filter { |
| if !partition_fields.is_empty() { |
| manifest_files.retain(|meta| { |
| let stats = meta.partition_stats(); |
| let min_values = BinaryRow::from_serialized_bytes(stats.min_values()).ok(); |
| let max_values = BinaryRow::from_serialized_bytes(stats.max_values()).ok(); |
| let null_counts = stats.null_counts().clone(); |
| let file_stats = FileStatsRows::for_manifest_partition( |
| meta.num_added_files() + meta.num_deleted_files(), |
| min_values, |
| max_values, |
| null_counts, |
| ); |
| pf.matches_manifest(&file_stats, partition_fields) |
| }); |
| } |
| } |
| |
| let manifest_path_prefix = format!("{}/{}", table_path.trim_end_matches('/'), MANIFEST_DIR); |
| let shared_cache = SharedSchemaCache::new(); |
| let all_entries: Vec<ManifestEntry> = futures::stream::iter(manifest_files) |
| .map(|meta| { |
| let path = format!("{}/{}", manifest_path_prefix, meta.file_name()); |
| let cache = shared_cache.clone(); |
| async move { |
| let input_file = file_io.new_input(&path)?; |
| let content = input_file.read().await?; |
| |
| // Per-task bucket cache (few distinct total_buckets values per manifest). |
| let mut bucket_cache: HashMap<i32, Option<HashSet<i32>>> = HashMap::new(); |
| |
| let entries = crate::spec::avro::from_manifest_bytes_filtered_shared( |
| &content, |
| &cache, |
| &mut |_kind, partition_bytes, bucket, total_buckets| { |
| // Bucket filter (negative bucket = unassigned) |
| if has_primary_keys && !scan_all_files && bucket < 0 { |
| return false; |
| } |
| if let Some(pred) = bucket_predicate { |
| let targets = bucket_cache.entry(total_buckets).or_insert_with(|| { |
| compute_target_buckets( |
| pred, |
| bucket_key_fields, |
| bucket_function_type, |
| total_buckets, |
| ) |
| }); |
| if let Some(targets) = targets { |
| if !targets.contains(&bucket) { |
| return false; |
| } |
| } |
| } |
| |
| // Partition filter |
| if let Some(pf) = partition_filter { |
| match pf.matches_entry(partition_bytes) { |
| Ok(false) => return false, |
| Ok(true) => {} |
| Err(_) => {} |
| } |
| } |
| |
| true |
| }, |
| )?; |
| |
| // Post-filter: level-0 and data predicates (need DataFileMeta) |
| let filtered: Vec<ManifestEntry> = entries |
| .into_iter() |
| .filter(|entry| { |
| if skip_level_zero && has_primary_keys && entry.file().level == 0 { |
| return false; |
| } |
| if !data_predicates.is_empty() |
| && !data_file_matches_predicates( |
| entry.file(), |
| data_predicates, |
| current_schema_id, |
| schema_fields, |
| ) |
| { |
| return false; |
| } |
| true |
| }) |
| .collect(); |
| Ok::<_, crate::Error>(filtered) |
| } |
| }) |
| .buffered(64) |
| .try_collect::<Vec<_>>() |
| .await? |
| .into_iter() |
| .flatten() |
| .collect(); |
| Ok(all_entries) |
| } |
| |
| /// Builds a map from (partition, bucket) to (data_file_name -> DeletionFile) from index manifest entries. |
| /// Only considers ADD entries with index_type "DELETION_VECTORS" and their deletion_vectors_ranges. |
| fn build_deletion_files_map( |
| index_entries: &[crate::spec::IndexManifestEntry], |
| table_path: &str, |
| ) -> HashMap<PartitionBucket, HashMap<String, DeletionFile>> { |
| use crate::spec::FileKind; |
| let table_path = table_path.trim_end_matches('/'); |
| let index_path_prefix = format!("{table_path}/{INDEX_DIR}"); |
| let mut map: HashMap<PartitionBucket, HashMap<String, DeletionFile>> = |
| HashMap::with_capacity(index_entries.len()); |
| for entry in index_entries { |
| if entry.kind != FileKind::Add { |
| continue; |
| } |
| if entry.index_file.index_type != "DELETION_VECTORS" { |
| continue; |
| } |
| let ranges = match &entry.index_file.deletion_vectors_ranges { |
| Some(r) if !r.is_empty() => r, |
| _ => continue, |
| }; |
| let key = PartitionBucket::new(entry.partition.clone(), entry.bucket); |
| let dv_path = format!("{}/{}", index_path_prefix, entry.index_file.file_name); |
| let per_bucket = map.entry(key).or_default(); |
| for (data_file_name, meta) in ranges { |
| per_bucket.insert( |
| data_file_name.clone(), |
| DeletionFile::new( |
| dv_path.clone(), |
| meta.offset as i64, |
| meta.length as i64, |
| meta.cardinality, |
| ), |
| ); |
| } |
| } |
| map |
| } |
| |
| /// Nets add/delete manifest entries for a scan, returning only the live ADD set. |
| /// |
| /// Mirrors Java `AbstractFileStoreScan.readAndMergeFileEntries`: first collect |
| /// the full [`Identifier`] of every DELETE entry, then keep the ADD entries |
| /// whose identifier is not in that set. The identity is the complete Paimon file |
| /// identity (`partition, bucket, level, file_name, extra_files, embedded_index, |
| /// external_path`, matching Java `FileEntry.Identifier`). |
| /// |
| /// Keying on file name alone is wrong: a single-run compaction upgrades a file |
| /// *in place* — `DELETE f@oldLevel` plus `ADD f@newLevel` with the same file |
| /// name (`PojoDataFileMeta.upgrade` reuses the name, only changing `level`). An |
| /// identity without `level` lets the DELETE cancel the upgraded ADD, dropping |
| /// the file from the scan and silently losing its rows on read. |
| /// |
| /// Collecting deletes first (rather than insert/remove while iterating) makes |
| /// the result independent of ADD/DELETE ordering, matching the Java scan path. |
| fn merge_manifest_entries(entries: Vec<ManifestEntry>) -> Vec<ManifestEntry> { |
| use crate::spec::Identifier; |
| let deleted: HashSet<Identifier> = entries |
| .iter() |
| .filter(|e| *e.kind() == FileKind::Delete) |
| .map(|e| e.identifier()) |
| .collect(); |
| entries |
| .into_iter() |
| .filter(|e| *e.kind() == FileKind::Add && !deleted.contains(&e.identifier())) |
| .collect() |
| } |
| |
| /// Whether scan-owned pruning still preserves `merged_row_count()` as a safe |
| /// row-count hint. |
| /// |
| /// Data predicates and row ranges can reduce rows within a split after planning, |
| /// so split-level row counts stop being a conservative bound for final rows. |
| pub(super) fn can_push_down_limit_hint_for_scan( |
| data_predicates: &[Predicate], |
| row_ranges: Option<&[RowRange]>, |
| ) -> bool { |
| data_predicates.is_empty() && row_ranges.is_none() |
| } |
| |
| fn should_skip_level_zero_for_scan( |
| scan_all_files: bool, |
| has_primary_keys: bool, |
| deletion_vectors_enabled: bool, |
| merge_engine: crate::Result<crate::spec::MergeEngine>, |
| ) -> bool { |
| if scan_all_files { |
| return false; |
| } |
| if !has_primary_keys { |
| return false; |
| } |
| |
| deletion_vectors_enabled || merge_engine.is_ok_and(|e| e == crate::spec::MergeEngine::FirstRow) |
| } |
| |
| /// TableScan for full table scan (no incremental, no predicate). |
| /// |
| /// Reference: [pypaimon.read.table_scan.TableScan](https://github.com/apache/paimon/blob/master/paimon-python/pypaimon/read/table_scan.py) |
| #[derive(Debug, Clone)] |
| pub struct TableScan<'a> { |
| table: &'a Table, |
| partition_filter: Option<PartitionFilter>, |
| data_predicates: Vec<Predicate>, |
| bucket_predicate: Option<Predicate>, |
| /// Optional limit on the number of rows to return. |
| /// When set, the scan will try to return only enough splits to satisfy the limit. |
| limit: Option<usize>, |
| row_ranges: Option<Vec<RowRange>>, |
| /// When true, disables level-0 filtering so all files are visible. |
| /// Used by non-read paths (overwrite, truncate, writer restore) that need |
| /// the complete file set. Normal read scans leave this as `false`. |
| scan_all_files: bool, |
| } |
| |
| impl<'a> TableScan<'a> { |
| pub(crate) fn new( |
| table: &'a Table, |
| partition_filter: Option<PartitionFilter>, |
| data_predicates: Vec<Predicate>, |
| bucket_predicate: Option<Predicate>, |
| limit: Option<usize>, |
| row_ranges: Option<Vec<RowRange>>, |
| ) -> Self { |
| Self { |
| table, |
| partition_filter, |
| data_predicates, |
| bucket_predicate, |
| limit, |
| row_ranges, |
| scan_all_files: false, |
| } |
| } |
| |
| /// Disable level-0 filtering so all files are visible. |
| /// |
| /// Used by non-read paths (overwrite, truncate, writer restore) that need |
| /// the complete file set regardless of merge engine or DV settings. |
| pub fn with_scan_all_files(mut self) -> Self { |
| self.scan_all_files = true; |
| self |
| } |
| |
| /// Set row ranges for scan-time filtering. |
| /// |
| /// This replaces any existing row_ranges. Typically used to inject |
| /// results from global index lookups (e.g. full-text search). |
| pub fn with_row_ranges(mut self, ranges: Vec<RowRange>) -> Self { |
| self.row_ranges = if ranges.is_empty() { |
| None |
| } else { |
| Some(ranges) |
| }; |
| self |
| } |
| |
| /// Plan the full scan: resolve snapshot (via options or latest), then read manifests and build DataSplits. |
| /// |
| /// Time travel is resolved from table options: |
| /// - only one of `scan.version`, `scan.timestamp-millis` may be set |
| /// - `scan.version` → tag name (if exists) → snapshot id (if parseable) → error |
| /// - `scan.timestamp-millis` → find the latest snapshot <= that timestamp |
| /// - otherwise → read the latest snapshot |
| /// |
| /// Reference: [TimeTravelUtil.tryTravelToSnapshot](https://github.com/apache/paimon/blob/master/paimon-core/src/main/java/org/apache/paimon/table/source/snapshot/TimeTravelUtil.java) |
| pub async fn plan(&self) -> crate::Result<Plan> { |
| let snapshot = match self.resolve_snapshot().await? { |
| Some(snapshot) => snapshot, |
| None => return Ok(Plan::new(Vec::new())), |
| }; |
| self.plan_snapshot(snapshot).await |
| } |
| |
| async fn resolve_snapshot(&self) -> crate::Result<Option<Snapshot>> { |
| // A table copy produced by `copy_with_time_travel` already resolved |
| // the selector in its options; reuse it instead of re-reading |
| // tag/snapshot files on every plan. |
| if let Some(snapshot) = self.table.travel_snapshot() { |
| return Ok(Some(snapshot.clone())); |
| } |
| // A time-travelled schema without its resolved snapshot means the |
| // selector was changed after the travel (`copy_with_options`). |
| // Resolving the new selector here would evolve a different snapshot's |
| // files to the stale historical schema, so fail instead. |
| if self.table.is_time_traveled() { |
| return Err(crate::Error::DataInvalid { |
| message: "Table options changed after time travel; \ |
| use copy_with_time_travel to re-resolve the snapshot and schema" |
| .to_string(), |
| source: None, |
| }); |
| } |
| |
| let file_io = self.table.file_io(); |
| let table_path = self.table.location(); |
| |
| match super::time_travel::travel_to_snapshot( |
| file_io, |
| table_path, |
| self.table.schema().options(), |
| ) |
| .await? |
| { |
| Some(snapshot) => Ok(Some(snapshot)), |
| None => { |
| let snapshot_manager = |
| SnapshotManager::new(file_io.clone(), table_path.to_string()); |
| snapshot_manager.get_latest_snapshot().await |
| } |
| } |
| } |
| |
| /// Apply a limit-pushdown hint to the generated splits. |
| /// |
| /// Mirrors Java `DataTableBatchScan#applyPushDownLimit`: splits whose |
| /// `merged_row_count()` is unknown (for example merge-needed PK splits or |
| /// unknown deletion cardinality) are skipped — they contribute an unknown |
| /// number of rows, so they cannot help satisfy the limit. Pruning is |
| /// committed only once the accumulated known row count reaches the limit; |
| /// if it never does, the original split list is returned unchanged. The |
| /// caller or query engine must still enforce the final LIMIT. |
| fn apply_limit_pushdown(&self, splits: Vec<DataSplit>) -> Vec<DataSplit> { |
| let limit = match self.limit { |
| Some(l) => l, |
| None => return splits, |
| }; |
| if limit == 0 { |
| return Vec::new(); |
| } |
| |
| if splits.is_empty() { |
| return splits; |
| } |
| |
| let mut limited_splits = Vec::new(); |
| let mut scanned_row_count: i64 = 0; |
| |
| for split in &splits { |
| if let Some(merged_count) = split.merged_row_count() { |
| limited_splits.push(split.clone()); |
| scanned_row_count += merged_count; |
| if scanned_row_count >= limit as i64 { |
| return limited_splits; |
| } |
| } |
| } |
| |
| splits |
| } |
| |
| /// Read all manifest entries from a snapshot, applying filters and merging. |
| /// |
| /// This is the shared entry point used by both `plan_snapshot` (scan) and |
| /// `TableCommit` (overwrite). Filters include partition predicate, data |
| /// predicates, and bucket predicate. |
| pub(crate) async fn plan_manifest_entries( |
| &self, |
| snapshot: &Snapshot, |
| ) -> crate::Result<Vec<ManifestEntry>> { |
| let file_io = self.table.file_io(); |
| let table_path = self.table.location(); |
| let core_options = CoreOptions::new(self.table.schema().options()); |
| let data_evolution_enabled = core_options.data_evolution_enabled(); |
| |
| let has_primary_keys = !self.table.schema().primary_keys().is_empty(); |
| let deletion_vectors_enabled = core_options.deletion_vectors_enabled(); |
| |
| // Skip level-0 files for PK tables when: |
| // - DV mode: level-0 files are unmerged, DV handles dedup at higher levels |
| // - FirstRow engine without DV: reads go through DataFileReader (no merge), |
| // so only compacted (level > 0) files are safe to read directly |
| // Deduplicate engine always uses KeyValueFileReader which handles level-0 |
| // via sort-merge, so level-0 files must remain visible. |
| // |
| // Non-read paths (overwrite, truncate, writer restore) set scan_all_files=true |
| // to see all files including level-0, matching Java's CommitScanner behavior. |
| let skip_level_zero = should_skip_level_zero_for_scan( |
| self.scan_all_files, |
| has_primary_keys, |
| deletion_vectors_enabled, |
| core_options.merge_engine(), |
| ); |
| |
| let partition_fields = self.table.schema().partition_fields(); |
| |
| let pushdown_data_predicates = if data_evolution_enabled { |
| &[][..] |
| } else { |
| self.data_predicates.as_slice() |
| }; |
| |
| let bucket_key_fields: Vec<DataField> = if self.bucket_predicate.is_none() { |
| Vec::new() |
| } else { |
| let bucket_keys = core_options.bucket_key().unwrap_or_else(|| { |
| if has_primary_keys { |
| self.table.schema().trimmed_primary_keys() |
| } else { |
| Vec::new() |
| } |
| }); |
| bucket_keys |
| .iter() |
| .filter_map(|key| { |
| self.table |
| .schema() |
| .fields() |
| .iter() |
| .find(|f| f.name() == key) |
| .cloned() |
| }) |
| .collect::<Vec<_>>() |
| }; |
| let bucket_function_type = core_options.bucket_function_type()?; |
| |
| let entries = read_all_manifest_entries( |
| file_io, |
| table_path, |
| snapshot, |
| skip_level_zero, |
| self.scan_all_files, |
| has_primary_keys, |
| self.partition_filter.as_ref(), |
| &partition_fields, |
| pushdown_data_predicates, |
| self.table.schema().id(), |
| self.table.schema().fields(), |
| self.bucket_predicate.as_ref(), |
| &bucket_key_fields, |
| bucket_function_type, |
| ) |
| .await?; |
| Ok(merge_manifest_entries(entries)) |
| } |
| |
| fn can_push_down_limit_hint(&self, row_ranges: Option<&[RowRange]>) -> bool { |
| can_push_down_limit_hint_for_scan(&self.data_predicates, row_ranges) |
| } |
| |
| async fn plan_snapshot(&self, snapshot: Snapshot) -> crate::Result<Plan> { |
| let file_io = self.table.file_io(); |
| let table_path = self.table.location(); |
| let core_options = CoreOptions::new(self.table.schema().options()); |
| let data_evolution_enabled = core_options.data_evolution_enabled(); |
| let target_split_size = core_options.source_split_target_size(); |
| let open_file_cost = core_options.source_split_open_file_cost(); |
| let partition_keys = self.table.schema().partition_keys(); |
| |
| let entries = self.plan_manifest_entries(&snapshot).await?; |
| if entries.is_empty() { |
| return Ok(Plan::new(Vec::new())); |
| } |
| |
| // For non-data-evolution tables, cross-schema files were kept (fail-open) |
| // by the pushdown. Apply the full schema-aware filter for those files. |
| let entries = if self.data_predicates.is_empty() || data_evolution_enabled { |
| entries |
| } else { |
| let current_schema_id = self.table.schema().id(); |
| let has_cross_schema = entries |
| .iter() |
| .any(|e| e.file().schema_id != current_schema_id); |
| if !has_cross_schema { |
| entries |
| } else { |
| let mut kept = Vec::with_capacity(entries.len()); |
| let mut schema_cache: HashMap<i64, Option<Arc<ResolvedStatsSchema>>> = |
| HashMap::new(); |
| for entry in entries { |
| if entry.file().schema_id == current_schema_id |
| || data_file_matches_predicates_for_table( |
| self.table, |
| entry.file(), |
| &self.data_predicates, |
| &mut schema_cache, |
| ) |
| .await |
| { |
| kept.push(entry); |
| } |
| } |
| kept |
| } |
| }; |
| if entries.is_empty() { |
| return Ok(Plan::new(Vec::new())); |
| } |
| |
| // Group by (partition, bucket), decomposing entries to avoid cloning partition. |
| let mut groups: HashMap<(Vec<u8>, i32), (i32, Vec<DataFileMeta>)> = |
| HashMap::with_capacity(entries.len()); |
| for e in entries { |
| let (partition, bucket, total_buckets, file) = e.into_parts(); |
| let entry = groups |
| .entry((partition, bucket)) |
| .or_insert_with(|| (total_buckets, Vec::new())); |
| entry.1.push(file); |
| } |
| |
| let snapshot_id = snapshot.id(); |
| let base_path = table_path.trim_end_matches('/'); |
| let mut splits = Vec::with_capacity(groups.len()); |
| |
| let partition_computer = if !partition_keys.is_empty() { |
| Some(PartitionComputer::new( |
| partition_keys, |
| self.table.schema().fields(), |
| core_options.partition_default_name(), |
| core_options.legacy_partition_name(), |
| )?) |
| } else { |
| None |
| }; |
| |
| // Primary-key tables must keep key-overlapping files in one split so the |
| // sort-merge reader sees every version of a key. The comparator decodes |
| // the trimmed-PK min/max keys written by the kv writer. |
| // |
| // Deletion-vector and first-row tables read without merging (stale rows |
| // are masked by DVs / level-0 is skipped), so they keep plain size-based |
| // packing like Java's MergeTreeSplitGenerator fast path. |
| let read_merges_overlapping_keys = !core_options.deletion_vectors_enabled() |
| && !matches!( |
| core_options.merge_engine(), |
| Ok(crate::spec::MergeEngine::FirstRow) |
| ); |
| let pk_comparator = if read_merges_overlapping_keys { |
| KeyComparator::from_table_schema(self.table.schema()) |
| } else { |
| None |
| }; |
| |
| // Read deletion vector index manifest once (like Java generateSplits / scanDvIndex). |
| let (deletion_files_map, effective_row_ranges) = |
| if let Some(index_manifest_name) = snapshot.index_manifest() { |
| let index_manifest_path = format!("{base_path}/{MANIFEST_DIR}"); |
| let path = format!("{index_manifest_path}/{index_manifest_name}"); |
| let index_entries = IndexManifest::read(file_io, &path).await?; |
| let dv_map = build_deletion_files_map(&index_entries, base_path); |
| |
| // Use pushed-down row_ranges first; otherwise try global index. |
| let row_ranges = if self.row_ranges.is_some() { |
| self.row_ranges.clone() |
| } else if data_evolution_enabled |
| && core_options.global_index_enabled() |
| && !self.data_predicates.is_empty() |
| { |
| super::global_index_scanner::evaluate_global_index( |
| file_io, |
| base_path, |
| &index_entries, |
| &self.data_predicates, |
| self.table.schema().fields(), |
| ) |
| .await? |
| } else { |
| None |
| }; |
| |
| (Some(dv_map), row_ranges) |
| } else { |
| (None, self.row_ranges.clone()) |
| }; |
| |
| for ((partition, bucket), (total_buckets, data_files)) in groups { |
| let partition_row = BinaryRow::from_serialized_bytes(&partition)?; |
| |
| let bucket_path = if let Some(ref computer) = partition_computer { |
| let partition_path = computer.generate_partition_path(&partition_row)?; |
| format!("{base_path}/{partition_path}{}", bucket_dir_name(bucket)) |
| } else { |
| format!("{base_path}/{}", bucket_dir_name(bucket)) |
| }; |
| |
| // Original `partition` Vec consumed by PartitionBucket for DV map lookup. |
| let per_bucket_deletion_map = deletion_files_map |
| .as_ref() |
| .and_then(|map| map.get(&PartitionBucket::new(partition, bucket))); |
| |
| // Data-evolution tables merge overlapping row-id groups column-wise during read. |
| // Keep that split boundary intact and only bin-pack single-file groups. |
| // Apply group-level predicate filtering after grouping by row_id range. |
| let file_groups: Vec<SplitGroup> = if data_evolution_enabled { |
| let row_id_groups = group_by_overlapping_row_id(data_files); |
| |
| // Filter groups by merged stats before splitting. |
| let row_id_groups: Vec<Vec<DataFileMeta>> = if self.data_predicates.is_empty() { |
| row_id_groups |
| } else { |
| row_id_groups |
| .into_iter() |
| .filter(|group| { |
| data_evolution_group_matches_predicates( |
| group, |
| &self.data_predicates, |
| self.table.schema().fields(), |
| ) |
| }) |
| .collect() |
| }; |
| |
| // Filter groups by row ID ranges. |
| let row_id_groups = if let Some(ref ranges) = effective_row_ranges { |
| row_id_groups |
| .into_iter() |
| .filter(|group| group.iter().any(|f| any_range_overlaps_file(ranges, f))) |
| .collect() |
| } else { |
| row_id_groups |
| }; |
| |
| let (singles, multis): (Vec<_>, Vec<_>) = row_id_groups |
| .into_iter() |
| .partition(|group| group.len() == 1); |
| |
| let mut result = Vec::new(); |
| for group in multis { |
| // Files sharing a row-id range hold column slices of the |
| // same logical rows; physical counts overcount them |
| // (Java DataEvolutionSplitGenerator: not raw convertible). |
| result.push(SplitGroup { |
| files: group, |
| raw_convertible: false, |
| }); |
| } |
| |
| let single_files: Vec<DataFileMeta> = singles.into_iter().flatten().collect(); |
| for file_group in split_for_batch(single_files, target_split_size, open_file_cost) { |
| result.push(SplitGroup { |
| files: file_group, |
| raw_convertible: true, |
| }); |
| } |
| |
| result |
| } else if let Some(ref comparator) = pk_comparator { |
| // Merge-tree path: keep key-overlapping files in one split and |
| // mark which splits the sort-merge reader can skip (mirrors |
| // Java MergeTreeSplitGenerator#splitForBatch). Only engines |
| // whose writer deduplicates at flush guarantee a file never |
| // holds two rows of one key, so only they may mark groups raw |
| // convertible; see merge_tree_split_for_batch. (First-row |
| // tables do not take this path today, but its writer dedups |
| // too, so keep the gate accurate.) |
| let file_keys_unique = matches!( |
| core_options.merge_engine(), |
| Ok(crate::spec::MergeEngine::Deduplicate) |
| | Ok(crate::spec::MergeEngine::FirstRow) |
| ); |
| merge_tree_split_for_batch( |
| data_files, |
| comparator, |
| target_split_size, |
| open_file_cost, |
| file_keys_unique, |
| ) |
| } else { |
| split_for_batch(data_files, target_split_size, open_file_cost) |
| .into_iter() |
| .map(|files| SplitGroup { |
| files, |
| raw_convertible: true, |
| }) |
| .collect() |
| }; |
| |
| for group in file_groups { |
| let SplitGroup { |
| files: file_group, |
| raw_convertible, |
| } = group; |
| let data_deletion_files = per_bucket_deletion_map.map(|per_bucket| { |
| file_group |
| .iter() |
| .map(|f| per_bucket.get(&f.file_name).cloned()) |
| .collect::<Vec<Option<DeletionFile>>>() |
| }); |
| |
| // Compute row_ranges before moving file_group to avoid clone |
| let split_row_ranges = if let Some(ref ranges) = effective_row_ranges { |
| let mut split_ranges = Vec::new(); |
| for file in &file_group { |
| split_ranges.extend(intersect_ranges_with_file(ranges, file)); |
| } |
| let split_ranges = merge_row_ranges(split_ranges); |
| if split_ranges.is_empty() { |
| None |
| } else { |
| Some(split_ranges) |
| } |
| } else { |
| None |
| }; |
| |
| let mut builder = DataSplitBuilder::new() |
| .with_snapshot(snapshot_id) |
| .with_partition(partition_row.clone()) |
| .with_bucket(bucket) |
| .with_bucket_path(bucket_path.clone()) |
| .with_total_buckets(total_buckets) |
| .with_data_files(file_group) |
| .with_raw_convertible(raw_convertible); |
| if let Some(files) = data_deletion_files { |
| builder = builder.with_data_deletion_files(files); |
| } |
| if let Some(row_ranges) = split_row_ranges { |
| builder = builder.with_row_ranges(row_ranges); |
| } |
| splits.push(builder.build()?); |
| } |
| } |
| |
| // With data predicates or row_ranges, merged_row_count() reflects pre-filter |
| // row counts, so stopping early could return fewer rows than the limit. |
| let splits = if self.can_push_down_limit_hint(effective_row_ranges.as_deref()) { |
| self.apply_limit_pushdown(splits) |
| } else { |
| splits |
| }; |
| |
| Ok(Plan::new(splits)) |
| } |
| } |
| |
| #[cfg(test)] |
| mod tests { |
| use super::{should_skip_level_zero_for_scan, TableScan}; |
| use crate::catalog::Identifier; |
| use crate::io::FileIOBuilder; |
| use crate::spec::{ |
| stats::BinaryTableStats, ArrayType, BinaryRow, BinaryRowBuilder, BucketFunctionType, |
| DataField, DataFileMeta, DataType, Datum, DeletionVectorMeta, FileKind, IndexFileMeta, |
| IndexManifestEntry, IntType, Predicate, PredicateBuilder, PredicateOperator, |
| Schema as PaimonSchema, TableSchema, VarCharType, |
| }; |
| use crate::table::bucket_filter::{compute_target_buckets, extract_predicate_for_keys}; |
| use crate::table::partition_filter::PartitionFilter; |
| use crate::table::source::{DataSplit, DataSplitBuilder, DeletionFile}; |
| use crate::table::stats_filter::{data_file_matches_predicates, group_by_overlapping_row_id}; |
| use crate::table::Table; |
| use crate::Error; |
| use chrono::{DateTime, Utc}; |
| use std::collections::HashSet; |
| |
| /// Helper to build a DataFileMeta with data evolution fields. |
| fn make_evo_file( |
| name: &str, |
| file_size: i64, |
| row_count: i64, |
| max_seq: i64, |
| first_row_id: Option<i64>, |
| ) -> DataFileMeta { |
| DataFileMeta { |
| file_name: name.to_string(), |
| file_size, |
| row_count, |
| min_key: Vec::new(), |
| max_key: Vec::new(), |
| key_stats: BinaryTableStats::new(Vec::new(), Vec::new(), Vec::new()), |
| value_stats: BinaryTableStats::new(Vec::new(), Vec::new(), Vec::new()), |
| min_sequence_number: 0, |
| max_sequence_number: max_seq, |
| schema_id: 0, |
| level: 0, |
| extra_files: Vec::new(), |
| creation_time: DateTime::<Utc>::from_timestamp(0, 0), |
| delete_row_count: None, |
| embedded_index: None, |
| first_row_id, |
| write_cols: None, |
| external_path: None, |
| file_source: None, |
| value_stats_cols: None, |
| } |
| } |
| |
| #[test] |
| fn test_merge_manifest_entries_keeps_in_place_upgraded_file() { |
| // Reproduces a single-run compaction "upgrade": the SAME file name is |
| // deleted at level 0 and re-added at a higher level (Paimon promotes a |
| // lone sorted run in place instead of rewriting it). Netting ADD/DELETE |
| // by file name alone (ignoring `level`) wrongly drops the upgraded file. |
| // Matches Java `FileEntry.Identifier`, which includes `level`. |
| use super::merge_manifest_entries; |
| use crate::spec::ManifestEntry; |
| |
| let entry = |kind: FileKind, name: &str, level: i32| -> ManifestEntry { |
| let mut file = make_evo_file(name, 1, 1, 1, None); |
| file.level = level; |
| ManifestEntry::new(kind, Vec::new(), 0, 1, file, 2) |
| }; |
| |
| let entries = vec![ |
| entry(FileKind::Add, "f.parquet", 0), // original level-0 write |
| entry(FileKind::Delete, "f.parquet", 0), // compaction removes the L0 version |
| entry(FileKind::Add, "f.parquet", 5), // same file upgraded to level 5 |
| entry(FileKind::Add, "g.parquet", 0), // unrelated fresh file |
| ]; |
| |
| let mut live: Vec<(String, i32)> = merge_manifest_entries(entries) |
| .into_iter() |
| .map(|e| (e.file().file_name.clone(), e.file().level)) |
| .collect(); |
| live.sort(); |
| assert_eq!( |
| live, |
| vec![("f.parquet".to_string(), 5), ("g.parquet".to_string(), 0)], |
| "upgraded file (f@L5) must survive; only f@L0 is cancelled by the DELETE" |
| ); |
| } |
| |
| fn file_names(groups: &[Vec<DataFileMeta>]) -> Vec<Vec<&str>> { |
| groups |
| .iter() |
| .map(|g| g.iter().map(|f| f.file_name.as_str()).collect()) |
| .collect() |
| } |
| |
| fn int_stats_row(value: Option<i32>) -> Vec<u8> { |
| let mut builder = BinaryRowBuilder::new(1); |
| match value { |
| Some(value) => builder.write_int(0, value), |
| None => builder.set_null_at(0), |
| } |
| builder.build_serialized() |
| } |
| |
| fn partition_string_field() -> Vec<DataField> { |
| vec![DataField::new( |
| 0, |
| "dt".to_string(), |
| DataType::VarChar(VarCharType::default()), |
| )] |
| } |
| |
| fn int_field() -> Vec<DataField> { |
| vec![DataField::new( |
| 0, |
| "id".to_string(), |
| DataType::Int(IntType::new()), |
| )] |
| } |
| |
| fn test_data_file_meta( |
| min_values: Vec<u8>, |
| max_values: Vec<u8>, |
| null_counts: Vec<Option<i64>>, |
| row_count: i64, |
| ) -> DataFileMeta { |
| test_data_file_meta_with_schema( |
| min_values, |
| max_values, |
| null_counts, |
| row_count, |
| 0, // default schema_id |
| ) |
| } |
| |
| fn test_data_file_meta_with_schema( |
| min_values: Vec<u8>, |
| max_values: Vec<u8>, |
| null_counts: Vec<Option<i64>>, |
| row_count: i64, |
| schema_id: i64, |
| ) -> DataFileMeta { |
| DataFileMeta { |
| file_name: "test.parquet".into(), |
| file_size: 128, |
| row_count, |
| min_key: Vec::new(), |
| max_key: Vec::new(), |
| key_stats: BinaryTableStats::new(Vec::new(), Vec::new(), Vec::new()), |
| value_stats: BinaryTableStats::new(min_values, max_values, null_counts), |
| min_sequence_number: 0, |
| max_sequence_number: 0, |
| schema_id, |
| level: 1, |
| extra_files: Vec::new(), |
| creation_time: Some(Utc::now()), |
| delete_row_count: None, |
| embedded_index: None, |
| first_row_id: None, |
| write_cols: None, |
| external_path: None, |
| file_source: None, |
| value_stats_cols: None, |
| } |
| } |
| |
| fn limit_test_table() -> Table { |
| let file_io = FileIOBuilder::new("file").build().unwrap(); |
| let schema = PaimonSchema::builder().build().unwrap(); |
| let table_schema = TableSchema::new(0, &schema); |
| Table::new( |
| file_io, |
| Identifier::new("test_db", "test_table"), |
| "/tmp/test-table".to_string(), |
| table_schema, |
| None, |
| ) |
| } |
| |
| fn limit_test_split(file_name: &str, row_count: i64) -> DataSplit { |
| let mut file = test_data_file_meta(Vec::new(), Vec::new(), Vec::new(), row_count); |
| file.file_name = file_name.to_string(); |
| |
| DataSplitBuilder::new() |
| .with_snapshot(1) |
| .with_partition(BinaryRow::new(0)) |
| .with_bucket(0) |
| .with_bucket_path(format!("file:/tmp/{file_name}")) |
| .with_total_buckets(1) |
| .with_data_files(vec![file]) |
| .build() |
| .unwrap() |
| } |
| |
| fn limit_test_split_with_unknown_merged_row_count( |
| file_name: &str, |
| row_count: i64, |
| ) -> DataSplit { |
| let mut file = test_data_file_meta(Vec::new(), Vec::new(), Vec::new(), row_count); |
| file.file_name = file_name.to_string(); |
| |
| DataSplitBuilder::new() |
| .with_snapshot(1) |
| .with_partition(BinaryRow::new(0)) |
| .with_bucket(0) |
| .with_bucket_path(format!("file:/tmp/{file_name}")) |
| .with_total_buckets(1) |
| .with_data_files(vec![file]) |
| .with_data_deletion_files(vec![Some(DeletionFile::new( |
| format!("file:/tmp/{file_name}.dv"), |
| 0, |
| 0, |
| None, |
| ))]) |
| .build() |
| .unwrap() |
| } |
| |
| fn split_file_names(splits: &[DataSplit]) -> Vec<&str> { |
| splits |
| .iter() |
| .map(|split| split.data_files()[0].file_name.as_str()) |
| .collect() |
| } |
| |
| #[test] |
| fn test_apply_limit_pushdown_zero_returns_empty() { |
| let table = limit_test_table(); |
| let scan = TableScan::new(&table, None, vec![], None, Some(0), None); |
| let splits = vec![ |
| limit_test_split("a.parquet", 2), |
| limit_test_split("b.parquet", 3), |
| ]; |
| |
| let pruned = scan.apply_limit_pushdown(splits); |
| |
| assert!(pruned.is_empty()); |
| } |
| |
| /// Java semantics: unknown-count splits are skipped — they cannot prove |
| /// progress toward the limit — and pruning commits once the counted |
| /// splits alone cover the limit. |
| #[test] |
| fn test_apply_limit_pushdown_skips_unknown_merged_row_count() { |
| let table = limit_test_table(); |
| let scan = TableScan::new(&table, None, vec![], None, Some(3), None); |
| let splits = vec![ |
| limit_test_split("a.parquet", 2), |
| limit_test_split_with_unknown_merged_row_count("b.parquet", 4), |
| limit_test_split("c.parquet", 3), |
| ]; |
| |
| let pruned = scan.apply_limit_pushdown(splits); |
| |
| assert_eq!(split_file_names(&pruned), vec!["a.parquet", "c.parquet"]); |
| } |
| |
| /// When counted splits never reach the limit, the original split list is |
| /// returned unchanged (mirrors Java `applyPushDownLimit`). |
| #[test] |
| fn test_apply_limit_pushdown_returns_all_when_limit_not_reached() { |
| let table = limit_test_table(); |
| let scan = TableScan::new(&table, None, vec![], None, Some(100), None); |
| let splits = vec![ |
| limit_test_split("a.parquet", 2), |
| limit_test_split_with_unknown_merged_row_count("b.parquet", 4), |
| limit_test_split("c.parquet", 3), |
| ]; |
| |
| let pruned = scan.apply_limit_pushdown(splits); |
| |
| assert_eq!( |
| split_file_names(&pruned), |
| vec!["a.parquet", "b.parquet", "c.parquet"] |
| ); |
| } |
| |
| /// A non-raw-convertible split (merge-needed PK split) has an unknown |
| /// merged row count: its physical rows overcount merged versions, so it |
| /// must not satisfy the limit on its own. |
| #[test] |
| fn test_apply_limit_pushdown_treats_merge_splits_as_unknown() { |
| let table = limit_test_table(); |
| let scan = TableScan::new(&table, None, vec![], None, Some(3), None); |
| |
| let mut merge_file = test_data_file_meta(Vec::new(), Vec::new(), Vec::new(), 10); |
| merge_file.file_name = "merge.parquet".to_string(); |
| let merge_split = DataSplitBuilder::new() |
| .with_snapshot(1) |
| .with_partition(BinaryRow::new(0)) |
| .with_bucket(0) |
| .with_bucket_path("file:/tmp/merge.parquet".to_string()) |
| .with_total_buckets(1) |
| .with_data_files(vec![merge_file]) |
| .with_raw_convertible(false) |
| .build() |
| .unwrap(); |
| assert_eq!(merge_split.merged_row_count(), None); |
| |
| let splits = vec![merge_split, limit_test_split("a.parquet", 3)]; |
| let pruned = scan.apply_limit_pushdown(splits); |
| |
| // The merge split is skipped; only the counted split commits pruning. |
| assert_eq!(split_file_names(&pruned), vec!["a.parquet"]); |
| } |
| |
| #[test] |
| fn test_first_row_skips_level_zero_by_default() { |
| assert!(should_skip_level_zero_for_scan( |
| false, |
| true, |
| false, |
| Ok(crate::spec::MergeEngine::FirstRow), |
| )); |
| } |
| |
| #[test] |
| fn test_scan_all_files_disables_first_row_level_zero_skip() { |
| assert!(!should_skip_level_zero_for_scan( |
| true, |
| true, |
| false, |
| Ok(crate::spec::MergeEngine::FirstRow), |
| )); |
| } |
| |
| #[test] |
| fn test_partition_filter_decode_failure_fails_open() { |
| let fields = partition_string_field(); |
| let predicate = PredicateBuilder::new(&fields) |
| .equal("dt", Datum::String("2024-01-01".into())) |
| .unwrap(); |
| |
| // Range predicate to force Predicate variant (fail-open path) |
| let filter = PartitionFilter::Predicate(predicate); |
| assert!(filter.matches_entry(&[0xFF, 0x00]).unwrap()); |
| } |
| |
| #[test] |
| fn test_partition_filter_eval_error_fails_fast() { |
| let mut builder = BinaryRowBuilder::new(1); |
| builder.write_string(0, "2024-01-01"); |
| let serialized = builder.build_serialized(); |
| |
| let predicate = Predicate::Leaf { |
| column: "dt".into(), |
| index: 0, |
| data_type: DataType::Array(ArrayType::new(DataType::Int(IntType::new()))), |
| op: PredicateOperator::Eq, |
| literals: vec![Datum::Int(42)], |
| }; |
| |
| let filter = PartitionFilter::Predicate(predicate); |
| let err = filter |
| .matches_entry(&serialized) |
| .expect_err("eval_row error should propagate"); |
| |
| assert!( |
| matches!(&err, Error::Unsupported { message } if message.contains("extract_datum")), |
| "Expected extract_datum unsupported error, got: {err:?}" |
| ); |
| } |
| |
| const TEST_SCHEMA_ID: i64 = 0; |
| fn test_schema_fields() -> Vec<DataField> { |
| int_field() |
| } |
| |
| #[test] |
| fn test_group_by_overlapping_row_id_empty() { |
| let result = group_by_overlapping_row_id(vec![]); |
| assert!(result.is_empty()); |
| } |
| |
| #[test] |
| fn test_group_by_overlapping_row_id_no_row_ids() { |
| let files = vec![ |
| make_evo_file("a", 10, 100, 1, None), |
| make_evo_file("b", 10, 100, 2, None), |
| ]; |
| let groups = group_by_overlapping_row_id(files); |
| assert_eq!(file_names(&groups), vec![vec!["b"], vec!["a"]]); |
| } |
| |
| #[test] |
| fn test_group_by_overlapping_row_id_same_range() { |
| let files = vec![ |
| make_evo_file("a", 10, 100, 2, Some(0)), |
| make_evo_file("b", 10, 100, 1, Some(0)), |
| ]; |
| let groups = group_by_overlapping_row_id(files); |
| assert_eq!(groups.len(), 1); |
| assert_eq!(file_names(&groups), vec![vec!["a", "b"]]); |
| } |
| |
| #[test] |
| fn test_group_by_overlapping_row_id_overlapping_ranges() { |
| let files = vec![ |
| make_evo_file("a", 10, 100, 1, Some(0)), |
| make_evo_file("b", 10, 100, 2, Some(50)), |
| ]; |
| let groups = group_by_overlapping_row_id(files); |
| assert_eq!(groups.len(), 1); |
| assert_eq!(file_names(&groups), vec![vec!["a", "b"]]); |
| } |
| |
| #[test] |
| fn test_group_by_overlapping_row_id_non_overlapping() { |
| let files = vec![ |
| make_evo_file("a", 10, 100, 1, Some(0)), |
| make_evo_file("b", 10, 100, 2, Some(100)), |
| ]; |
| let groups = group_by_overlapping_row_id(files); |
| assert_eq!(groups.len(), 2); |
| assert_eq!(file_names(&groups), vec![vec!["a"], vec!["b"]]); |
| } |
| |
| #[test] |
| fn test_group_by_overlapping_row_id_mixed() { |
| let files = vec![ |
| make_evo_file("a", 10, 100, 1, Some(0)), |
| make_evo_file("b", 10, 100, 2, Some(0)), |
| make_evo_file("c", 10, 100, 3, None), |
| make_evo_file("d", 10, 100, 4, Some(200)), |
| ]; |
| let groups = group_by_overlapping_row_id(files); |
| assert_eq!( |
| file_names(&groups), |
| vec![vec!["c"], vec!["b", "a"], vec!["d"]] |
| ); |
| } |
| |
| #[test] |
| fn test_group_by_overlapping_row_id_sorted_by_seq() { |
| let files = vec![ |
| make_evo_file("a", 10, 100, 1, Some(0)), |
| make_evo_file("b", 10, 100, 3, Some(0)), |
| make_evo_file("c", 10, 100, 2, Some(0)), |
| ]; |
| let groups = group_by_overlapping_row_id(files); |
| assert_eq!(groups.len(), 1); |
| assert_eq!(file_names(&groups), vec![vec!["b", "c", "a"]]); |
| } |
| |
| #[test] |
| fn test_data_file_matches_eq_prunes_out_of_range() { |
| let fields = int_field(); |
| let file = test_data_file_meta( |
| int_stats_row(Some(10)), |
| int_stats_row(Some(20)), |
| vec![Some(0)], |
| 5, |
| ); |
| let predicate = PredicateBuilder::new(&fields) |
| .equal("id", Datum::Int(30)) |
| .unwrap(); |
| |
| assert!(!data_file_matches_predicates( |
| &file, |
| &[predicate], |
| TEST_SCHEMA_ID, |
| &test_schema_fields(), |
| )); |
| } |
| |
| #[test] |
| fn test_data_file_matches_is_null_prunes_when_null_count_is_zero() { |
| let fields = int_field(); |
| let file = test_data_file_meta( |
| int_stats_row(Some(10)), |
| int_stats_row(Some(20)), |
| vec![Some(0)], |
| 5, |
| ); |
| let predicate = PredicateBuilder::new(&fields).is_null("id").unwrap(); |
| |
| assert!(!data_file_matches_predicates( |
| &file, |
| &[predicate], |
| TEST_SCHEMA_ID, |
| &test_schema_fields(), |
| )); |
| } |
| |
| #[test] |
| fn test_data_file_matches_is_not_null_prunes_all_null_file() { |
| let fields = int_field(); |
| let file = test_data_file_meta(int_stats_row(None), int_stats_row(None), vec![Some(5)], 5); |
| let predicate = PredicateBuilder::new(&fields).is_not_null("id").unwrap(); |
| |
| assert!(!data_file_matches_predicates( |
| &file, |
| &[predicate], |
| TEST_SCHEMA_ID, |
| &test_schema_fields(), |
| )); |
| } |
| |
| #[test] |
| fn test_data_file_matches_unsupported_predicate_fails_open() { |
| let fields = int_field(); |
| let file = test_data_file_meta( |
| int_stats_row(Some(10)), |
| int_stats_row(Some(20)), |
| vec![Some(0)], |
| 5, |
| ); |
| let pb = PredicateBuilder::new(&fields); |
| let predicate = Predicate::or(vec![ |
| pb.less_than("id", Datum::Int(5)).unwrap(), |
| pb.greater_than("id", Datum::Int(25)).unwrap(), |
| ]); |
| |
| assert!(data_file_matches_predicates( |
| &file, |
| &[predicate], |
| TEST_SCHEMA_ID, |
| &test_schema_fields(), |
| )); |
| } |
| |
| #[test] |
| fn test_data_file_matches_corrupt_stats_fails_open() { |
| let fields = int_field(); |
| let file = test_data_file_meta(Vec::new(), Vec::new(), vec![Some(0)], 5); |
| let predicate = PredicateBuilder::new(&fields) |
| .equal("id", Datum::Int(30)) |
| .unwrap(); |
| |
| assert!(data_file_matches_predicates( |
| &file, |
| &[predicate], |
| TEST_SCHEMA_ID, |
| &test_schema_fields(), |
| )); |
| } |
| |
| #[test] |
| fn test_data_file_matches_schema_mismatch_fails_open() { |
| let fields = int_field(); |
| let file = test_data_file_meta_with_schema( |
| int_stats_row(Some(10)), |
| int_stats_row(Some(20)), |
| vec![Some(0)], |
| 5, |
| 5, |
| ); |
| let predicate = PredicateBuilder::new(&fields) |
| .equal("id", Datum::Int(30)) |
| .unwrap(); |
| |
| assert!(data_file_matches_predicates( |
| &file, |
| &[predicate], |
| TEST_SCHEMA_ID, |
| &test_schema_fields(), |
| )); |
| } |
| |
| #[test] |
| fn test_data_file_matches_always_false_prunes_despite_schema_mismatch() { |
| let file = test_data_file_meta_with_schema( |
| int_stats_row(Some(10)), |
| int_stats_row(Some(20)), |
| vec![Some(0)], |
| 5, |
| 99, |
| ); |
| |
| assert!(!data_file_matches_predicates( |
| &file, |
| &[Predicate::AlwaysFalse], |
| TEST_SCHEMA_ID, |
| &test_schema_fields(), |
| )); |
| } |
| |
| #[test] |
| fn test_data_file_matches_always_true_keeps_file_despite_schema_mismatch() { |
| let file = test_data_file_meta_with_schema( |
| int_stats_row(Some(10)), |
| int_stats_row(Some(20)), |
| vec![Some(0)], |
| 5, |
| 99, |
| ); |
| |
| assert!(data_file_matches_predicates( |
| &file, |
| &[Predicate::AlwaysTrue], |
| TEST_SCHEMA_ID, |
| &test_schema_fields(), |
| )); |
| } |
| |
| #[test] |
| fn test_build_deletion_files_map_preserves_cardinality() { |
| let entries = vec![IndexManifestEntry { |
| version: 1, |
| kind: FileKind::Add, |
| partition: vec![1, 2, 3], |
| bucket: 7, |
| index_file: IndexFileMeta { |
| index_type: "DELETION_VECTORS".into(), |
| file_name: "index-file".into(), |
| file_size: 128, |
| row_count: 1, |
| deletion_vectors_ranges: Some(indexmap::IndexMap::from([( |
| "data-file.parquet".into(), |
| DeletionVectorMeta { |
| offset: 11, |
| length: 22, |
| cardinality: Some(33), |
| }, |
| )])), |
| global_index_meta: None, |
| }, |
| }]; |
| |
| let map = super::build_deletion_files_map(&entries, "file:/tmp/table"); |
| |
| let by_bucket = map |
| .get(&super::PartitionBucket::new(vec![1, 2, 3], 7)) |
| .expect("partition bucket should exist"); |
| let deletion_file = by_bucket |
| .get("data-file.parquet") |
| .expect("deletion file should exist"); |
| |
| assert_eq!( |
| deletion_file, |
| &DeletionFile::new("file:/tmp/table/index/index-file".into(), 11, 22, Some(33)) |
| ); |
| } |
| |
| // ======================== Bucket predicate filtering ======================== |
| |
| fn bucket_key_fields() -> Vec<DataField> { |
| vec![DataField::new( |
| 0, |
| "id".to_string(), |
| DataType::Int(IntType::new()), |
| )] |
| } |
| |
| #[test] |
| fn test_extract_predicate_for_keys_eq() { |
| let fields = vec![ |
| DataField::new(0, "id".to_string(), DataType::Int(IntType::new())), |
| DataField::new( |
| 1, |
| "name".to_string(), |
| DataType::VarChar(VarCharType::default()), |
| ), |
| ]; |
| let pb = PredicateBuilder::new(&fields); |
| let filter = Predicate::and(vec![ |
| pb.equal("id", Datum::Int(42)).unwrap(), |
| pb.equal("name", Datum::String("alice".into())).unwrap(), |
| ]); |
| |
| let keys = vec!["id".to_string()]; |
| let extracted = extract_predicate_for_keys(&filter, &fields, &keys); |
| assert!(extracted.is_some()); |
| match extracted.unwrap() { |
| Predicate::Leaf { |
| column, index, op, .. |
| } => { |
| assert_eq!(column, "id"); |
| assert_eq!(index, 0); // remapped to key index |
| assert_eq!(op, PredicateOperator::Eq); |
| } |
| other => panic!("expected Leaf, got {other:?}"), |
| } |
| } |
| |
| #[test] |
| fn test_extract_predicate_for_keys_no_match() { |
| let fields = vec![ |
| DataField::new(0, "id".to_string(), DataType::Int(IntType::new())), |
| DataField::new( |
| 1, |
| "name".to_string(), |
| DataType::VarChar(VarCharType::default()), |
| ), |
| ]; |
| let pb = PredicateBuilder::new(&fields); |
| let filter = pb.equal("name", Datum::String("alice".into())).unwrap(); |
| |
| let keys = vec!["id".to_string()]; |
| let extracted = extract_predicate_for_keys(&filter, &fields, &keys); |
| assert!(extracted.is_none()); |
| } |
| |
| #[test] |
| fn test_compute_target_buckets_single_eq() { |
| let fields = bucket_key_fields(); |
| // Build a bucket predicate (already projected to bucket key space, index=0) |
| let pred = Predicate::Leaf { |
| column: "id".into(), |
| index: 0, |
| data_type: DataType::Int(IntType::new()), |
| op: PredicateOperator::Eq, |
| literals: vec![Datum::Int(42)], |
| }; |
| |
| let buckets = compute_target_buckets(&pred, &fields, BucketFunctionType::Default, 4); |
| assert!(buckets.is_some()); |
| let buckets = buckets.unwrap(); |
| assert_eq!(buckets.len(), 1); |
| // The bucket should be deterministic |
| let bucket = *buckets.iter().next().unwrap(); |
| assert!((0..4).contains(&bucket)); |
| } |
| |
| #[test] |
| fn test_compute_target_buckets_in_predicate() { |
| let fields = bucket_key_fields(); |
| let pred = Predicate::Leaf { |
| column: "id".into(), |
| index: 0, |
| data_type: DataType::Int(IntType::new()), |
| op: PredicateOperator::In, |
| literals: vec![Datum::Int(1), Datum::Int(2), Datum::Int(3)], |
| }; |
| |
| let buckets = compute_target_buckets(&pred, &fields, BucketFunctionType::Default, 4); |
| assert!(buckets.is_some()); |
| let buckets = buckets.unwrap(); |
| // Should have at most 3 buckets (could be fewer if some hash to the same bucket) |
| assert!(!buckets.is_empty()); |
| assert!(buckets.len() <= 3); |
| for &b in &buckets { |
| assert!((0..4).contains(&b)); |
| } |
| } |
| |
| #[test] |
| fn test_compute_target_buckets_range_returns_none() { |
| let fields = bucket_key_fields(); |
| let pred = Predicate::Leaf { |
| column: "id".into(), |
| index: 0, |
| data_type: DataType::Int(IntType::new()), |
| op: PredicateOperator::Gt, |
| literals: vec![Datum::Int(10)], |
| }; |
| |
| let buckets = compute_target_buckets(&pred, &fields, BucketFunctionType::Default, 4); |
| assert!( |
| buckets.is_none(), |
| "Range predicates cannot determine target buckets" |
| ); |
| } |
| |
| #[test] |
| fn test_compute_target_buckets_composite_key() { |
| let fields = vec![ |
| DataField::new(0, "a".to_string(), DataType::Int(IntType::new())), |
| DataField::new(1, "b".to_string(), DataType::Int(IntType::new())), |
| ]; |
| let pred = Predicate::And(vec![ |
| Predicate::Leaf { |
| column: "a".into(), |
| index: 0, |
| data_type: DataType::Int(IntType::new()), |
| op: PredicateOperator::Eq, |
| literals: vec![Datum::Int(1)], |
| }, |
| Predicate::Leaf { |
| column: "b".into(), |
| index: 1, |
| data_type: DataType::Int(IntType::new()), |
| op: PredicateOperator::Eq, |
| literals: vec![Datum::Int(2)], |
| }, |
| ]); |
| |
| let buckets = compute_target_buckets(&pred, &fields, BucketFunctionType::Default, 8); |
| assert!(buckets.is_some()); |
| let buckets = buckets.unwrap(); |
| assert_eq!(buckets.len(), 1); |
| let bucket = *buckets.iter().next().unwrap(); |
| assert!((0..8).contains(&bucket)); |
| } |
| |
| #[test] |
| fn test_compute_target_buckets_partial_key_returns_none() { |
| // Only one of two bucket key fields has an eq predicate |
| let fields = vec![ |
| DataField::new(0, "a".to_string(), DataType::Int(IntType::new())), |
| DataField::new(1, "b".to_string(), DataType::Int(IntType::new())), |
| ]; |
| let pred = Predicate::Leaf { |
| column: "a".into(), |
| index: 0, |
| data_type: DataType::Int(IntType::new()), |
| op: PredicateOperator::Eq, |
| literals: vec![Datum::Int(1)], |
| }; |
| |
| let buckets = compute_target_buckets(&pred, &fields, BucketFunctionType::Default, 8); |
| assert!( |
| buckets.is_none(), |
| "Partial bucket key should not determine target buckets" |
| ); |
| } |
| |
| #[test] |
| fn test_compute_target_buckets_string_key() { |
| let fields = vec![DataField::new( |
| 0, |
| "name".to_string(), |
| DataType::VarChar(VarCharType::default()), |
| )]; |
| let pred = Predicate::Leaf { |
| column: "name".into(), |
| index: 0, |
| data_type: DataType::VarChar(VarCharType::default()), |
| op: PredicateOperator::Eq, |
| literals: vec![Datum::String("alice".into())], |
| }; |
| |
| let buckets = compute_target_buckets(&pred, &fields, BucketFunctionType::Default, 4); |
| assert!(buckets.is_some()); |
| let buckets = buckets.unwrap(); |
| assert_eq!(buckets.len(), 1); |
| let bucket = *buckets.iter().next().unwrap(); |
| assert!((0..4).contains(&bucket)); |
| } |
| |
| #[test] |
| fn test_compute_target_buckets_mod_function() { |
| let fields = bucket_key_fields(); |
| let pred = Predicate::Leaf { |
| column: "id".into(), |
| index: 0, |
| data_type: DataType::Int(IntType::new()), |
| op: PredicateOperator::Eq, |
| literals: vec![Datum::Int(-3)], |
| }; |
| |
| let buckets = compute_target_buckets(&pred, &fields, BucketFunctionType::Mod, 5); |
| assert_eq!(buckets, Some(HashSet::from([2]))); |
| } |
| |
| #[test] |
| fn test_compute_target_buckets_hive_function() { |
| let fields = vec![ |
| DataField::new(0, "id".to_string(), DataType::Int(IntType::new())), |
| DataField::new( |
| 1, |
| "name".to_string(), |
| DataType::VarChar(VarCharType::default()), |
| ), |
| ]; |
| let pred = Predicate::And(vec![ |
| Predicate::Leaf { |
| column: "id".into(), |
| index: 0, |
| data_type: DataType::Int(IntType::new()), |
| op: PredicateOperator::Eq, |
| literals: vec![Datum::Int(7)], |
| }, |
| Predicate::Leaf { |
| column: "name".into(), |
| index: 1, |
| data_type: DataType::VarChar(VarCharType::default()), |
| op: PredicateOperator::Eq, |
| literals: vec![Datum::String("hello".into())], |
| }, |
| ]); |
| |
| let buckets = compute_target_buckets(&pred, &fields, BucketFunctionType::Hive, 8); |
| assert_eq!(buckets, Some(HashSet::from([3]))); |
| } |
| |
| #[test] |
| fn test_compute_target_buckets_is_null() { |
| let fields = bucket_key_fields(); |
| let pred = Predicate::Leaf { |
| column: "id".into(), |
| index: 0, |
| data_type: DataType::Int(IntType::new()), |
| op: PredicateOperator::IsNull, |
| literals: vec![], |
| }; |
| |
| let buckets = compute_target_buckets(&pred, &fields, BucketFunctionType::Default, 4); |
| assert!(buckets.is_some(), "IsNull should determine a target bucket"); |
| let buckets = buckets.unwrap(); |
| assert_eq!(buckets.len(), 1); |
| let bucket = *buckets.iter().next().unwrap(); |
| assert!((0..4).contains(&bucket)); |
| |
| // Verify it matches the expected bucket from a null BinaryRow |
| let mut builder = BinaryRowBuilder::new(1); |
| builder.set_null_at(0); |
| let expected = (builder.build().hash_code() % 4).abs(); |
| assert_eq!(bucket, expected); |
| } |
| |
| #[test] |
| fn test_compute_target_buckets_composite_key_with_null() { |
| let fields = vec![ |
| DataField::new(0, "a".to_string(), DataType::Int(IntType::new())), |
| DataField::new(1, "b".to_string(), DataType::Int(IntType::new())), |
| ]; |
| // a = 1 AND b IS NULL |
| let pred = Predicate::And(vec![ |
| Predicate::Leaf { |
| column: "a".into(), |
| index: 0, |
| data_type: DataType::Int(IntType::new()), |
| op: PredicateOperator::Eq, |
| literals: vec![Datum::Int(1)], |
| }, |
| Predicate::Leaf { |
| column: "b".into(), |
| index: 1, |
| data_type: DataType::Int(IntType::new()), |
| op: PredicateOperator::IsNull, |
| literals: vec![], |
| }, |
| ]); |
| |
| let buckets = compute_target_buckets(&pred, &fields, BucketFunctionType::Default, 8); |
| assert!( |
| buckets.is_some(), |
| "Composite key with IsNull should determine a target bucket" |
| ); |
| let buckets = buckets.unwrap(); |
| assert_eq!(buckets.len(), 1); |
| let bucket = *buckets.iter().next().unwrap(); |
| assert!((0..8).contains(&bucket)); |
| } |
| } |