[CELEBORN-194] Introduce client side metrics for celeborn - #3740
[CELEBORN-194] Introduce client side metrics for celeborn#3740AmandeepSingh285 wants to merge 15 commits into
Conversation
|
Hi @SteNicholas , @RexXiong could you please help with a high level review on the implementation design for change adding client side metrics. |
Codecov Report❌ Patch coverage is 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. 🚀 New features to boost your workflow:
|
|
Gentle ping @SteNicholas , @RexXiong could you please help with a high level review of the approach. Thanks! |
|
Gentle follow-up ping @SteNicholas @RexXiong . Would appreciate a high-level review of the proposed approach whenever you have some time. Thanks! |
There was a problem hiding this comment.
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 aclientMetricsmap of{name -> (value, type)}. - Add client and master metric sources (
CelebornClientSource,ApplicationMetricsSource) plus wiring inLifecycleManager/Master. - Introduce
celeborn.client.metrics.enabledconfig 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.
There was a problem hiding this comment.
@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:
- Client metrics source + cleaner thread leak (per Spark driver).
clientSourceis created unconditionally and never destroyed — see inline onLifecycleManager.scala:226. - Master re-registers dead apps' metrics → permanent leak.
Masteris a plainRpcEndpoint(notThreadSafeRpcEndpoint), so its Inbox runs withenableConcurrent = trueand heartbeats are processed concurrently. A heartbeat racing/afterhandleAppLostresurrects the app's per-app gauges/counters, which are then never cleaned. See inline onApplicationMetricsSource.updateApplicationMetrics. - 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. - Unbounded master cardinality + silent truncation. Per-
applicationIdlabeling has no top-N cap and no master-side enable flag;AbstractSource.getMetricstruncates atmetricsCapacity(4096) and emits counters last, so app counters drop first. The codebase already solved this for worker per-app metrics viaceleborn.metrics.worker.app.topResourceConsumption.count(default 0/off). See inline onMaster.scala. - Gauge flaps to 0 on removal (cache cleared before the gauge is unregistered) — inline on
removeApplicationMetrics. - Metric semantics:
ClientReviveFailCountis incremented bychangePartitions.sizein one place but by +1 inhandleRevive, andClientShuffleDataLostCountis bumped in bothhandleMapPartitionEndandReducePartitionCommitHandler.stageEnd— mixed units / possible double-count. Inline onChangePartitionManager. fromPbsilently maps unknownPbMetricTypeto Gauge (forward-compat trap) — inline onControlMessages.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).
|
Thanks @SteNicholas for the review. Still working on improving this PR. Will take into account all the updated you mentioned. Thanks! |
|
@AmandeepSingh285, please firstly resolve conflicts. |
36a9603 to
438f541
Compare
ca1b783 to
8b718f1
Compare
There was a problem hiding this comment.
@AmandeepSingh285, thanks for updates. I found three remaining issues in the latest head: ambiguous heartbeat acknowledgements can duplicate counter deltas, the one-time cardinality check still performs a full series scan per heartbeat, and each enabled client starts an unnecessary timer-cleaner thread. Details inline.
There was a problem hiding this comment.
Pull request overview
Copilot reviewed 20 out of 20 changed files in this pull request and generated no new comments.
Comments suppressed due to low confidence (1)
master/src/main/scala/org/apache/celeborn/service/deploy/master/ApplicationMetricsSource.scala:73
metricLabelscome from the client heartbeat and are later rendered into the Prometheus text format viaMetricLabels.labelString(which interpolates keys/values verbatim ask="v"). Without escaping, a label value containing",\, or a newline can break the/metricsoutput (or inject additional lines), and an invalid label key can produce an invalid exposition format. Consider sanitizing/escaping client-provided labels before registering/updating metrics.
def updateApplicationMetrics(
appId: String,
metricLabels: Map[String, String],
metrics: JMap[String, ClientMetric]): Unit = {
if (!masterClientMetricsEnabled || metricLabels.isEmpty) {
There was a problem hiding this comment.
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 undermetrics, so it may not appear in master-scoped config docs/lists consistently. Consider categorizing it under bothmasterandmetricsfor 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
handleHeartbeatFromApplicationalways convertsmetricLabelsto a ScalaMapand callsupdateApplicationMetricseven 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.enabledis true and no appLabels are set, even ifceleborn.metrics.enabledis 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.mddocuments only the cardinality warning threshold, but not the other new master-side client-metrics configs (...clientMetrics.enabledand...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 undermetrics, so it may not appear in master-scoped config docs/lists consistently (e.g., compared to...seriesCardinality.warnThreshold). Consider categorizing it under bothmasterandmetrics.
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.")
| // Atomically get-or-create so the returned NamedCounter is guaranteed to be the instance | ||
| // currently resident in namedCounters. Callers must never re-.get(metricNameWithLabel) | ||
| // afterwards, since a concurrent removeCounter could have removed it in the meantime. | ||
| namedCounters.computeIfAbsent( |
| TrackedGauge(namedGauge, holder, ConcurrentHashMap.newKeySet[String]()) | ||
| } | ||
| tracked.contributingAppIds.add(appId) | ||
| tracked.handle.set(value) |
There was a problem hiding this comment.
Gauges become last-writer-wins across applications sharing a label set. The documented low-cardinality labels (env, team, etc.) intentionally collapse multiple applications into one series, but set(value) makes ClientActiveShuffleCount and the other gauges report whichever application heartbeated last rather than an aggregate. It also leaves a removed application value behind when that application was the last writer and other contributors remain. Please retain per-application gauge values and define an aggregation such as sum/max, or require an application-unique label and document that contract.
| def updateApplicationMetrics( | ||
| appId: String, | ||
| metricLabels: Map[String, String], | ||
| metrics: JMap[String, ClientMetric]): Unit = { |
There was a problem hiding this comment.
Validate heartbeat-provided labels before registering metrics. These labels come from the wire and are eventually rendered as key="value" without escaping; the client-side config check only verifies key/value parsing and can also be bypassed by another client implementation. Invalid keys or values containing quotes, backslashes, or newlines can corrupt the master Prometheus response for every scrape. Please enforce Prometheus label-name rules on the master and correctly escape, or reject, unsafe values.
|
@AmandeepSingh285, thanks for updates. I left some comments for updates. PTAL. |
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?
Does this PR introduce any user-facing change?
How was this patch tested?
Local setup and unit tests