fix: add native metrics from every plan a task runs instead of keeping the last - #6416
Conversation
…g the last Native plans publish absolute metric values, and CometMetricNode.set wrote them into the SQL metric with set. When one task runs several native plans over the same metric tree, as a coalesce without a shuffle does over a Comet RDD, each plan overwrote the one before it, so 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 them were short too. Each native plan now gets its own copy of the metric tree, sharing the same SQL metrics, which remembers the last value that plan reported and adds only the increase, so periodic updates are not counted twice and several plans add up. The peak memory gauges keep the largest value instead.
|
@andygrove could you approve a CI run on |
sunchao
left a comment
There was a problem hiding this comment.
Summary
- Prior state and problem: Native plans overwrote shared SQL metrics, so coalesced tasks reported only their last partition’s contribution.
- Design approach:
newInstance()gives each native plan independent reporting history while sharing the originalSQLMetricaccumulators. - Correctness / compatibility analysis: Checked Spark’s accumulator, coalesce, scan-row and input-byte semantics across 3.4.3, 3.5.9, 4.0.4, 4.1.3 and 4.2.0. No introduced P1/P2 issues found within this review.
- Key design decisions: Counters accumulate positive deltas. Repeated or decreasing snapshots add nothing.
peak_mem_usedandbuild_mem_usedretain maxima. - Implementation sketch:
CometExecIteratorcopies the metrics tree once per native plan. Bookkeeping scales with tree nodes and reported metric keys, without adding JNI calls or copying batch data. Existing task listeners retain the same accumulator identities. - Behavioral changes worth calling out: Metrics now include every coalesced partition. Documentation and regression tests cover accumulation, nested native blocks and preservation of the lifecycle test’s injected failure.
- Suggested improvements: None meet the requested P1/P2 threshold.
Reviewed the entire five-file diff from 961fbfe768e034a339c0090d99bfc5402e6e7dad to 8fe2af29fdcb5782d1e6ecbf1d52bdf775374dce. The PR is not a draft. Routed skills: review-comet-pr and review-comet-ffi-pr. Existing discussion contains no unresolved substantiated P1/P2 concerns.
Validation: Freshly compiled the exact-head CometMetricNode and regenerated protobuf in a disposable harness using real Spark 3.5.9 and 4.1.3 dependencies, with unrelated plan-type stubs. Four focused tests passed on each version, covering snapshots, interleaved plans, both memory gauges and task-listener totals. The four-plan reproduction reported 2,500 rows at base and 10,000 at head.
Exact-head CI: labeling passed. Comet CI and CodeQL remain action_required, with no completed build/test verdict. Local validation did not include a native build, the full integration/lifecycle suites or performance benchmarks. Project files remain unchanged.
sunchao
left a comment
There was a problem hiding this comment.
Summary
- Prior state and problem: Native plans overwrote shared SQL metrics, so a coalesced task reported only its last partition’s contribution.
- Design approach:
newInstance()shares the originalSQLMetricaccumulators while giving each native plan independent reporting history. - Correctness / compatibility analysis: Checked Spark’s accumulator, coalesce, scan-row and input-byte semantics across 3.4.3, 3.5.9, 4.0.4, 4.1.3 and 4.2.0. No introduced P1/P2 issues found within this review.
- Key design decisions: Counters add positive deltas. Repeated or decreasing snapshots contribute nothing.
peak_mem_usedandbuild_mem_usedretain maxima. Completion listeners keep the same accumulator identities. - Implementation sketch:
CometExecIteratorcopies the metric tree once per native plan. Additional bookkeeping scales with tree nodes and reported metric keys, without adding JNI calls or copying batch data. The implementation stays localized. - Behavioral changes worth calling out: Compared with
branch-1.1at992c806a7e38c2e88bd018aa5774164b0850e1fa, accumulation across plans is an intended correction. Documentation and regression tests cover coalesced totals, nested native blocks and preservation of the lifecycle test’s injected failure. - Suggested improvements: None meet the requested P1/P2 threshold.
Reviewed the entire five-file diff from fef94f6cd78b18151dff57b7a936798385356de5 to 93dd97180e17c556abc6cbcf5b98ef71e14e1f15. The PR is not a draft. Routed skills: review-comet-pr and review-comet-ffi-pr. Existing reviews and discussion contain no unresolved substantiated P1/P2 concerns.
Validation: Freshly compiled the current CometMetricNode and regenerated protobuf in a disposable harness using real Spark 3.5.9 and 4.1.3 dependencies, with unrelated plan-type stubs. Four focused tests passed on each version, covering snapshots, interleaved plans, both memory gauges and task-listener totals. The four-plan reproduction reported 2,500 rows at base and 10,000 at head.
Exact-head CI: Comet CI remains in progress. There are 17 successful checks, including the native build, lint checks and CodeQL, with no failures observed. Rust tests, four Spark 4.1 suite groups and TPC-H/TPC-DS verification remain running. Fourteen checks were skipped, including Spark SQL, Iceberg and macOS jobs. Local validation did not include a native build, full integration/lifecycle suites or performance benchmarks. Project files remain unchanged.
sunchao
left a comment
There was a problem hiding this comment.
Summary
- Prior state and problem: Native plans overwrote shared SQL metrics, so coalesced tasks reported only their last partition’s contribution.
- Design approach:
newInstance()preserves the sharedSQLMetricaccumulators while giving each native plan independent reporting history. - Correctness / compatibility analysis: Checked Spark’s accumulator, coalesce, scan-row and input-byte semantics across 3.4.3, 3.5.9, 4.0.4, 4.1.3 and 4.2.0. No introduced P1/P2 issues found within this review.
- Key design decisions: Counters add positive deltas. Repeated or decreasing snapshots contribute nothing.
peak_mem_usedandbuild_mem_usedretain maxima. Task listeners preserve accumulator identity and deduplication. - Implementation sketch: The change stays localized to metric bookkeeping and the
createPlanargument. It adds one tree copy per native plan and history per reported metric, without additional JNI calls or batch copies. - Behavioral changes worth calling out: Compared with
branch-1.1at992c806a7e38c2e88bd018aa5774164b0850e1fa, accumulating every plan’s contribution is an intended reporting correction. Documentation and regression tests cover coalesced totals, nested native blocks and preservation of the lifecycle test’s injected failure. - Suggested improvements: None meet the requested P1/P2 threshold.
Reviewed the entire five-file diff from fef94f6cd78b18151dff57b7a936798385356de5 to 93dd97180e17c556abc6cbcf5b98ef71e14e1f15. The PR is not a draft. Routed skills: review-comet-pr and review-comet-ffi-pr. Existing reviews and discussion contain no unresolved substantiated P1/P2 concerns.
Validation: Freshly compiled the exact-head CometMetricNode and regenerated protobuf in a disposable harness using real Spark 3.5.9 and 4.1.3 dependencies, with unrelated plan-type stubs. Four focused checks passed on each version, covering snapshots, interleaved plans, both memory gauges and task-listener totals. The four-plan reproduction reported 2,500 rows at base and 10,000 at head.
Exact-head CI: Comet CI remains in progress. There are 17 successful checks, including the native build, lint and CodeQL, with no failures observed. Rust tests, four Spark 4.1 suite groups and TPC-H/TPC-DS verification remain running. Fourteen checks were skipped, including Spark SQL, Iceberg and macOS jobs. Local validation did not include a native build, full integration/lifecycle suites or performance benchmarks. Project files remain unchanged.
Since apache#6416, CometMetricNode reports a sum metric with add instead of set. The test metric threw only from set, so the two cleanup-failure tests in CometNativeWriteSuite stopped injecting their failure and failed on Spark 3.4 and 3.5.
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.setwrote them into the sharedSQLMetricwithset. 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. ACOALESCE(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:
FileSourceScanExecdoesnumOutputRows += batch.numRows()for each batch the task produces, andFileScanRDDincrementsrecordsReadper batch and folds the bytes of earlier partitions intobytesRead(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 sameSQLMetrics but keeps its own record of the last value each metric reported.CometExecIteratorhands one copy to everyNative.createPlan.CometMetricNode.setadds 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_usedandbuild_mem_usedare high-water marks, so they keep the maximum across the task's plans rather than the sum.How are these changes tested?
CometTaskMetricsSuitefeed 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.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'soutput_rowsandbytes_scannedand the task input metrics match the uncoalesced query and Spark's own coalesced result. It fails without the fix.CometExecIteratorLifecycleSuitekeeps its throwing metric node on the copied tree so the injected failure still fires.