dwsmith1983 opened a new pull request, #6416:
URL: https://github.com/apache/datafusion-comet/pull/6416

   ## Which issue does this PR close?
   
   Closes #5879.
   
   ## Rationale for this change
   
   Native plans publish the absolute metric values of their own plan, and 
`CometMetricNode.set` wrote them into the shared `SQLMetric` with `set`. When 
one task runs several native plans over the same metric tree, which a coalesce 
without a shuffle does over a Comet RDD (one native plan per parent partition), 
each plan overwrote the one before it and the task reported only the last plan. 
A `COALESCE(1)` over four files showed 2500 of 10000 scanned rows and one 
file's bytes, and the task input metrics derived from those values were short 
by the same amount. Spill counters went through the same path and were 
overwritten the same way.
   
   Spark adds every coalesced partition into the same metric: 
`FileSourceScanExec` does `numOutputRows += batch.numRows()` for each batch the 
task produces, and `FileScanRDD` increments `recordsRead` per batch and folds 
the bytes of earlier partitions into `bytesRead` (SPARK-13071). A coalesced 
query should therefore report the same totals as the query without the coalesce.
   
   ## What changes are included in this PR?
   
   - `CometMetricNode.newInstance()` returns a copy of the tree that updates 
the same `SQLMetric`s but keeps its own record of the last value each metric 
reported. `CometExecIterator` hands one copy to every `Native.createPlan`.
   - `CometMetricNode.set` adds the increase since that instance's last report 
instead of replacing the value, so periodic updates within one plan are not 
counted twice and plans that run one after another in a task add up. A value 
below an earlier report adds nothing.
   - `peak_mem_used` and `build_mem_used` are high-water marks, so they keep 
the maximum across the task's plans rather than the sum.
   - The metrics page of the user guide notes that native metrics accumulate 
per task and names the two metrics that keep the maximum.
   
   ## How are these changes tested?
   
   - Unit tests in `CometTaskMetricsSuite` feed native-style metric updates 
through two instances of one tree and check that counters add up, that peak 
memory keeps the maximum, that a first reported zero marks a size metric as 
set, and that a dip followed by a recovery is counted once.
   - An end-to-end test runs `COALESCE(1)` over a four-file table with one scan 
partition per file, with and without the periodic metrics update, and checks 
that the scan's `output_rows` and `bytes_scanned` and the task input metrics 
match the uncoalesced query and Spark's own coalesced result. It fails without 
the fix.
   - A second end-to-end test covers native blocks that feed each other through 
a JVM input (a project over a coalesce, an aggregate over a union, and an 
aggregate over a shuffle) and checks that each operator's rows are counted once.
   - `CometExecIteratorLifecycleSuite` keeps its throwing metric node on the 
copied tree so the injected failure still fires.
   - The suites pass on Spark 3.5 and 4.1.
   


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