parthchandra opened a new issue, #6524:
URL: https://github.com/apache/datafusion-comet/issues/6524
### 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` on`IcebergScanExec` 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.
### Notes
- Reported by @andygrove during review of #5331.
- Distinct from #5343 (which is about bounding how many files the merge
opens at once); this is
about the delete-loader being rebuilt per `execute`. Filing separately so
the regression shape and
the measurement are not lost.
### Steps to reproduce
_No response_
### Expected behavior
_No response_
### Additional context
_No response_
--
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]