brgr-s opened a new issue, #2950:
URL: https://github.com/apache/iceberg-rust/issues/2950

   ### Apache Iceberg Rust version
   
   None
   
   ### Describe the bug
   
   Each equality delete file is turned into a predicate expression tree. In
   `crates/iceberg/src/arrow/caching_delete_file_loader.rs`,
   `parse_equality_deletes_record_batch_stream` emits, for every deleted row, 
an OR across that
   row's key columns of `(col IS NULL OR col != v)` — `IS NOT NULL` when the 
delete value is
   null — and then puts all row predicates into an AND-tree.
   
   This results in one leaf per deleted row, evaluated per data row. On an 
unpartitioned table every later equality delete is applicable, so the
   predicate grows with the entire delete history.
   
   If found this when testing a compaction implementation based on PR #2620, 
but the results are unrelated to that particular scenario. Test tabled where 
produced by Flink upserts (one data file + one equality-delete
   file per commit, unpartitioned, so N commits give N·(N−1)/2 delete 
references):
   
   ```
                      data/eq-delete files   rows      refs      read+rewrite 
phase
   small                   120 / 120        56,223     7,140          27,536 ms
   larger                  500 / 500       130,387   124,750         295,268 ms
   ```
   
   ### To Reproduce
   
   Any v2 table on which equality-delete files accumulate reproduces this. The 
table needs to be unpartitioned, or the deletes to share a partition, so that
   each delete file is applicable to every earlier data file.
   
   Take such a table with N commits, each adding one data file and one 
equality-delete file, and scan it:
   
   ```rust
   let mut stream = table.scan().build()?.to_arrow().await?;
   while let Some(_batch) = stream.try_next().await? {}
   ```
   
   Scan wall time grows quadratically in N.
   
   The applicable-delete count that drives it is visible from the plan alone, 
with no instrumentation:
   
   ```rust
   let tasks: Vec<FileScanTask> = table
       .scan()
       .build()?
       .plan_files()
       .await?
       .try_collect()
       .await?;
   
   let refs: usize = tasks.iter().map(|t| t.deletes.len()).sum();
   let max:  usize = tasks.iter().map(|t| t.deletes.len()).max().unwrap_or(0);
   ```
   
   For N = 500 that gives `refs = 124,750 = 500·499/2` and `max = 499` — the 
oldest data file has every later delete applicable. Timing the tasks 
individually shows the cost concentrated on those oldest files, which carry the 
most deletes and contribute the fewest surviving rows.
   
   There is already a (stale) PR that aims to implement Javas' behaviour: PR 
#2343 by @t3hw. The original author states
   
   > I never got around to getting this to the dev mailing list, and my org 
dropped iceberg-rust in favor of a heavier use of spark, so i guess no one is 
actively pushing for it to be merged
   
   I have a slightly different implementation that is inspired by @t3hw 
original work, but instead of filtering after a batch is decoded, I build a 
`RowFilter` with an `ArrowPredicateFn` which is applied during decoding.
   
   The changes amount to +781/-172, so will write to the dev mailing list 
before filing a PR.
   
   ### Expected behavior
   
   The Java implementation keys equality deletes into a `StructLikeSet` and 
checks it per row, which is O(1) per row per key layout. This scales with (data 
rows + delete keys).
   
   ### Willingness to contribute
   
   I can contribute a fix for this bug independently


-- 
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]

Reply via email to