From 63b72697be81c82f2ab7a44658b81b284db3239b Mon Sep 17 00:00:00 2001 From: Li Jiajia Date: Fri, 31 Jul 2026 22:55:55 -0400 Subject: [PATCH] fix(predicate): resolve _ROW_ID by name, not by its placeholder leaf index --- crates/paimon/src/arrow/filtering.rs | 34 +++ crates/paimon/src/arrow/format/orc.rs | 7 +- crates/paimon/src/arrow/format/parquet.rs | 57 ++++- crates/paimon/src/arrow/residual.rs | 166 +++++++++--- crates/paimon/src/predicate_stats.rs | 78 +++--- crates/paimon/src/spec/predicate.rs | 76 +++++- crates/paimon/src/spec/schema.rs | 37 ++- .../paimon/src/table/data_evolution_reader.rs | 239 +++++++++++++++++- crates/paimon/src/table/kv_file_reader.rs | 78 ++++++ crates/paimon/src/table/read_builder.rs | 164 +++++++++++- crates/paimon/src/table/row_id_predicate.rs | 110 +++++--- crates/paimon/src/table/stats_filter.rs | 16 +- crates/paimon/src/table/table_scan.rs | 29 +++ .../paimon/src/table/vector_search_builder.rs | 17 +- 14 files changed, 944 insertions(+), 164 deletions(-) diff --git a/crates/paimon/src/arrow/filtering.rs b/crates/paimon/src/arrow/filtering.rs index 16170317b..1466bd7fd 100644 --- a/crates/paimon/src/arrow/filtering.rs +++ b/crates/paimon/src/arrow/filtering.rs @@ -134,3 +134,37 @@ fn normalize_field_mapping(mapping: Option>, num_fields: usize) -> Vec< }) .unwrap_or_else(|| identity_field_mapping(num_fields)) } + +#[cfg(test)] +mod tests { + use super::*; + use crate::spec::{ + DataType, Datum, IntType, PredicateBuilder, PredicateOperator, ROW_ID_FIELD_NAME, + }; + + #[test] + fn test_a_row_id_branch_can_die_during_per_file_remapping() { + let table_fields = vec![ + DataField::new(0, "id".to_string(), DataType::Int(IntType::new())), + DataField::new(1, "added".to_string(), DataType::Int(IntType::new())), + ]; + let file_fields = vec![table_fields[0].clone()]; + let filter = Predicate::or(vec![ + PredicateBuilder::new(&table_fields) + .is_null("added") + .unwrap(), + Predicate::Leaf { + column: ROW_ID_FIELD_NAME.to_string(), + index: 0, + data_type: DataType::BigInt(crate::spec::BigIntType::new()), + op: PredicateOperator::Eq, + literals: vec![Datum::Long(5)], + }, + ]); + + assert_eq!( + remap_predicates_to_file(&[filter], &table_fields, &file_fields), + vec![Predicate::AlwaysTrue] + ); + } +} diff --git a/crates/paimon/src/arrow/format/orc.rs b/crates/paimon/src/arrow/format/orc.rs index 8c600b0d3..6ebbeb5fd 100644 --- a/crates/paimon/src/arrow/format/orc.rs +++ b/crates/paimon/src/arrow/format/orc.rs @@ -17,7 +17,7 @@ use super::{FilePredicates, FormatFileReader}; use crate::io::FileRead; -use crate::spec::{DataField, DataType, Datum, Predicate, PredicateOperator}; +use crate::spec::{is_row_id_column, DataField, DataType, Datum, Predicate, PredicateOperator}; use crate::table::{ArrowRecordBatchStream, RowRange}; use crate::Error; use async_trait::async_trait; @@ -222,6 +222,7 @@ fn build_orc_leaf_predicate( file_fields: &[DataField], ) -> Option { let Predicate::Leaf { + column, index, op, literals, @@ -230,6 +231,10 @@ fn build_orc_leaf_predicate( else { return None; }; + // Not in the file, and its index would push the wrong column down. + if is_row_id_column(column) { + return None; + } let file_field = file_fields.get(*index)?; let column = file_field.name(); diff --git a/crates/paimon/src/arrow/format/parquet.rs b/crates/paimon/src/arrow/format/parquet.rs index c851299f3..d46194726 100644 --- a/crates/paimon/src/arrow/format/parquet.rs +++ b/crates/paimon/src/arrow/format/parquet.rs @@ -24,8 +24,8 @@ use crate::arrow::{RowFilter, RowFilterContext}; use crate::io::{FileRead, OutputFile}; use crate::spec::stats::BinaryTableStats; use crate::spec::{ - BinaryRowBuilder, CoreOptions, DataField, DataType, Datum, MetadataStatsMode, Predicate, - PredicateOperator, + is_row_id_column, BinaryRowBuilder, CoreOptions, DataField, DataType, Datum, MetadataStatsMode, + Predicate, PredicateOperator, }; use crate::table::{ArrowRecordBatchStream, RowRange}; use crate::Error; @@ -574,11 +574,11 @@ fn build_parquet_arrow_predicate( // the union of referenced Parquet roots, ordered exactly as the projected // RecordBatch. This preserves OR/NOT semantics; splitting it into leaf // RowFilters would incorrectly turn the expression into a conjunction. - let mut field_indices = Vec::new(); - crate::arrow::residual::collect_predicate_field_indices(predicate, &mut field_indices); - let mut projected = field_indices + let mut leaf_refs = Vec::new(); + crate::arrow::residual::collect_predicate_leaf_refs(predicate, &mut leaf_refs); + let mut projected = leaf_refs .into_iter() - .filter_map(|index| { + .filter_map(|(_, index)| { let field = file_fields.get(index)?; parquet_root_index(parquet_schema, field.name()).map(|root| (root, field.clone())) }) @@ -633,11 +633,19 @@ fn parquet_predicate_row_filter_accepted( match predicate { Predicate::AlwaysTrue | Predicate::AlwaysFalse => Ok(true), Predicate::Leaf { + column, index, op, literals, .. - } => parquet_leaf_row_filter_accepted(parquet_schema, *index, *op, literals, file_fields), + } => parquet_leaf_row_filter_accepted( + parquet_schema, + column, + *index, + *op, + literals, + file_fields, + ), Predicate::And(children) | Predicate::Or(children) => { for child in children { if !parquet_predicate_row_filter_accepted(parquet_schema, child, file_fields)? { @@ -660,6 +668,7 @@ fn parquet_predicate_row_filter_accepted( /// an unsupported (but well-formed) leaf yields `Ok(false)`. fn parquet_leaf_row_filter_accepted( parquet_schema: &parquet::schema::types::SchemaDescriptor, + column: &str, index: usize, op: PredicateOperator, literals: &[Datum], @@ -668,6 +677,11 @@ fn parquet_leaf_row_filter_accepted( if !predicate_supported_for_parquet_row_filter(op) { return Ok(false); } + // Not in the file, so the decoder cannot evaluate it. Rejecting the leaf + // rejects any enclosing predicate too, leaving it to the post-scan residual. + if is_row_id_column(column) { + return Ok(false); + } let Some(file_field) = file_fields.get(index) else { return Ok(false); }; @@ -2043,6 +2057,35 @@ mod tests { assert!(row_filter.is_some()); } + fn row_id_leaf() -> Predicate { + Predicate::Leaf { + column: crate::spec::ROW_ID_FIELD_NAME.to_string(), + index: 0, + data_type: DataType::BigInt(crate::spec::BigIntType::new()), + op: super::PredicateOperator::Eq, + literals: vec![Datum::Long(1)], + } + } + + #[test] + fn test_row_id_predicate_builds_no_decoder_row_filter() { + let fields = test_fields(); + let schema = test_parquet_schema(); + assert!(build_parquet_row_filter(&schema, &[row_id_leaf()], &fields) + .expect("row filter should build") + .is_none()); + + let mixed = Predicate::or(vec![ + row_id_leaf(), + PredicateBuilder::new(&fields) + .equal("score", Datum::Int(7)) + .expect("leaf should build"), + ]); + assert!(build_parquet_row_filter(&schema, &[mixed], &fields) + .expect("row filter should build") + .is_none()); + } + // ----------------------------------------------------------------------- // String predicate tests (StartsWith / EndsWith / Contains) // ----------------------------------------------------------------------- diff --git a/crates/paimon/src/arrow/residual.rs b/crates/paimon/src/arrow/residual.rs index 56b77050e..ad38d6e5b 100644 --- a/crates/paimon/src/arrow/residual.rs +++ b/crates/paimon/src/arrow/residual.rs @@ -45,7 +45,7 @@ //! must not reference any vortex-specific types. use crate::arrow::format::FilePredicates; -use crate::spec::{DataField, DataType, Datum, Predicate, PredicateOperator}; +use crate::spec::{is_row_id_column, DataField, DataType, Datum, Predicate, PredicateOperator}; use crate::Error; use arrow_array::{ Array, ArrayRef, BinaryArray, BooleanArray, Date32Array, Datum as ArrowDatum, Decimal128Array, @@ -159,41 +159,55 @@ fn evaluate_predicate_mask( }))) } Predicate::Leaf { + column, index, op, literals, data_type: predicate_data_type, - .. } => { - let Some(file_field) = file_fields.get(*index) else { - return Ok(None); + // Resolve the batch column by NAME, but pick that name carefully: a + // real column may be renamed in the file, so its batch name comes + // from `file_fields[index]`; `_ROW_ID` is absent from `file_fields` + // and only its own name is meaningful. Java's `PredicateRemapper` + // rebinds by name too. + let file_field = if is_row_id_column(column) { + None + } else { + match file_fields.get(*index) { + Some(field) => Some(field), + None => return Ok(None), + } }; - // Resolve the predicate column in the batch by NAME against the batch's - // own schema. We must not index by the column's position in - // `scan_fields`: a reader's emitted batch may order columns by its file - // schema (e.g. ORC `ProjectionMask::named_roots`), not by `scan_fields` - // order, so positional indexing can select the wrong column (and - // compare mismatched types). `scan_fields` is used only to detect the - // Gap-A "predicate column not scanned" bug below. - let column = batch + let field_name = file_field.map_or(column.as_str(), DataField::name); + let field_type = file_field.map_or(predicate_data_type, DataField::data_type); + // Never index by position in `scan_fields`: a reader's emitted batch + // may order columns by its file schema (e.g. ORC + // `ProjectionMask::named_roots`) rather than by `scan_fields` order. + // `scan_fields` is used only for the Gap-A guard below. + let array = batch .schema() - .index_of(file_field.name()) + .index_of(field_name) .ok() .map(|batch_index| batch.column(batch_index)); - let Some(column) = column else { - // The predicate column exists in the file schema but is absent - // from the batch actually scanned — this is the Gap-A bug (a - // reader that did not widen its scan to include predicate columns - // before filtering). It must never happen. Fail loudly in - // debug/test builds; degrade to a skip (rather than panic) in - // release. `scan_fields` is unused for resolution now (we look up - // by name in the batch), so touch it here only to keep the guard - // message informative. + let Some(array) = array else { let _ = scan_fields; + if file_field.is_none() { + // A read that cannot synthesize `_ROW_ID`; skipping the + // predicate would return rows that do not match it. Covers + // residuals applied to a predicate-free reader (PK merge + // output, vector search); other readers reject earlier. + return Err(Error::Unsupported { + message: format!( + "filtering on '{field_name}' is not supported by this read; it is \ + available on data-evolution reads, or via row ranges" + ), + }); + } + // A real column missing here is a reader bug: it did not widen + // its scan to the predicate columns. debug_assert!( false, - "residual predicate column '{}' exists in file_fields but is missing from the scanned batch; the reader must widen its scan to include predicate columns", - file_field.name() + "residual predicate column '{field_name}' is missing from the scanned batch; the reader must widen its scan to include predicate columns" ); return Ok(None); }; @@ -206,16 +220,14 @@ fn evaluate_predicate_mask( // column up to the predicate type first — then the literal is always // representable and the comparison is exact. let predicate_arrow_type = crate::arrow::paimon_type_to_arrow(predicate_data_type)?; - let mask = if column.data_type() == &predicate_arrow_type { - evaluate_exact_leaf_predicate(column, file_field.data_type(), *op, literals) + let mask = if array.data_type() == &predicate_arrow_type { + evaluate_exact_leaf_predicate(array, field_type, *op, literals) } else { - let cast_column = arrow_cast::cast(column, &predicate_arrow_type).map_err(|e| { + let cast_column = arrow_cast::cast(array, &predicate_arrow_type).map_err(|e| { Error::DataInvalid { message: format!( - "Failed to cast residual column '{}' from {:?} to {:?}: {e}", - file_field.name(), - column.data_type(), - predicate_arrow_type + "Failed to cast residual column '{field_name}' from {:?} to {predicate_arrow_type:?}: {e}", + array.data_type() ), source: Some(Box::new(e)), } @@ -251,11 +263,15 @@ pub(crate) fn widen_scan_fields( let mut fields = read_fields.to_vec(); if let Some(fp) = predicates { - let mut predicate_indices = Vec::new(); + let mut refs = Vec::new(); for predicate in &fp.predicates { - collect_predicate_field_indices(predicate, &mut predicate_indices); + collect_predicate_leaf_refs(predicate, &mut refs); } - for index in predicate_indices { + for (name, index) in refs { + // Not read from the file, so there is nothing to widen with. + if is_row_id_column(name) { + continue; + } if let Some(field) = fp.file_fields.get(index) { push_unique_scan_field(&mut fields, field); } @@ -265,15 +281,22 @@ pub(crate) fn widen_scan_fields( fields } -pub(crate) fn collect_predicate_field_indices(predicate: &Predicate, indices: &mut Vec) { +/// Collect every leaf as `(column name, leaf index)`. +/// +/// Callers that resolve a leaf positionally must check the name first — see +/// [`crate::spec::is_row_id_column`]. +pub(crate) fn collect_predicate_leaf_refs<'a>( + predicate: &'a Predicate, + refs: &mut Vec<(&'a str, usize)>, +) { match predicate { - Predicate::Leaf { index, .. } => indices.push(*index), + Predicate::Leaf { column, index, .. } => refs.push((column.as_str(), *index)), Predicate::And(children) | Predicate::Or(children) => { for child in children { - collect_predicate_field_indices(child, indices); + collect_predicate_leaf_refs(child, refs); } } - Predicate::Not(inner) => collect_predicate_field_indices(inner, indices), + Predicate::Not(inner) => collect_predicate_leaf_refs(inner, refs), Predicate::AlwaysTrue | Predicate::AlwaysFalse => {} } } @@ -1038,7 +1061,7 @@ fn float64_literal(literal: &Datum) -> Option { #[cfg(test)] mod tests { use super::*; - use crate::spec::{IntType, VarCharType}; + use crate::spec::{IntType, VarCharType, ROW_ID_FIELD_NAME}; use arrow_array::{Int32Array, StringArray}; use arrow_schema::{DataType as ArrowDataType, Field as ArrowField, Schema as ArrowSchema}; use std::sync::Arc; @@ -1520,6 +1543,71 @@ mod tests { let _ = filter_record_batch_by_predicates(batch, &fp, &scan_fields); } + fn row_id_leaf(op: PredicateOperator, literals: Vec) -> Predicate { + Predicate::Leaf { + column: ROW_ID_FIELD_NAME.to_string(), + index: 0, + data_type: DataType::BigInt(crate::spec::BigIntType::new()), + op, + literals, + } + } + + fn batch_with_row_id(values: Vec, row_ids: Vec) -> RecordBatch { + let schema = Arc::new(ArrowSchema::new(vec![ + ArrowField::new("v", ArrowDataType::Int32, true), + ArrowField::new(ROW_ID_FIELD_NAME, ArrowDataType::Int64, true), + ])); + RecordBatch::try_new( + schema, + vec![ + Arc::new(Int32Array::from(values)), + Arc::new(arrow_array::Int64Array::from(row_ids)), + ], + ) + .unwrap() + } + + #[test] + fn test_row_id_leaf_resolves_by_name_not_by_placeholder_index() { + let v = int_field(1, "v"); + let batch = batch_with_row_id(vec![99, 7, 7], vec![1, 2, 3]); + let pred = row_id_leaf(PredicateOperator::Eq, vec![Datum::Long(1)]); + let fp = file_predicates(vec![pred], vec![v.clone()]); + let out = filter_record_batch_by_predicates(batch, &fp, &[v]).unwrap(); + assert_eq!(int_values(&out), vec![99]); + } + + #[test] + fn test_row_id_leaf_inside_a_disjunction_resolves_by_name() { + let v = int_field(1, "v"); + let batch = batch_with_row_id(vec![99, 7, 5], vec![1, 2, 3]); + let pred = Predicate::or(vec![ + row_id_leaf(PredicateOperator::Eq, vec![Datum::Long(1)]), + leaf( + 0, + DataType::Int(IntType::new()), + PredicateOperator::Eq, + vec![Datum::Int(7)], + ), + ]); + let fp = file_predicates(vec![pred], vec![v.clone()]); + let out = filter_record_batch_by_predicates(batch, &fp, &[v]).unwrap(); + assert_eq!(int_values(&out), vec![99, 7]); + } + + #[test] + fn test_widen_scan_fields_skips_system_columns() { + let other = int_field(1, "other"); + let v = int_field(2, "v"); + let fp = file_predicates( + vec![row_id_leaf(PredicateOperator::Eq, vec![Datum::Long(1)])], + vec![other, v.clone()], + ); + let widened = widen_scan_fields(std::slice::from_ref(&v), Some(&fp)); + assert_eq!(widened, vec![v]); + } + #[test] fn test_filter_when_batch_column_order_differs_from_scan_fields() { // Regression: a reader (e.g. ORC `ProjectionMask::named_roots`) may emit diff --git a/crates/paimon/src/predicate_stats.rs b/crates/paimon/src/predicate_stats.rs index 846970b22..aec0c55b1 100644 --- a/crates/paimon/src/predicate_stats.rs +++ b/crates/paimon/src/predicate_stats.rs @@ -15,7 +15,7 @@ // specific language governing permissions and limitations // under the License. -use crate::spec::{DataField, DataType, Datum, Predicate, PredicateOperator}; +use crate::spec::{is_row_id_column, DataField, DataType, Datum, Predicate, PredicateOperator}; use std::cmp::Ordering; pub(crate) trait StatsAccessor { @@ -427,27 +427,33 @@ fn predicate_may_match_with_schema( !predicate_must_match_with_schema(inner, stats, field_mapping, file_fields) } Predicate::Leaf { + column, index, data_type, op, literals, - .. - } => match field_mapping.get(*index).copied().flatten() { - Some(file_index) => { - let Some(file_field) = file_fields.get(file_index) else { - return true; - }; - data_leaf_may_match( - file_index, - file_field.data_type(), - data_type, - *op, - literals, - stats, - ) + } => { + // `_ROW_ID` has no column stats, so never prune on it. + if is_row_id_column(column) { + return true; + } + match field_mapping.get(*index).copied().flatten() { + Some(file_index) => { + let Some(file_field) = file_fields.get(file_index) else { + return true; + }; + data_leaf_may_match( + file_index, + file_field.data_type(), + data_type, + *op, + literals, + stats, + ) + } + None => missing_field_may_match(*op, stats.row_count()), } - None => missing_field_may_match(*op, stats.row_count()), - }, + } } } @@ -470,27 +476,33 @@ fn predicate_must_match_with_schema( !predicate_may_match_with_schema(inner, stats, field_mapping, file_fields) } Predicate::Leaf { + column, index, data_type, op, literals, - .. - } => match field_mapping.get(*index).copied().flatten() { - Some(file_index) => { - let Some(file_field) = file_fields.get(file_index) else { - return false; - }; - data_leaf_must_match( - file_index, - file_field.data_type(), - data_type, - *op, - literals, - stats, - ) + } => { + // Stats cannot decide `_ROW_ID`, so it never provably matches. + if is_row_id_column(column) { + return false; + } + match field_mapping.get(*index).copied().flatten() { + Some(file_index) => { + let Some(file_field) = file_fields.get(file_index) else { + return false; + }; + data_leaf_must_match( + file_index, + file_field.data_type(), + data_type, + *op, + literals, + stats, + ) + } + None => missing_field_must_match(*op, stats.row_count()), } - None => missing_field_must_match(*op, stats.row_count()), - }, + } } } diff --git a/crates/paimon/src/spec/predicate.rs b/crates/paimon/src/spec/predicate.rs index f91573b7f..ce4c19bcd 100644 --- a/crates/paimon/src/spec/predicate.rs +++ b/crates/paimon/src/spec/predicate.rs @@ -26,7 +26,7 @@ use crate::error::*; use crate::spec::binary_row::BinaryRow; use crate::spec::types::DataType; -use crate::spec::DataField; +use crate::spec::{is_row_id_column, DataField}; use std::cmp::Ordering; use std::fmt; @@ -440,6 +440,9 @@ impl Predicate { op, literals, } => { + if is_row_id_column(column) { + return None; + } let new_index = (*mapping.get(*index)?)?; Some(Predicate::Leaf { column: column.clone(), @@ -478,7 +481,9 @@ impl Predicate { /// retained as a residual data predicate after partition projection. pub(crate) fn references_only_mapped_fields(&self, mapping: &[Option]) -> bool { match self { - Predicate::Leaf { index, .. } => mapping.get(*index).is_some_and(Option::is_some), + Predicate::Leaf { column, index, .. } => { + !is_row_id_column(column) && mapping.get(*index).is_some_and(Option::is_some) + } Predicate::And(children) | Predicate::Or(children) => children .iter() .all(|child| child.references_only_mapped_fields(mapping)), @@ -489,6 +494,8 @@ impl Predicate { /// Project leaf field indices from table schema space into a smaller field space. /// + /// A `_ROW_ID` leaf never maps — see [`is_row_id_column`]. + /// /// Unlike [`Self::remap_field_index`], mixed `AND` subtrees keep the children /// that can be projected and drop the rest. `OR` and `NOT` still require all /// children to be projectable to preserve correctness. @@ -507,6 +514,9 @@ impl Predicate { op, literals, } => { + if is_row_id_column(column) { + return None; + } let new_index = (*mapping.get(*index)?)?; Some(Predicate::Leaf { column: column.clone(), @@ -1726,6 +1736,68 @@ mod tests { ] } + fn row_id_leaf() -> Predicate { + Predicate::Leaf { + column: crate::spec::ROW_ID_FIELD_NAME.to_string(), + index: 0, + data_type: DataType::BigInt(BigIntType::new()), + op: PredicateOperator::GtEq, + literals: vec![Datum::Long(10)], + } + } + + #[test] + fn test_system_column_never_maps_onto_a_field_space() { + let fields = test_fields(); + let partition_first = [ + fields[2].clone(), + fields[0].clone(), + fields[1].clone(), + fields[3].clone(), + ]; + let mapping = + field_idx_to_partition_idx(&partition_first, &["dt".to_string(), "hr".to_string()]); + assert_eq!(mapping[0], Some(0), "field 0 is the dt partition key"); + + let leaf = row_id_leaf(); + assert_eq!(leaf.project_field_index_inclusive(&mapping), None); + assert_eq!(leaf.remap_field_index(&mapping), None); + assert!(!leaf.references_only_mapped_fields(&mapping)); + } + + #[test] + fn test_system_column_does_not_drag_a_conjunction_onto_a_field_space() { + let fields = test_fields(); + let mapping = field_idx_to_partition_idx(&fields, &["hr".to_string()]); + let dt = PredicateBuilder::new(&fields) + .greater_or_equal("hr", Datum::Int(3)) + .unwrap(); + + let conjunction = Predicate::and(vec![row_id_leaf(), dt.clone()]); + let projected = conjunction.project_field_index_inclusive(&mapping).unwrap(); + assert!(!matches!(projected, Predicate::And(_)), "only hr survives"); + assert!(!conjunction.references_only_mapped_fields(&mapping)); + + let disjunction = Predicate::or(vec![row_id_leaf(), dt]); + assert_eq!(disjunction.project_field_index_inclusive(&mapping), None); + } + + #[test] + fn test_other_reserved_names_are_still_ordinary_columns() { + let fields = vec![ + DataField::new(0, "rowkind".to_string(), DataType::Int(IntType::new())), + DataField::new(1, "hr".to_string(), DataType::Int(IntType::new())), + ]; + let mapping = field_idx_to_partition_idx(&fields, &["rowkind".to_string()]); + let leaf = PredicateBuilder::new(&fields) + .equal("rowkind", Datum::Int(1)) + .unwrap(); + + assert!(leaf.project_field_index_inclusive(&mapping).is_some()); + assert!(leaf.remap_field_index(&mapping).is_some()); + assert!(leaf.references_only_mapped_fields(&mapping)); + } + // ======================== PredicateBuilder basics ======================== #[test] diff --git a/crates/paimon/src/spec/schema.rs b/crates/paimon/src/spec/schema.rs index b3c48d13f..e94c9f496 100644 --- a/crates/paimon/src/spec/schema.rs +++ b/crates/paimon/src/spec/schema.rs @@ -203,14 +203,6 @@ impl TableSchema { /// colliding with a system field (e.g. `_ROW_ID`) is otherwise excluded /// from the physical read and silently filled with the system value. fn validate_no_reserved_fields(&self) -> crate::Result<()> { - // Java SpecialFields.SYSTEM_FIELD_NAMES. - const SYSTEM_FIELD_NAMES: [&str; 5] = [ - SEQUENCE_NUMBER_FIELD_NAME, - VALUE_KIND_FIELD_NAME, - "_LEVEL", - ROW_KIND_FIELD_NAME, - ROW_ID_FIELD_NAME, - ]; const KEY_FIELD_PREFIX: &str = "_KEY_"; // Java SpecialFields.SYSTEM_FIELD_ID_START = Integer.MAX_VALUE / 2. const SYSTEM_FIELD_ID_START: i32 = i32::MAX / 2; @@ -983,6 +975,35 @@ pub const VALUE_KIND_FIELD_ID: i32 = i32::MAX - 2; pub const ROW_KIND_FIELD_NAME: &str = "rowkind"; +/// Java `SpecialFields.SYSTEM_FIELD_NAMES`. +const SYSTEM_FIELD_NAMES: [&str; 5] = [ + SEQUENCE_NUMBER_FIELD_NAME, + VALUE_KIND_FIELD_NAME, + "_LEVEL", + ROW_KIND_FIELD_NAME, + ROW_ID_FIELD_NAME, +]; + +/// A row's global id. Nullable: a data file lacking `first_row_id` yields nulls. +pub(crate) fn row_id_data_field() -> DataField { + DataField::new( + ROW_ID_FIELD_ID, + ROW_ID_FIELD_NAME.to_string(), + DataType::BigInt(crate::spec::BigIntType::with_nullable(true)), + ) +} + +/// `_ROW_ID` is synthesized by the reader and is not a table column, so +/// `PredicateBuilder` cannot resolve it and callers hand-build the leaf with a +/// placeholder index. Every index-based resolution must recognize it by name +/// instead, or it binds the predicate to whatever field sits at that index. +/// +/// The other reserved names are excluded on purpose: `Schema::builder` accepts +/// them as real columns, since only `validate_no_reserved_fields` rejects them. +pub(crate) fn is_row_id_column(name: &str) -> bool { + name == ROW_ID_FIELD_NAME +} + /// Must match Java Paimon's `SpecialFields.ROW_KIND` (Integer.MAX_VALUE - 4). pub const ROW_KIND_FIELD_ID: i32 = i32::MAX - 4; diff --git a/crates/paimon/src/table/data_evolution_reader.rs b/crates/paimon/src/table/data_evolution_reader.rs index 3e62cd3e0..c20c64287 100644 --- a/crates/paimon/src/table/data_evolution_reader.rs +++ b/crates/paimon/src/table/data_evolution_reader.rs @@ -131,7 +131,7 @@ impl DataEvolutionReader { blob_view_resolve_enabled: bool, blob_view_rest_env: Option, ) -> crate::Result { - let row_id_index = read_type.iter().position(|f| f.name() == ROW_ID_FIELD_NAME); + let projected_row_id_index = read_type.iter().position(|f| f.name() == ROW_ID_FIELD_NAME); let file_read_type: Vec = read_type .iter() .filter(|f| f.name() != ROW_ID_FIELD_NAME) @@ -152,11 +152,25 @@ impl DataEvolutionReader { let wide_file_read_type = crate::arrow::residual::widen_scan_fields(&file_read_type, file_predicates.as_ref()); // Wide batches at the _ROW_ID attach point: original read_type columns - // (caller order, _ROW_ID at row_id_index) followed by the extras. - // row_id_index <= file_read_type.len(), so inserting _ROW_ID never - // displaces a trailing extra column. + // (caller order, _ROW_ID at row_id_index) followed by the extras. A + // projected row_id_index is <= file_read_type.len(), so inserting + // _ROW_ID never displaces a trailing extra column. let mut wide_read_type = read_type; wide_read_type.extend_from_slice(&wide_file_read_type[file_read_type.len()..]); + // A residual on `_ROW_ID` needs the column when `filter_wide_batch` + // runs, even unprojected. `widen_scan_fields` cannot supply a + // synthesized column, so append it and let `project_output` trim it. + let row_id_index = match projected_row_id_index { + Some(index) => Some(index), + None if predicates + .iter() + .any(super::row_id_predicate::references_row_id) => + { + wide_read_type.push(crate::spec::row_id_data_field()); + Some(wide_read_type.len() - 1) + } + None => None, + }; let wide_output_schema = build_target_arrow_schema(&wide_read_type)?; Ok(Self { @@ -390,8 +404,9 @@ impl DataEvolutionReader { /// /// Layout invariant: the first `output_schema.fields().len()` columns of /// `batch` are exactly the original read_type columns. Extras were appended - /// at the end by `widen_scan_fields`, and `_ROW_ID` insertion at - /// `row_id_index` keeps them trailing. + /// at the end — by `widen_scan_fields`, plus an unprojected `_ROW_ID` a + /// residual needs — and `_ROW_ID` insertion at `row_id_index` keeps them + /// trailing. fn project_output(&self, filtered: RecordBatch) -> crate::Result { let final_width = self.output_schema.fields().len(); if filtered.num_columns() == final_width { @@ -836,6 +851,11 @@ fn predicate_references_any_field( ) -> bool { match predicate { Predicate::Leaf { column, index, .. } => { + // Never a BLOB column; resolving its placeholder index would force + // every BLOB to be resolved before filtering. + if crate::spec::is_row_id_column(column) { + return false; + } field_names.contains(column) || table_fields .get(*index) @@ -6502,6 +6522,213 @@ mod tests { ); } + #[tokio::test] + async fn test_evolution_read_applies_row_id_residual_to_the_row_ids() { + let tempdir = tempdir().unwrap(); + let table_path = local_file_path(tempdir.path()); + let bucket_dir = tempdir.path().join("bucket-0"); + fs::create_dir_all(&bucket_dir).unwrap(); + + let parquet_path = bucket_dir.join("data.parquet"); + write_int_parquet_file(&parquet_path, vec![("id", vec![1, 2, 3, 4])], None); + + let table = two_col_evolution_table(table_path); + let split = DataSplitBuilder::new() + .with_snapshot(1) + .with_partition(BinaryRow::new(0)) + .with_bucket(0) + .with_bucket_path(local_file_path(&bucket_dir)) + .with_total_buckets(1) + .with_data_files(vec![data_file_meta_with_path( + "data.parquet", + 100, + 4, + 1, + parquet_path.metadata().unwrap().len() as i64, + Some(vec!["id"]), + )]) + .build() + .unwrap(); + + let predicate = Predicate::Leaf { + column: crate::spec::ROW_ID_FIELD_NAME.to_string(), + index: 0, + data_type: DataType::BigInt(crate::spec::BigIntType::new()), + op: crate::spec::PredicateOperator::NotEq, + literals: vec![Datum::Long(102)], + }; + + let mut builder = table.new_read_builder(); + builder.with_projection(&["id"]).unwrap(); + builder.with_filter(predicate); + let read = builder.new_read().unwrap(); + let batches = read + .to_arrow(&[split]) + .unwrap() + .try_collect::>() + .await + .unwrap(); + + assert_eq!(collect_int_values(&batches, "id"), vec![1, 2, 4]); + assert_eq!(batches[0].num_columns(), 1); + } + + #[test] + fn test_row_id_filter_is_not_blob_dependent() { + let table_fields = vec![ + DataField::new(0, "payload".to_string(), DataType::Blob(BlobType::new())), + DataField::new(1, "v".to_string(), DataType::Int(IntType::new())), + ]; + let blob_fields: HashSet = ["payload".to_string()].into_iter().collect(); + let row_id = Predicate::Leaf { + column: crate::spec::ROW_ID_FIELD_NAME.to_string(), + index: 0, + data_type: DataType::BigInt(crate::spec::BigIntType::new()), + op: crate::spec::PredicateOperator::Between, + literals: vec![Datum::Long(10), Datum::Long(20)], + }; + + assert!(!predicate_references_any_field( + &row_id, + &blob_fields, + &table_fields + )); + let on_blob = PredicateBuilder::new(&table_fields) + .is_null("payload") + .unwrap(); + assert!(predicate_references_any_field( + &on_blob, + &blob_fields, + &table_fields + )); + } + + #[tokio::test] + async fn test_evolution_read_enforces_a_row_id_leaf_nested_in_a_disjunction() { + let tempdir = tempdir().unwrap(); + let table_path = local_file_path(tempdir.path()); + let bucket_dir = tempdir.path().join("bucket-0"); + fs::create_dir_all(&bucket_dir).unwrap(); + let parquet_path = bucket_dir.join("data.parquet"); + write_int_parquet_file(&parquet_path, vec![("id", vec![1, 2, 3, 4])], None); + + let table = two_col_evolution_table(table_path); + let split = DataSplitBuilder::new() + .with_snapshot(1) + .with_partition(BinaryRow::new(0)) + .with_bucket(0) + .with_bucket_path(local_file_path(&bucket_dir)) + .with_total_buckets(1) + .with_data_files(vec![data_file_meta_with_path( + "data.parquet", + 100, + 4, + 1, + parquet_path.metadata().unwrap().len() as i64, + Some(vec!["id"]), + )]) + .build() + .unwrap(); + + let pb = PredicateBuilder::new(table.schema().fields()); + let mut builder = table.new_read_builder(); + builder.with_projection(&["id"]).unwrap(); + let rid = |op, v| Predicate::Leaf { + column: crate::spec::ROW_ID_FIELD_NAME.to_string(), + index: 0, + data_type: DataType::BigInt(crate::spec::BigIntType::new()), + op, + literals: vec![Datum::Long(v)], + }; + builder.with_filter(Predicate::and(vec![ + rid(crate::spec::PredicateOperator::GtEq, 100), + Predicate::or(vec![ + Predicate::and(vec![ + rid(crate::spec::PredicateOperator::Eq, 999), + pb.equal("id", Datum::Int(2)).unwrap(), + ]), + pb.equal("id", Datum::Int(4)).unwrap(), + ]), + ])); + let batches = builder + .new_read() + .unwrap() + .to_arrow(&[split]) + .unwrap() + .try_collect::>() + .await + .unwrap(); + + assert_eq!(collect_int_values(&batches, "id"), vec![4]); + } + + #[tokio::test] + async fn test_row_id_filter_without_data_evolution_is_rejected() { + let tempdir = tempdir().unwrap(); + let table_path = local_file_path(tempdir.path()); + let bucket_dir = tempdir.path().join("bucket-0"); + fs::create_dir_all(&bucket_dir).unwrap(); + let parquet_path = bucket_dir.join("data.parquet"); + write_int_parquet_file(&parquet_path, vec![("id", vec![1, 2, 3, 4])], None); + + let file_io = FileIOBuilder::new("file").build().unwrap(); + let table_schema = TableSchema::new( + 0, + &Schema::builder() + .column("id", DataType::Int(IntType::new())) + .column("value", DataType::Int(IntType::new())) + .build() + .unwrap(), + ); + let table = Table::new( + file_io, + Identifier::new("default", "plain_t"), + table_path, + table_schema, + None, + ); + let split = DataSplitBuilder::new() + .with_snapshot(1) + .with_partition(BinaryRow::new(0)) + .with_bucket(0) + .with_bucket_path(local_file_path(&bucket_dir)) + .with_total_buckets(1) + .with_data_files(vec![data_file_meta_with_path( + "data.parquet", + 100, + 4, + 1, + parquet_path.metadata().unwrap().len() as i64, + Some(vec!["id"]), + )]) + .build() + .unwrap(); + + let predicate = Predicate::Leaf { + column: crate::spec::ROW_ID_FIELD_NAME.to_string(), + index: 0, + data_type: DataType::BigInt(crate::spec::BigIntType::new()), + op: crate::spec::PredicateOperator::NotEq, + literals: vec![Datum::Long(102)], + }; + let mut builder = table.new_read_builder(); + builder.with_projection(&["id"]).unwrap(); + builder.with_filter(predicate); + let err = builder + .new_read() + .unwrap() + .to_arrow(&[split]) + .unwrap() + .try_collect::>() + .await + .unwrap_err(); + + assert!( + matches!(&err, crate::Error::Unsupported { message } if message.contains("_ROW_ID")), + "unexpected error: {err:?}" + ); + } + /// _ROW_ID + predicate, merge branch: same guarantee across a column merge. #[tokio::test] async fn test_evolution_read_row_id_with_predicate_merge_branch() { diff --git a/crates/paimon/src/table/kv_file_reader.rs b/crates/paimon/src/table/kv_file_reader.rs index 64ffce73b..d6f03e6d7 100644 --- a/crates/paimon/src/table/kv_file_reader.rs +++ b/crates/paimon/src/table/kv_file_reader.rs @@ -583,6 +583,84 @@ mod tests { use futures::TryStreamExt; use std::sync::Arc; + #[tokio::test] + async fn test_row_id_filter_on_a_primary_key_table_is_rejected() { + let file_io = test_file_io(); + let table_path = "memory:/kv_row_id_filter"; + setup_dirs(&file_io, table_path).await; + let table = pk_table(&file_io, table_path, &[]); + + write_commit( + &table, + &int_batch(vec![1, 2, 3], vec![Some(10), Some(20), Some(30)]), + ) + .await; + write_commit( + &table, + &int_batch(vec![1, 2, 3], vec![Some(11), Some(21), Some(31)]), + ) + .await; + + let row_id = Predicate::Leaf { + column: crate::spec::ROW_ID_FIELD_NAME.to_string(), + index: 0, + data_type: DataType::BigInt(crate::spec::BigIntType::new()), + op: crate::spec::PredicateOperator::NotEq, + literals: vec![Datum::Long(102)], + }; + let mut read_builder = table.new_read_builder(); + read_builder.with_filter(row_id); + let plan = read_builder.new_scan().plan().await.unwrap(); + let err = read_builder + .new_read() + .unwrap() + .to_arrow(plan.splits()) + .unwrap() + .try_collect::>() + .await + .unwrap_err(); + + assert!( + matches!(&err, Error::Unsupported { message } if message.contains("_ROW_ID")), + "unexpected error: {err:?}" + ); + } + + #[test] + fn test_row_id_conjunct_is_not_treated_as_a_primary_key_conjunct() { + let fields = vec![ + DataField::new(0, "k".to_string(), DataType::Int(IntType::new())), + DataField::new(1, "v".to_string(), DataType::Int(IntType::new())), + ]; + let row_id = Predicate::Leaf { + column: crate::spec::ROW_ID_FIELD_NAME.to_string(), + index: 0, + data_type: DataType::BigInt(crate::spec::BigIntType::new()), + op: crate::spec::PredicateOperator::Eq, + literals: vec![Datum::Long(5)], + }; + let on_key = PredicateBuilder::new(&fields) + .equal("k", Datum::Int(1)) + .unwrap(); + + assert_eq!( + retain_primary_key_conjuncts( + std::slice::from_ref(&row_id), + &fields, + &["k".to_string()] + ), + Vec::new() + ); + assert_eq!( + retain_primary_key_conjuncts( + &[Predicate::and(vec![row_id, on_key.clone()])], + &fields, + &["k".to_string()], + ), + vec![on_key] + ); + } + fn test_file_io() -> FileIO { FileIOBuilder::new("memory").build().unwrap() } diff --git a/crates/paimon/src/table/read_builder.rs b/crates/paimon/src/table/read_builder.rs index de432a5cd..8fed6b7f2 100644 --- a/crates/paimon/src/table/read_builder.rs +++ b/crates/paimon/src/table/read_builder.rs @@ -290,6 +290,10 @@ struct PaimonReadBuilder<'a> { filter: NormalizedFilter, limit: Option, row_ranges: Option>, + /// Whether `row_ranges` was derived from the current filter rather than set + /// by the caller. A derived range must go away when the filter it came from + /// is replaced; an explicit one is the caller's and stays. + row_ranges_from_filter: bool, case_sensitive: bool, } @@ -302,6 +306,7 @@ impl<'a> PaimonReadBuilder<'a> { filter: NormalizedFilter::default(), limit: None, row_ranges: None, + row_ranges_from_filter: false, case_sensitive: true, } } @@ -363,6 +368,12 @@ impl<'a> PaimonReadBuilder<'a> { /// the full predicate with an exact post-merge residual filter. pub fn with_filter(&mut self, filter: Predicate) -> &mut Self { self.filter = normalize_filter(self.table, filter); + // A range derived from the replaced filter goes with it; an explicit one + // is the caller's and stays. + if self.row_ranges_from_filter { + self.row_ranges = None; + self.row_ranges_from_filter = false; + } self.try_extract_row_id_ranges(); self } @@ -386,6 +397,7 @@ impl<'a> PaimonReadBuilder<'a> { } else { Some(ranges) }; + self.row_ranges_from_filter = false; self } @@ -398,12 +410,12 @@ impl<'a> PaimonReadBuilder<'a> { let combined = Predicate::and(self.filter.data_predicates.clone()); if let Some(ranges) = super::row_id_predicate::extract_row_id_ranges(&combined) { self.row_ranges = Some(ranges); - self.filter.data_predicates = self - .filter + self.row_ranges_from_filter = true; + // Ranges are a superset, so a conjunct leaves the residual only when + // they represent it exactly. + self.filter .data_predicates - .iter() - .filter_map(super::row_id_predicate::remove_row_id_filter) - .collect(); + .retain(|conjunct| !super::row_id_predicate::ranges_represent_conjunct(conjunct)); } } @@ -631,11 +643,16 @@ fn projected_read_field_ids_with_predicates( table_fields: &[DataField], ) -> Option> { let mut field_ids = projected_read_field_ids(read_type)?; - let mut predicate_indices = Vec::new(); + let mut refs = Vec::new(); for predicate in predicates { - crate::arrow::residual::collect_predicate_field_indices(predicate, &mut predicate_indices); + crate::arrow::residual::collect_predicate_leaf_refs(predicate, &mut refs); } - for index in predicate_indices { + for (name, index) in refs { + // Not in `table_fields`, and excluded from this set anyway — files never + // list it in `write_cols`. + if crate::spec::is_row_id_column(name) { + continue; + } let Some(field) = table_fields.get(index) else { // A malformed predicate must not make scan planning discard files. return None; @@ -881,6 +898,137 @@ mod tests { ); } + use crate::table::RowRange; + + fn row_id_leaf(op: crate::spec::PredicateOperator, v: i64) -> Predicate { + Predicate::Leaf { + column: crate::spec::ROW_ID_FIELD_NAME.to_string(), + index: 0, + data_type: DataType::BigInt(crate::spec::BigIntType::new()), + op, + literals: vec![crate::spec::Datum::Long(v)], + } + } + + #[test] + fn test_replacing_a_filter_drops_the_range_it_derived() { + use crate::spec::PredicateOperator; + let table = simple_table(); + let mut builder = PaimonReadBuilder::new(&table); + + builder.with_filter(row_id_leaf(PredicateOperator::Eq, 10)); + assert_eq!(builder.row_ranges, Some(vec![RowRange::new(10, 10)])); + + builder.with_filter( + PredicateBuilder::new(table.schema().fields()) + .equal("id", crate::spec::Datum::Int(7)) + .unwrap(), + ); + assert_eq!(builder.row_ranges, None); + } + + #[test] + fn test_an_explicit_range_survives_a_filter_replacement() { + use crate::spec::PredicateOperator; + let table = simple_table(); + let mut builder = PaimonReadBuilder::new(&table); + + builder.with_row_ranges(vec![RowRange::new(0, 20)]); + builder.with_filter(row_id_leaf(PredicateOperator::Eq, 10)); + assert_eq!(builder.row_ranges, Some(vec![RowRange::new(0, 20)])); + builder.with_filter( + PredicateBuilder::new(table.schema().fields()) + .equal("id", crate::spec::Datum::Int(7)) + .unwrap(), + ); + assert_eq!(builder.row_ranges, Some(vec![RowRange::new(0, 20)])); + } + + #[test] + fn test_inexactly_extracted_row_id_conjuncts_stay_residuals() { + use crate::spec::PredicateOperator; + let table = simple_table(); + let mut builder = PaimonReadBuilder::new(&table); + + let contradiction = Predicate::and(vec![ + row_id_leaf(PredicateOperator::Eq, 1), + row_id_leaf(PredicateOperator::NotEq, 1), + ]); + builder.with_filter(contradiction); + assert_eq!(builder.row_ranges, Some(vec![RowRange::new(1, 1)])); + assert_eq!( + builder.filter.data_predicates, + vec![row_id_leaf(PredicateOperator::NotEq, 1)] + ); + + let mut exact = PaimonReadBuilder::new(&table); + exact.with_filter(row_id_leaf(PredicateOperator::GtEq, 10)); + assert_eq!(exact.row_ranges, Some(vec![RowRange::new(10, i64::MAX)])); + assert!(exact.filter.data_predicates.is_empty()); + } + + #[test] + fn test_a_mixed_conjunct_is_never_replaced_by_ranges() { + use crate::spec::{Datum, PredicateOperator}; + let table = simple_table(); + let pb = PredicateBuilder::new(table.schema().fields()); + let disjunction = Predicate::or(vec![ + Predicate::and(vec![ + row_id_leaf(PredicateOperator::Eq, 1), + pb.equal("id", Datum::Int(7)).unwrap(), + ]), + pb.equal("dt", Datum::String("x".into())).unwrap(), + ]); + let mut builder = PaimonReadBuilder::new(&table); + builder.with_filter(Predicate::and(vec![ + row_id_leaf(PredicateOperator::GtEq, 10), + disjunction.clone(), + ])); + + assert_eq!(builder.row_ranges, Some(vec![RowRange::new(10, i64::MAX)])); + assert_eq!(builder.filter.data_predicates, vec![disjunction]); + } + + #[test] + fn test_row_id_filter_never_projects_onto_a_partition_key() { + let table = simple_table(); + let row_id = Predicate::Leaf { + column: crate::spec::ROW_ID_FIELD_NAME.to_string(), + index: 0, + data_type: DataType::BigInt(crate::spec::BigIntType::new()), + op: crate::spec::PredicateOperator::GtEq, + literals: vec![crate::spec::Datum::Long(10)], + }; + + let (partition_predicate, data_predicates) = + super::split_scan_predicates(&table, row_id.clone()); + assert_eq!(partition_predicate, None); + assert_eq!(data_predicates, vec![row_id.clone()]); + + assert!(!ReadBuilder::new(&table).is_exact_filter_pushdown(&row_id)); + } + + #[test] + fn test_projected_read_field_ids_ignore_system_predicate_fields() { + let fields = vec![ + DataField::new(1, "id".to_string(), DataType::Int(IntType::new())), + DataField::new(2, "payload".to_string(), DataType::Int(IntType::new())), + ]; + let read_type = Some(vec![fields[1].clone()]); + let row_id = Predicate::Leaf { + column: crate::spec::ROW_ID_FIELD_NAME.to_string(), + index: 0, + data_type: DataType::BigInt(crate::spec::BigIntType::new()), + op: crate::spec::PredicateOperator::Eq, + literals: vec![crate::spec::Datum::Long(1)], + }; + + assert_eq!( + super::projected_read_field_ids_with_predicates(&read_type, &[row_id], &fields), + Some(HashSet::from([2])) + ); + } + #[test] fn test_with_projection_validates_unknown_projection() { // A column that cannot match under any case sensitivity is an obvious diff --git a/crates/paimon/src/table/row_id_predicate.rs b/crates/paimon/src/table/row_id_predicate.rs index 353ed69cf..8316d4b28 100644 --- a/crates/paimon/src/table/row_id_predicate.rs +++ b/crates/paimon/src/table/row_id_predicate.rs @@ -62,32 +62,39 @@ pub(crate) fn extract_row_id_ranges(predicate: &Predicate) -> Option Option { - match predicate { - Predicate::Leaf { column, .. } if column == ROW_ID_FIELD_NAME => None, - Predicate::And(children) => { - let filtered: Vec = - children.iter().filter_map(remove_row_id_filter).collect(); - match filtered.len() { - 0 => None, - 1 => Some(filtered.into_iter().next().unwrap()), - _ => Some(Predicate::and(filtered)), - } +/// Whether the extracted row ranges represent `conjunct` exactly, so it can be +/// dropped from the residual. +/// +/// Only a conjunct built entirely from convertible `_ROW_ID` conditions +/// qualifies. Extraction ignores whatever it cannot convert, so its ranges are +/// merely a superset of the matching rows — sound for pruning, but anything +/// mixed leaves a remainder the ranges do not carry. Dropping such a conjunct +/// widens the predicate: `_ROW_ID >= 10 AND ((_ROW_ID = 1 AND v = 7) OR w = 8)` +/// would lose `_ROW_ID = 1` and return rows with `v = 7` at any row id. +pub(crate) fn ranges_represent_conjunct(conjunct: &Predicate) -> bool { + match conjunct { + Predicate::Leaf { + column, + op, + literals, + .. + } => column == ROW_ID_FIELD_NAME && leaf_to_ranges(*op, literals).is_some(), + Predicate::And(children) | Predicate::Or(children) => { + children.iter().all(ranges_represent_conjunct) } - Predicate::Or(children) => { - let filtered: Vec = - children.iter().filter_map(remove_row_id_filter).collect(); - if filtered.len() != children.len() { - // If any child was entirely _ROW_ID, the OR semantics change; - // conservatively keep the whole OR. - Some(predicate.clone()) - } else { - Some(Predicate::or(filtered)) - } + _ => false, + } +} + +/// Whether any leaf of `predicate` references `_ROW_ID`. +pub(crate) fn references_row_id(predicate: &Predicate) -> bool { + match predicate { + Predicate::Leaf { column, .. } => column == ROW_ID_FIELD_NAME, + Predicate::And(children) | Predicate::Or(children) => { + children.iter().any(references_row_id) } - other => Some(other.clone()), + Predicate::Not(inner) => references_row_id(inner), + Predicate::AlwaysTrue | Predicate::AlwaysFalse => false, } } @@ -179,6 +186,23 @@ mod tests { } } + #[test] + fn test_references_row_id_finds_nested_leaves() { + let other = Predicate::Leaf { + column: "v".to_string(), + index: 0, + data_type: DataType::BigInt(BigIntType::new()), + op: PredicateOperator::Eq, + literals: vec![Datum::Long(7)], + }; + let nested = Predicate::Not(Box::new(Predicate::or(vec![ + row_id_leaf(PredicateOperator::Eq, vec![Datum::Long(1)]), + other.clone(), + ]))); + assert!(references_row_id(&nested)); + assert!(!references_row_id(&other)); + } + fn data_leaf() -> Predicate { Predicate::Leaf { column: "value".to_string(), @@ -240,24 +264,26 @@ mod tests { } #[test] - fn test_remove_row_id_filter_leaf() { - let p = row_id_leaf(PredicateOperator::Eq, vec![Datum::Long(10)]); - assert!(remove_row_id_filter(&p).is_none()); - } - - #[test] - fn test_remove_row_id_filter_and() { - let p = Predicate::and(vec![ - row_id_leaf(PredicateOperator::GtEq, vec![Datum::Long(10)]), + fn test_ranges_represent_only_a_pure_row_id_conjunct() { + assert!(ranges_represent_conjunct(&row_id_leaf( + PredicateOperator::GtEq, + vec![Datum::Long(10)] + ))); + assert!(!ranges_represent_conjunct(&row_id_leaf( + PredicateOperator::NotEq, + vec![Datum::Long(10)] + ))); + assert!(!ranges_represent_conjunct(&data_leaf())); + assert!(!ranges_represent_conjunct(&Predicate::or(vec![ + Predicate::and(vec![ + row_id_leaf(PredicateOperator::Eq, vec![Datum::Long(1)]), + data_leaf(), + ]), data_leaf(), - ]); - let result = remove_row_id_filter(&p).unwrap(); - assert_eq!(result, data_leaf()); - } - - #[test] - fn test_remove_row_id_filter_keeps_non_row_id() { - let p = data_leaf(); - assert_eq!(remove_row_id_filter(&p).unwrap(), data_leaf()); + ]))); + assert!(ranges_represent_conjunct(&Predicate::or(vec![ + row_id_leaf(PredicateOperator::Eq, vec![Datum::Long(1)]), + row_id_leaf(PredicateOperator::Eq, vec![Datum::Long(2)]), + ]))); } } diff --git a/crates/paimon/src/table/stats_filter.rs b/crates/paimon/src/table/stats_filter.rs index 2eb18e556..15d08d18a 100644 --- a/crates/paimon/src/table/stats_filter.rs +++ b/crates/paimon/src/table/stats_filter.rs @@ -23,7 +23,9 @@ use crate::predicate_stats::{ data_leaf_may_match, data_leaf_must_match, missing_field_may_match, missing_field_must_match, predicates_may_match_with_schema, StatsAccessor, }; -use crate::spec::{extract_datum, BinaryRow, DataField, DataFileMeta, DataType, Datum, Predicate}; +use crate::spec::{ + extract_datum, is_row_id_column, BinaryRow, DataField, DataFileMeta, DataType, Datum, Predicate, +}; use std::collections::HashMap; use std::sync::Arc; @@ -400,12 +402,16 @@ fn data_evolution_predicate_may_match( row_count, ), Predicate::Leaf { + column, index, data_type, op, literals, - .. } => { + // `_ROW_ID` has no column stats, so never prune a group on it. + if is_row_id_column(column) { + return true; + } let Some(source) = field_sources.get(*index).copied().flatten() else { return missing_field_may_match(*op, row_count); }; @@ -463,12 +469,16 @@ fn data_evolution_predicate_must_match( row_count, ), Predicate::Leaf { + column, index, data_type, op, literals, - .. } => { + // Stats cannot decide `_ROW_ID`, so it never provably matches. + if is_row_id_column(column) { + return false; + } let Some(source) = field_sources.get(*index).copied().flatten() else { return missing_field_must_match(*op, row_count); }; diff --git a/crates/paimon/src/table/table_scan.rs b/crates/paimon/src/table/table_scan.rs index 632b6ee35..43f8173ca 100644 --- a/crates/paimon/src/table/table_scan.rs +++ b/crates/paimon/src/table/table_scan.rs @@ -4062,6 +4062,35 @@ mod tests { )); } + #[test] + fn test_data_evolution_group_is_not_pruned_by_a_row_id_predicate() { + 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 row_id = Predicate::Leaf { + column: crate::spec::ROW_ID_FIELD_NAME.to_string(), + index: 0, + data_type: DataType::BigInt(crate::spec::BigIntType::new()), + op: crate::spec::PredicateOperator::GtEq, + literals: vec![Datum::Long(100)], + }; + + assert!(data_evolution_group_matches_predicates( + std::slice::from_ref(&file), + std::slice::from_ref(&row_id), + &fields, + )); + assert!(data_evolution_group_matches_predicates( + &[file], + &[Predicate::negate(row_id)], + &fields, + )); + } + #[test] fn test_data_evolution_group_matches_not_prunes_when_inner_must_match() { let fields = int_field(); diff --git a/crates/paimon/src/table/vector_search_builder.rs b/crates/paimon/src/table/vector_search_builder.rs index 82183f381..00d1c81c6 100644 --- a/crates/paimon/src/table/vector_search_builder.rs +++ b/crates/paimon/src/table/vector_search_builder.rs @@ -23,8 +23,8 @@ use crate::lumina::{ is_lumina_index_type, LuminaIndexMeta, LuminaVectorIndexOptions, LuminaVectorMetric, }; use crate::spec::{ - BigIntType, CoreOptions, DataField, DataType, FileKind, GlobalIndexSearchMode, IndexFileMeta, - IndexManifest, IndexManifestEntry, Predicate, ROW_ID_FIELD_ID, ROW_ID_FIELD_NAME, + row_id_data_field, CoreOptions, DataField, DataType, FileKind, GlobalIndexSearchMode, + IndexFileMeta, IndexManifest, IndexManifestEntry, Predicate, ROW_ID_FIELD_NAME, }; use crate::table::bucket_filter::split_partition_and_data_predicates; use crate::table::data_file_reader::DataFileReader; @@ -2039,19 +2039,6 @@ pub(crate) fn reorder_and_strip_position( Ok(vec![projected]) } -/// The `_ROW_ID` field to append to a data-evolution read type so the reader -/// fills each row's global id. Mirrors the field the `DataEvolutionReader` -/// recognizes (Int64 / `BigInt`, nullable): a data file lacking `first_row_id` -/// yields nulls here, which `attach_scores_by_row_id` then fails loud on rather -/// than mis-aligning scores. -fn row_id_data_field() -> DataField { - DataField::new( - ROW_ID_FIELD_ID, - ROW_ID_FIELD_NAME.to_string(), - DataType::BigInt(BigIntType::with_nullable(true)), - ) -} - /// Collect materialized DE rows, join each row's `(rank, score)` by its global /// `_ROW_ID`, reorder to the search rank order, append the `__paimon_search_score` /// column, and drop `_ROW_ID`. Every row must map to a search candidate and the