Skip to content

[CELEBORN-194] Introduce client side metrics for celeborn - #3740

Open
AmandeepSingh285 wants to merge 18 commits into
apache:mainfrom
AmandeepSingh285:adding-client-side-metrics
Open

AmandeepSingh285 wants to merge 18 commits into
apache:mainfrom
AmandeepSingh285:adding-client-side-metrics

Conversation

@AmandeepSingh285

@AmandeepSingh285 AmandeepSingh285 commented Jun 15, 2026 •

Copy link
Copy Markdown
Contributor

What changes were proposed in this pull request?

Adding client side metrics for Celeborn via heartbeat to master. Adding gauge metrics.

Why are the changes needed?

These changes help increase observability for Celeborn clients.

Does this PR resolve a correctness bug?

  • Yes

Does this PR introduce any user-facing change?

  • Yes

How was this patch tested?

Local setup and unit tests

image

@AmandeepSingh285 AmandeepSingh285 changed the title [CELEBORN-194] WIP Adding client side metrics for celeborn [CELEBORN-194] [WIP] Adding client side metrics for celeborn Jun 15, 2026
@AmandeepSingh285

Copy link
Copy Markdown
Contributor Author

Hi @SteNicholas , @RexXiong could you please help with a high level review on the implementation design for change adding client side metrics.
Thanks!

@codecov

codecov Bot commented Jun 17, 2026 •

Copy link
Copy Markdown

Codecov Report

❌ Patch coverage is 81.86813% with 33 lines in your changes missing coverage. Please review.
✅ Project coverage is 58.77%. Comparing base (17159eb) to head (fe1aa14).
⚠️ Report is 1 commits behind head on main.

Files with missing lines Patch % Lines
.../org/apache/celeborn/client/LifecycleManager.scala 41.38% 15 Missing and 2 partials ⚠️
...born/common/protocol/message/ControlMessages.scala 76.00% 2 Missing and 4 partials ⚠️
...pache/celeborn/client/ChangePartitionManager.scala 0.00% 3 Missing ⚠️
...pache/celeborn/client/ApplicationHeartbeater.scala 77.78% 1 Missing and 1 partial ⚠️
.../apache/celeborn/common/metrics/ClientMetric.scala 50.00% 1 Missing and 1 partial ⚠️
...eleborn/common/metrics/source/AbstractSource.scala 95.66% 1 Missing and 1 partial ⚠️
...n/client/commit/ReducePartitionCommitHandler.scala 0.00% 1 Missing ⚠️
Additional details and impacted files
@@             Coverage Diff              @@
##               main    #3740      +/-   ##
============================================
+ Coverage     57.73%   58.77%   +1.04%     
- Complexity      214      319     +105     
============================================
  Files           397      399       +2     
  Lines         27880    28056     +176     
  Branches       2714     2729      +15     
============================================
+ Hits          16095    16488     +393     
+ Misses        10635    10384     -251     
- Partials       1150     1184      +34     

☔ View full report in Codecov by Harness.
📢 Have feedback on the report? Share it here.

🚀 New features to boost your workflow:
  • ❄️ Test Analytics: Detect flaky tests, report on failures, and find test suite problems.
  • 📦 JS Bundle Analysis: Save yourself from yourself by tracking and limiting bundle sizes in JS merges.

@AmandeepSingh285

Copy link
Copy Markdown
Contributor Author

Gentle ping @SteNicholas , @RexXiong could you please help with a high level review of the approach. Thanks!

@AmandeepSingh285 AmandeepSingh285 changed the title [CELEBORN-194] [WIP] Adding client side metrics for celeborn [CELEBORN-194] Adding client side metrics for celeborn Jun 22, 2026
@AmandeepSingh285

Copy link
Copy Markdown
Contributor Author

Gentle follow-up ping @SteNicholas @RexXiong . Would appreciate a high-level review of the proposed approach whenever you have some time. Thanks!

@AmandeepSingh285 AmandeepSingh285 changed the title [CELEBORN-194] Adding client side metrics for celeborn [CELEBORN-194][WIP] Adding client side metrics for celeborn Jun 22, 2026
@SteNicholas
SteNicholas requested a review from Copilot June 22, 2026 10:13

Copilot AI left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

Pull request overview

Note

Copilot couldn't run its full agentic review because no GitHub Actions runner was available. Make sure your repository has a runner available to run Copilot's review, or add a copilot-setup-steps.yml file specifying one with the runs-on attribute. See the docs for more details.

Adds client-side metrics collection in Celeborn clients and ships those metrics to the master via application heartbeats, where they are re-exposed on the master Prometheus endpoint labeled by applicationId.

Changes:

  • Extend HeartbeatFromApplication (and protobuf serde) to carry a clientMetrics map of {name -> (value, type)}.
  • Add client and master metric sources (CelebornClientSource, ApplicationMetricsSource) plus wiring in LifecycleManager/Master.
  • Introduce celeborn.client.metrics.enabled config and add/unit-test coverage for serde + source behavior.

Reviewed changes

Copilot reviewed 17 out of 17 changed files in this pull request and generated 6 comments.

Show a summary per file
File Description
master/src/test/scala/org/apache/celeborn/service/deploy/master/ApplicationMetricsSourceSuite.scala Adds unit tests for master-side application metrics source behavior.
master/src/main/scala/org/apache/celeborn/service/deploy/master/Master.scala Registers the new application metrics source and plumbs heartbeat clientMetrics through.
master/src/main/scala/org/apache/celeborn/service/deploy/master/ApplicationMetricsSource.scala Implements master-side cache + Prometheus re-export of client metrics by applicationId.
docs/configuration/metrics.md Documents new celeborn.client.metrics.enabled config.
common/src/test/scala/org/apache/celeborn/common/util/UtilsSuite.scala Adds serde round-trip test for clientMetrics in heartbeats.
common/src/main/scala/org/apache/celeborn/common/protocol/message/ControlMessages.scala Extends heartbeat message, protobuf encoding/decoding for client metrics.
common/src/main/scala/org/apache/celeborn/common/metrics/source/AbstractSource.scala Adds Role.CLIENT label behavior and counterExists helper.
common/src/main/scala/org/apache/celeborn/common/metrics/ClientMetric.scala Introduces ClientMetric + MetricType shared representation.
common/src/main/scala/org/apache/celeborn/common/CelebornConf.scala Adds celeborn.client.metrics.enabled config entry and accessor.
common/src/main/proto/TransportMessages.proto Adds clientMetrics field + metric type/message definitions to heartbeat protobuf.
client/src/test/scala/org/apache/celeborn/client/WorkerStatusTrackerSuite.scala Adds test coverage for excluded-worker metrics behavior.
client/src/test/scala/org/apache/celeborn/client/CelebornClientSourceSuite.scala Adds unit tests for client metric source counters/gauges + snapshot types.
client/src/main/scala/org/apache/celeborn/client/commit/ReducePartitionCommitHandler.scala Increments client “shuffle data lost” metric on lost-file conditions when enabled.
client/src/main/scala/org/apache/celeborn/client/LifecycleManager.scala Creates client metrics source, registers gauges, increments counters, and supplies snapshots to heartbeats.
client/src/main/scala/org/apache/celeborn/client/ChangePartitionManager.scala Increments revive-failure metrics when change partition assignment fails.
client/src/main/scala/org/apache/celeborn/client/CelebornClientSource.scala Implements client-side metrics source + snapshot export for heartbeat payload.
client/src/main/scala/org/apache/celeborn/client/ApplicationHeartbeater.scala Adds callback to attach client metrics to each HeartbeatFromApplication.

💡 Add Copilot custom instructions for smarter, more guided reviews. Learn how to get started.

Comment thread client/src/main/scala/org/apache/celeborn/client/LifecycleManager.scala Outdated
Comment thread client/src/main/scala/org/apache/celeborn/client/CelebornClientSource.scala Outdated
Comment thread client/src/test/scala/org/apache/celeborn/client/WorkerStatusTrackerSuite.scala Outdated

@SteNicholas SteNicholas 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.

@AmandeepSingh285, thanks for working on client-side metrics — the overall shape (client AbstractSource snapshot → heartbeat → master re-expose) is reasonable and the serde/config/docs are wired up. Since it's marked [WIP], I'm leaving review comments rather than approving; there are a few correctness/lifecycle issues worth resolving first (inline). Summary, most-impactful first:

  1. Client metrics source + cleaner thread leak (per Spark driver). clientSource is created unconditionally and never destroyed — see inline on LifecycleManager.scala:226.
  2. Master re-registers dead apps' metrics → permanent leak. Master is a plain RpcEndpoint (not ThreadSafeRpcEndpoint), so its Inbox runs with enableConcurrent = true and heartbeats are processed concurrently. A heartbeat racing/after handleAppLost resurrects the app's per-app gauges/counters, which are then never cleaned. See inline on ApplicationMetricsSource.updateApplicationMetrics.
  3. Non-atomic counter delta + fragile absolute→delta conversion. Concurrent heartbeats for one app double-count; app-restart-with-same-id or heartbeat reordering corrupt the delta. See inline on updateCounter.
  4. Unbounded master cardinality + silent truncation. Per-applicationId labeling has no top-N cap and no master-side enable flag; AbstractSource.getMetrics truncates at metricsCapacity (4096) and emits counters last, so app counters drop first. The codebase already solved this for worker per-app metrics via celeborn.metrics.worker.app.topResourceConsumption.count (default 0/off). See inline on Master.scala.
  5. Gauge flaps to 0 on removal (cache cleared before the gauge is unregistered) — inline on removeApplicationMetrics.
  6. Metric semantics: ClientReviveFailCount is incremented by changePartitions.size in one place but by +1 in handleRevive, and ClientShuffleDataLostCount is bumped in both handleMapPartitionEnd and ReducePartitionCommitHandler.stageEnd — mixed units / possible double-count. Inline on ChangePartitionManager.
  7. fromPb silently maps unknown PbMetricType to Gauge (forward-compat trap) — inline on ControlMessages.scala.

Minor / cleanup (no inline needed): the if (clientMetricsEnabled) clientSource.incCounter(...) guard is copy-pasted ~12×, and ReducePartitionCommitHandler recomputes the gate inline instead of reusing the cached clientMetricsEnabled field — a single incClientMetric(name, n) helper would centralize the gate and avoid the semantic drift in (6). ClientMetric/MetricType also duplicate proto PbClientMetric/PbMetricType (4 spots to keep in lockstep). Test gap: the counter-delta path in ApplicationMetricsSource is untested (ApplicationMetricsSourceSuite only sends MetricType.Gauge).

Comment thread client/src/main/scala/org/apache/celeborn/client/LifecycleManager.scala Outdated
Comment thread master/src/main/scala/org/apache/celeborn/service/deploy/master/Master.scala Outdated
Comment thread client/src/main/scala/org/apache/celeborn/client/ChangePartitionManager.scala Outdated
@AmandeepSingh285

Copy link
Copy Markdown
Contributor Author

Thanks @SteNicholas for the review. Still working on improving this PR. Will take into account all the updated you mentioned. Thanks!

@AmandeepSingh285 AmandeepSingh285 changed the title [CELEBORN-194][WIP] Adding client side metrics for celeborn [CELEBORN-194] Adding client side metrics for celeborn Jul 2, 2026
@SteNicholas SteNicholas changed the title [CELEBORN-194] Adding client side metrics for celeborn [CELEBORN-194] Introduce client side metrics for celeborn Jul 2, 2026
@SteNicholas
SteNicholas requested a review from Copilot July 2, 2026 08:30
@SteNicholas

Copy link
Copy Markdown
Member

@AmandeepSingh285, please firstly resolve conflicts.

Copilot AI left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

Pull request overview

Copilot reviewed 19 out of 19 changed files in this pull request and generated 3 comments.

Comment thread client/src/main/scala/org/apache/celeborn/client/CelebornClientSource.scala Outdated
Comment thread client/src/test/scala/org/apache/celeborn/client/WorkerStatusTrackerSuite.scala Outdated

Copilot AI left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

Pull request overview

Copilot reviewed 19 out of 19 changed files in this pull request and generated 6 comments.

Comment thread client/src/test/scala/org/apache/celeborn/client/WorkerStatusTrackerSuite.scala Outdated
Comment thread client/src/main/scala/org/apache/celeborn/client/ApplicationHeartbeater.scala Outdated
Comment thread client/src/main/scala/org/apache/celeborn/client/LifecycleManager.scala Outdated

Copilot AI left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

Pull request overview

Copilot reviewed 19 out of 19 changed files in this pull request and generated 2 comments.

Comment thread common/src/main/scala/org/apache/celeborn/common/CelebornConf.scala

Copilot AI left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

Pull request overview

Copilot reviewed 20 out of 20 changed files in this pull request and generated 1 comment.

Suppressed comments (5)

common/src/main/scala/org/apache/celeborn/common/CelebornConf.scala:6000

  • This config key is master-specific (celeborn.metrics.master.*) but is categorized only under metrics, so it may not appear in master-scoped config docs/lists consistently. Consider categorizing it under both master and metrics for consistency with other master-only configs.
  val MASTER_CLIENT_METRICS_REMOVED_APP_RETENTION: ConfigEntry[Long] =
    buildConf("celeborn.metrics.master.clientMetrics.removedApp.retention")
      .categories("metrics")
      .doc("How long to retain removed application IDs in the client metrics source to " +
        "reject late heartbeats after an application is lost. Entries older than this are " +
        "periodically evicted.")

master/src/main/scala/org/apache/celeborn/service/deploy/master/Master.scala:1257

  • handleHeartbeatFromApplication always converts metricLabels to a Scala Map and calls updateApplicationMetrics even when master-side client metrics are disabled. Since this runs on every application heartbeat, guard the call to avoid per-heartbeat allocations and work when the feature is off (or when labels are empty).
    applicationMetricsSource.updateApplicationMetrics(
      appId,
      metricLabels.asScala.toMap,
      clientMetrics)

client/src/main/scala/org/apache/celeborn/client/ApplicationHeartbeater.scala:56

  • This warning is logged whenever celeborn.client.metrics.enabled is true and no appLabels are set, even if celeborn.metrics.enabled is false (in which case LifecycleManager won't start client metrics at all). Tightening the condition avoids noisy/misleading warnings for configurations where client metrics cannot be emitted.
  if (conf.clientMetricsEnabled && appMetricLabels.isEmpty) {

docs/configuration/master.md:91

  • master.md documents only the cardinality warning threshold, but not the other new master-side client-metrics configs (...clientMetrics.enabled and ...removedApp.retention). Adding them here keeps the master configuration reference complete and prevents users from missing the enable/retention knobs when reading only the master config page.
| celeborn.metrics.master.clientMetrics.seriesCardinality.warnThreshold | 1000 | false | Client metric series are keyed only by their (low-cardinality) label set and are only reclaimed when an application is lost, so a high-cardinality `celeborn.client.metrics.appLabels` configuration can grow the number of distinct series without bound. If the number of tracked series exceeds this threshold, the master logs a one-time warning. | 0.7.0 |  | 

common/src/main/scala/org/apache/celeborn/common/CelebornConf.scala:5990

  • This config key is master-specific (celeborn.metrics.master.*) but is categorized only under metrics, so it may not appear in master-scoped config docs/lists consistently (e.g., compared to ...seriesCardinality.warnThreshold). Consider categorizing it under both master and metrics.

This issue also appears on line 5995 of the same file.

  val MASTER_CLIENT_METRICS_ENABLED: ConfigEntry[Boolean] =
    buildConf("celeborn.metrics.master.clientMetrics.enabled")
      .categories("metrics")
      .doc("When true, the master exposes client-side metrics forwarded in application " +
        "heartbeats on its Prometheus endpoint.")

@SteNicholas

Copy link
Copy Markdown
Member

@AmandeepSingh285, thanks for updates. I left some comments for updates. PTAL.

@AmandeepSingh285

Copy link
Copy Markdown
Contributor Author

Hi @SteNicholas, @RexXiong gentle ping for review.

@AmandeepSingh285

Copy link
Copy Markdown
Contributor Author

Hi @SteNicholas , @RexXiong gentle ping for review

@AmandeepSingh285

Copy link
Copy Markdown
Contributor Author

@SteNicholas PTAL.

@SteNicholas SteNicholas 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.

@AmandeepSingh285, thanks for updates. I have left some comments for updates. BTW, could you please rebase the latest main branch to resolve conflicts?

"reject late heartbeats after an application is lost. Entries older than this are " +
"periodically evicted.")
.version("0.7.0")
.timeConf(TimeUnit.MINUTES)

@SteNicholas SteNicholas Sep 18, 2026 •

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.

Read the retention interval in milliseconds. timeConf(TimeUnit.MINUTES) returns a count of minutes, but masterClientMetricsRemovedAppRetentionMs passes that value unchanged to the millisecond cutoff and scheduleWithFixedDelay(..., TimeUnit.MILLISECONDS). Consequently the default 5min expires tombstones after about 5ms and runs cleanup roughly 200 times per second, allowing late heartbeats to recreate removed metrics. A valid setting such as 30s converts to zero and causes IllegalArgumentException during ApplicationMetricsSource construction, preventing master startup when this feature is enabled. Please use timeConf(TimeUnit.MILLISECONDS) and test unit-bearing configuration strings; the current test's typed Long setter does not catch this mismatch.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Updated

sum,
JavaUtils.newConcurrentHashMap[String, java.lang.Long]())
})
tracked.updateAppValue(appId, value)

@SteNicholas SteNicholas Sep 18, 2026 •

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.

Update the contributing app while holding the per-key map operation. The removal race remains in this gauge-only implementation because computeIfAbsent returns before updateAppValue runs. If app A is the last contributor, app B's heartbeat can obtain its TrackedGauge, then concurrent removal of A can empty perAppValues and deregister/remove that gauge. B subsequently updates the orphaned object, so its accepted value is absent from both the tracked map and metrics registry until its next heartbeat. The final removed-app check does not help because B is still live. I reproduced this interleaving with the added methods. Please perform creation and the app-value update inside the same per-key compute operation used to coordinate with removal, and add a concurrent regression test.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Fixed. Thanks for pointing out

Comment on lines +76 to +79
val reason = rejectReason(appId, metricLabels)
if (reason != null) {
logWarning(s"Ignoring client metrics from $appId: $reason")
return

@SteNicholas SteNicholas Sep 18, 2026 •

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.

Skip normal disabled/empty heartbeats without warning. Master.handleHeartbeatFromApplication calls this method unconditionally, including for older clients and clients with metrics disabled. With the default master configuration, rejectReason therefore emits Ignoring client metrics ... master client metrics are disabled on every application heartbeat. Enabling the master flag still produces no metric labels provided warnings for ordinary clients that have not opted in. This adds continuous warning traffic proportional to the application count even when no client metrics were sent. Please return silently for disabled collection or an empty metrics payload, and reserve warnings for invalid non-empty submissions.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Fixed

@AmandeepSingh285
AmandeepSingh285 force-pushed the adding-client-side-metrics branch from c85cece to 9962fb6 Compare September 19, 2026 13:07
key,
new java.util.function.BiFunction[String, TrackedGauge, TrackedGauge] {
override def apply(k: String, existing: TrackedGauge): TrackedGauge = {
if (isAppRemoved(appId)) {

@SteNicholas SteNicholas Sep 29, 2026 •

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.

First gauge registration can escape application cleanup. For a new series, compute checks isAppRemoved(appId) before publishing the map entry. If removeApplicationMetrics runs after this check but before insertion completes, removeAppFromMetrics takes a keySet snapshot that cannot see the in-progress entry. The heartbeat then publishes and registers the gauge after cleanup has completed, leaving a terminated application's contribution exposed indefinitely.

I reproduced this with a controlled interleaving in an isolated Scala harness using the unchanged gauge methods: removal completed with zero tracked gauges, then the first heartbeat completed with one registered gauge retaining value 7. Please recheck removal after compute returns, or coordinate registration and cleanup per application. A regression test should race removal against the first insertion of an absent series; the current test starts with an existing key and misses this case.

workersAssignedToApp.remove(appId)
statusSystem.handleAppLost(appId, requestId)
quotaManager.handleAppLost(appId)
applicationMetricsSource.removeApplicationMetrics(appId)

@SteNicholas SteNicholas Sep 29, 2026 •

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.

Clean up application metrics across HA leadership changes. This cleanup runs only in the leader's RPC handler, while ApplicationMetricsSource is local to each master and remains registered with its metrics endpoint. After leadership moves from A to B, A keeps exporting its last client gauge values. If an application terminates on B, the replicated AppLost path in MetaHandler only calls updateAppLostMeta; it never removes the contribution from A's metrics source. The stale values can survive even if A becomes leader again, because the application has already been removed from the heartbeat-timeout map.

Please clear/rebuild the source on leadership changes and ensure replicated application removal cleans local metric state, or otherwise reconcile the source with replicated application liveness. A regression test should report metrics to leader A, transfer leadership to B, terminate the application on B, and verify A no longer exports that contribution, including after becoming leader again.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Thanks @SteNicholas for the review. I agree that this could be an issue.
One possible approach to address this would be to persist these metrics through Ratis, ensuring that the values remain consistent across all Master nodes.
What do you think about this approach?

This branch has not been deployed

No deployments
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Projects

None yet

Development

Successfully merging this pull request may close these issues.

4 participants