Skip to content

fix: add native metrics from every plan a task runs instead of keeping the last - #6416

Merged
andygrove merged 7 commits into
apache:mainfrom
dwsmith1983:fix/5879-native-metrics-accumulate
Oct 4, 2026
Merged

andygrove merged 7 commits into
apache:mainfrom
dwsmith1983:fix/5879-native-metrics-accumulate

Conversation

@dwsmith1983

Copy link
Copy Markdown
Contributor

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

…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.
@github-actions github-actions Bot added the bug Something isn't working label Sep 29, 2026
@dwsmith1983

Copy link
Copy Markdown
Contributor Author

@andygrove could you approve a CI run on 5b226f288, and review when you have time? It fixes the coalesced scan metrics in #5879.

@sunchao sunchao left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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 original SQLMetric accumulators.
  • 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_used and build_mem_used retain maxima.
  • Implementation sketch: CometExecIterator copies 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.

@andygrove andygrove left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

LGTM. Thanks @dwsmith1983

@andygrove
andygrove enabled auto-merge October 3, 2026 17:43

@sunchao sunchao left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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 original SQLMetric accumulators 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_used and build_mem_used retain maxima. Completion listeners keep the same accumulator identities.
  • Implementation sketch: CometExecIterator copies 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.1 at 992c806a7e38c2e88bd018aa5774164b0850e1fa, 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 sunchao left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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 shared SQLMetric accumulators 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_used and build_mem_used retain maxima. Task listeners preserve accumulator identity and deduplication.
  • Implementation sketch: The change stays localized to metric bookkeeping and the createPlan argument. 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.1 at 992c806a7e38c2e88bd018aa5774164b0850e1fa, 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.

@andygrove
andygrove added this pull request to the merge queue Oct 4, 2026
Merged via the queue into apache:main with commit ff468b1 Oct 4, 2026
40 checks passed
@dwsmith1983
dwsmith1983 deleted the fix/5879-native-metrics-accumulate branch October 4, 2026 16:53
peterxcli added a commit to peterxcli/datafusion-comet that referenced this pull request Oct 5, 2026
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.
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

bug Something isn't working

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Native metrics from several plan instances in one task overwrite each other, so a coalesced scan reports only its last partition

3 participants