Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
49 changes: 48 additions & 1 deletion crates/paimon/src/arrow/filtering.rs
Original file line number Diff line number Diff line change
Expand Up @@ -17,7 +17,7 @@

use crate::arrow::schema_evolution::create_index_mapping;
pub(crate) use crate::predicate_stats::{predicates_may_match_with_schema, StatsAccessor};
use crate::spec::{DataField, Predicate, PredicateOperator};
use crate::spec::{is_row_id_column, DataField, Predicate, PredicateOperator};

/// Remap predicates from table-level indices to file-level indices.
/// Predicates referencing fields not present in the file are resolved based on
Expand All @@ -44,6 +44,12 @@ fn remap_predicate(predicate: &Predicate, mapping: &[Option<usize>]) -> Predicat
op,
literals,
} => {
// `_ROW_ID` is not a file column and has no per-file position, so
// mapping its placeholder index would collapse the leaf to a
// constant. Keep it; the residual resolves it by name.
if is_row_id_column(column) {
return predicate.clone();
}
match mapping.get(*index).copied().flatten() {
Some(file_index) => Predicate::Leaf {
column: column.clone(),
Expand Down Expand Up @@ -134,3 +140,44 @@ fn normalize_field_mapping(mapping: Option<Vec<i32>>, 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};

#[test]
fn test_a_row_id_leaf_survives_per_file_remapping() {
let table_fields = vec![
DataField::new(1, "added".to_string(), DataType::Int(IntType::new())),
DataField::new(0, "base".to_string(), DataType::Int(IntType::new())),
];
let file_fields = vec![table_fields[1].clone()];
let leaf = crate::spec::row_id_leaf(PredicateOperator::NotEq, vec![Datum::Long(102)]);

assert_eq!(
remap_predicates_to_file(std::slice::from_ref(&leaf), &table_fields, &file_fields),
vec![leaf]
);
}

#[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(),
crate::spec::row_id_leaf(PredicateOperator::Eq, vec![Datum::Long(5)]),
]);

assert_eq!(
remap_predicates_to_file(&[filter], &table_fields, &file_fields),
vec![Predicate::AlwaysTrue]
);
}
}
7 changes: 6 additions & 1 deletion crates/paimon/src/arrow/format/blob.rs
Original file line number Diff line number Diff line change
Expand Up @@ -135,10 +135,15 @@ impl FormatFileReader for BlobFormatReader {
reader: Box<dyn FileRead>,
file_size: u64,
read_fields: &[DataField],
_predicates: Option<&FilePredicates>,
predicates: Option<&FilePredicates>,
batch_size: Option<usize>,
row_selection: Option<Vec<RowRange>>,
) -> crate::Result<ArrowRecordBatchStream> {
// This reader evaluates no predicate at all, so nothing would enforce a
// `_ROW_ID` one.
if let Some(fp) = predicates {
crate::table::row_id_predicate::reject_row_id_filter(&fp.predicates, "blob files")?;
}
let field_kind = validate_read_fields(read_fields)?;

let target_schema = build_target_arrow_schema(read_fields)?;
Expand Down
27 changes: 27 additions & 0 deletions crates/paimon/src/arrow/format/mosaic.rs
Original file line number Diff line number Diff line change
Expand Up @@ -53,6 +53,11 @@ impl FormatFileReader for MosaicFormatReader {
batch_size: Option<usize>,
row_selection: Option<Vec<RowRange>>,
) -> crate::Result<ArrowRecordBatchStream> {
// This reader only prunes row groups by stats, and stats cannot decide a
// `_ROW_ID` predicate, so nothing would enforce it.
if let Some(fp) = predicates {
crate::table::row_id_predicate::reject_row_id_filter(&fp.predicates, "mosaic files")?;
}
let handle = tokio::runtime::Handle::try_current().map_err(|e| Error::UnexpectedError {
message: "Mosaic reader requires a Tokio runtime".to_string(),
source: Some(Box::new(e)),
Expand Down Expand Up @@ -955,6 +960,28 @@ mod tests {
assert_eq!(ids.value(4), 5);
}

#[tokio::test]
async fn test_row_id_predicate_is_rejected() {
let data = write_mosaic(&sample_batch());
let fields = data_fields();
let predicates = FilePredicates {
predicates: vec![crate::spec::row_id_leaf(
crate::spec::PredicateOperator::NotEq,
vec![Datum::Long(101)],
)],
row_filter_factory: None,
file_fields: fields.clone(),
};
let err = read_batches_with_predicates(data, &fields, Some(&predicates), None)
.await
.unwrap_err();

assert!(
matches!(&err, Error::Unsupported { message } if message.contains("_ROW_ID")),
"unexpected error: {err:?}"
);
}

#[tokio::test]
async fn test_read_projection_order() {
let fields = data_fields();
Expand Down
7 changes: 6 additions & 1 deletion crates/paimon/src/arrow/format/orc.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -222,6 +222,7 @@ fn build_orc_leaf_predicate(
file_fields: &[DataField],
) -> Option<orc_rust::predicate::Predicate> {
let Predicate::Leaf {
column,
index,
op,
literals,
Expand All @@ -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();

Expand Down
54 changes: 47 additions & 7 deletions crates/paimon/src/arrow/format/parquet.rs
Original file line number Diff line number Diff line change
Expand Up @@ -25,8 +25,8 @@ use crate::arrow::{ParquetReadBudget, 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;
Expand Down Expand Up @@ -767,11 +767,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()))
})
Expand Down Expand Up @@ -826,11 +826,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)? {
Expand All @@ -853,6 +861,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],
Expand All @@ -861,6 +870,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);
};
Expand Down Expand Up @@ -2243,6 +2257,32 @@ mod tests {
assert!(row_filter.is_some());
}

#[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,
&[crate::spec::row_id_leaf(
super::PredicateOperator::Eq,
vec![Datum::Long(1)]
)],
&fields
)
.expect("row filter should build")
.is_none());

let mixed = Predicate::or(vec![
crate::spec::row_id_leaf(super::PredicateOperator::Eq, vec![Datum::Long(1)]),
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)
// -----------------------------------------------------------------------
Expand Down
5 changes: 5 additions & 0 deletions crates/paimon/src/arrow/format/vortex.rs
Original file line number Diff line number Diff line change
Expand Up @@ -149,6 +149,11 @@ fn read_vortex_batches(
})?;

if scan_fields.is_empty() {
// `_ROW_ID` never widens the scan, so this fast path has nothing to
// evaluate it against and would treat the predicate as matching.
if let Some(fp) = predicates.as_ref() {
crate::table::row_id_predicate::reject_row_id_filter(&fp.predicates, "this read")?;
}
let row_count = if constant_predicates_match(predicates.as_ref()) {
match &row_selection {
Some(ranges) => ranges.iter().map(|r| r.count() as usize).sum(),
Expand Down
Loading
Loading