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]

Reply via email to