mbutrovich commented on code in PR #2961:
URL: https://github.com/apache/iceberg-rust/pull/2961#discussion_r4125703885
##########
crates/iceberg/src/arrow/reader/row_filter.rs:
##########
@@ -62,8 +67,131 @@ impl ArrowReader {
// creates the projection mask for the Arrow predicates.
let projection_mask = ProjectionMask::leaves(parquet_schema,
column_indices.clone());
let predicate_func = visit(&mut converter, predicates)?;
- let arrow_predicate = ArrowPredicateFn::new(projection_mask,
predicate_func);
- Ok(RowFilter::new(vec![Box::new(arrow_predicate)]))
+ Ok(Box::new(ArrowPredicateFn::new(
+ projection_mask,
+ predicate_func,
+ )))
+ }
+
+ /// Builds one Arrow row-filter predicate per equality-delete set. The
predicate is based
+ /// on a hash-set lookup (see `EqDeleteSet`). It keeps a row unless its
key tuple is present
+ /// in that set. A row is deleted when it matches any set (the predicates
are AND-ed by the `RowFilter`).
+ pub(super) fn build_equality_delete_predicates(
+ sets: &[Arc<EqDeleteSet>],
+ parquet_schema: &SchemaDescriptor,
+ arrow_schema: &ArrowSchemaRef,
+ use_position_fallback: bool,
+ ) -> Result<Vec<Box<dyn ArrowPredicate>>> {
+ let field_id_map =
+ Self::resolve_field_id_map(parquet_schema, arrow_schema,
use_position_fallback)?;
+
+ let mut predicates: Vec<Box<dyn ArrowPredicate>> = Vec::new();
+ for set in sets {
+ if set.is_empty() {
+ continue;
+ }
+
+ // Parquet leaf index for each key column, in `fields` order; a
column dropped
+ // from this file by schema evolution has no entry.
+ let leaf_indices: Vec<Option<usize>> = set
+ .fields
+ .iter()
+ .map(|(_, id, _)| field_id_map.get(id).copied())
+ .collect();
+
+ let mut column_indices: Vec<usize> =
leaf_indices.iter().flatten().copied().collect();
+ column_indices.sort_unstable();
+ column_indices.dedup();
+ let projection_mask = ProjectionMask::leaves(parquet_schema,
column_indices.clone());
+
+ // Position of each key column within the projected batch
(parquet-rs presents the
+ // masked leaves in ascending leaf-index order).
+ let batch_positions: Vec<Option<usize>> = leaf_indices
+ .iter()
+ .map(|leaf| leaf.and_then(|idx|
column_indices.binary_search(&idx).ok()))
+ .collect();
+
+ let target_types: Vec<Type> = set.fields.iter().map(|(_, _, ty)|
ty.clone()).collect();
+ let num_cols = set.fields.len();
+ let set = set.clone();
+
+ let predicate_func =
+ move |batch: RecordBatch| -> std::result::Result<BooleanArray,
ArrowError> {
+ let num_rows = batch.num_rows();
+
+ // Change each key column into `Datum`s once, promoting to
the
+ // table type so the keys match the parsed delete keys
under schema
+ // evolution. A column absent from this file reads as
all-null.
+ let mut columns: Vec<Vec<Option<Datum>>> =
Vec::with_capacity(num_cols);
+ for (i, target_type) in target_types.iter().enumerate() {
+ let Some(pos) = batch_positions[i] else {
+ columns.push(vec![None; num_rows]);
+ continue;
+ };
Review Comment:
> it is spec divergent, it needs to be made obvious, and it needs to be
fixed (but maybe not on this PR). I think this would be worth filing an issue,
even if this PR does not get accepted.
The comment on `test_eq_delete_on_column_absent_from_data_file` makes the
divergence visible. Could you file that issue and put its number in the test
comment, so whoever fixes it finds the test?
##########
crates/iceberg/src/arrow/delete_filter.rs:
##########
@@ -163,68 +162,74 @@ impl DeleteFilter {
}
}
- /// Retrieve the equality delete predicate for a given eq delete file path
- pub(crate) async fn get_equality_delete_predicate_for_delete_file_path(
+ /// Retrieve the equality delete set for a given eq delete file path
+ pub(crate) async fn get_equality_delete_set_for_delete_file_path(
&self,
file_path: &str,
- ) -> Option<Predicate> {
+ ) -> Option<Arc<EqDeleteSet>> {
let notifier = {
match self.state.read().unwrap().equality_deletes.get(file_path) {
None => return None,
Some(EqDelState::Loading(notifier)) => notifier.clone(),
- Some(EqDelState::Loaded(predicate)) => {
- return Some(predicate.clone());
+ Some(EqDelState::Loaded(set)) => {
+ return Some(set.clone());
}
}
};
notifier.notified().await;
match self.state.read().unwrap().equality_deletes.get(file_path) {
- Some(EqDelState::Loaded(predicate)) => Some(predicate.clone()),
+ Some(EqDelState::Loaded(set)) => Some(set.clone()),
_ => unreachable!("Cannot be any other state than loaded"),
}
}
- /// Builds eq delete predicate for the provided task.
- pub(crate) async fn build_equality_delete_predicate(
+ /// Builds the equality-delete sets applicable to the given task, one per
distinct
+ /// equality-column layout.
+ pub(crate) async fn build_equality_delete_sets(
&self,
file_scan_task: &FileScanTask,
- ) -> Result<Option<BoundPredicate>> {
- // * Filter the task's deletes into just the Equality deletes
- // * Retrieve the unbound predicate for each from
self.state.equality_deletes
- // * Logical-AND them all together to get a single combined `Predicate`
- // * Bind the predicate to the task's schema to get a `BoundPredicate`
-
- let mut combined_predicate = AlwaysTrue;
+ ) -> Result<Vec<Arc<EqDeleteSet>>> {
+ let mut groups: HashMap<Vec<i32>, Vec<Arc<EqDeleteSet>>> =
HashMap::new();
for delete in &file_scan_task.deletes {
if !is_equality_delete(delete) {
continue;
}
- let Some(predicate) = self
-
.get_equality_delete_predicate_for_delete_file_path(&delete.file_path)
+ let Some(set) = self
+
.get_equality_delete_set_for_delete_file_path(&delete.file_path)
.await
else {
return Err(Error::new(
ErrorKind::Unexpected,
format!(
- "Missing predicate for equality delete file '{}'",
+ "Missing equality delete set for delete file '{}'",
delete.file_path
),
));
};
- combined_predicate = combined_predicate.and(predicate);
+ let layout = set.fields.iter().map(|(_, id, _)| *id).collect();
+ groups.entry(layout).or_default().push(set);
}
- if combined_predicate == AlwaysTrue {
- return Ok(None);
+ let mut result = Vec::with_capacity(groups.len());
+ for mut sets in groups.into_values() {
+ if sets.len() == 1 {
+ result.push(sets.pop().unwrap());
+ } else {
+ let mut combined = (*sets[0]).clone();
+ for other in &sets[1..] {
+ // `union` checks if `other`s' layout matches `combined`,
+ // which is currently always the case. This fails should a
change
+ // break this current invariant.
+ combined.union(other)?;
+ }
+ result.push(Arc::new(combined));
+ }
}
Review Comment:
> `isInDeleteSets` lives in a `DeleteFilter` instance, which is constructed
per data file. I might be looking at the wrong thing, however?
You're right, I had that wrong. `RowDataReader` builds a new
`SparkDeleteFilter` for every task
([RowDataReader.java#L98-L99](https://github.com/apache/iceberg/blob/212aa894e919ff9fdf89160da4aadcbe6858ac01/spark/v4.0/spark/src/main/java/org/apache/iceberg/spark/source/RowDataReader.java#L98-L99)),
and `BaseDeleteLoader.loadEqualityDeletes` builds a fresh `StructLikeSet` from
the cached per-file rows on every call
([BaseDeleteLoader.java#L100-L105](https://github.com/apache/iceberg/blob/212aa894e919ff9fdf89160da4aadcbe6858ac01/data/src/main/java/org/apache/iceberg/data/BaseDeleteLoader.java#L100-L105)).
So merging per task matches Java.
I also measured the heap with a counting allocator to check the per-task
copy isn't a memory regression. For 100k single-column `long` keys, the cached
`EqDeleteSet` retains 8.1 MB (80 B/key). The predicate tree on `main` retained
38.8 MB (387 B/key), and each task cloned and bound it for another 61.4 MB.
With a `long` plus a 10-character string key, the numbers are 23.5 MB for the
set against 78.4 MB for the tree and 124.5 MB per task on `main`. The union
copy here is smaller than what it replaces, and a task with one delete file per
layout now just clones an `Arc`.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]