Skip to content

Native Iceberg sorted-merge scan rebuilds the delete-file cache per file (shared deletes re-read) #6524

Description

@parthchandra

Describe the bug

Native Iceberg reader (native/core/src/execution/operators/iceberg_scan.rs) - on the sorted-merge read path added in #5331, IcebergScanExec is multi-partition — one file
per Spark partition, one execute() call per file, with a SortPreservingMergeExec above. Each
execute_with_tasks builds a fresh ArrowReaderBuilder(...).build(), and that is where
iceberg-rust constructs the CachingDeleteFileLoader. Its delete_filter is the shared state the
crate documents as existing "to allow caching loaded deletes across multiple calls to
load_deletes (e.g., across multiple file scan tasks)."

Because the reader is rebuilt per file, that cache is now per-file. A delete file that applies to
the whole partition is downloaded and parsed once per data file instead of once for the partition.
fill_delete_file_sizes de-dups only within a single call, so it also issues one HEAD per data
file for the same delete file.

Rows are correct either way — this is purely I/O and CPU, not a correctness problem.

Measurement

Two data files sharing one positional delete file:

  • Unordered path: bytes_scanned = 2717
  • Ordered (merge) path: bytes_scanned = 4256

The 1539-byte delta is exactly the delete file being read a second time. At the default cap
(spark.comet.scan.icebergNative.sortMerge.maxFilesPerPartition = 64) this is up to ~64x on a
merge-on-read table, and equality deletes are the worst case since they are always partition-scoped
and the expensive ones to parse.

Suggested fix

Same as in file_io : hold one ArrowReader onIcebergScanExec and clone it per partition so the CachingDeleteFileLoader (and its delete_filter) is shared across the partition's files. It needs batch_size off the TaskContext, so it has to be a OnceLock filled on first execute rather than built in new.
ArrowReader is Clone, and read() calls ScanMetrics::new() per invocation and re-bases the
loader's metrics through with_scan_metrics, so per-partition metrics stay separate while the
delete cache is shared. With that prototype the ordered path drops back to bytes_scanned = 2717
with identical rows.

Additional context

Metadata

Metadata

Assignees

Labels

bugSomething isn't working

Type

No type

Projects

No projects

    Milestone

    No milestone

    Relationships

    None yet

    Development

    No branches or pull requests

    Issue actions