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]
