Skip to content

Commit b68623b

Browse files
authored
[da-vinci][server] Add OTel metrics to HeartbeatMonitoringServiceStats (linkedin#2592)
Add dual Tehuti+OTel metric recording to HeartbeatMonitoringServiceStats using the joint MetricEntityStateOneEnum pattern. Two OTel metrics (exception_count, heartbeat_count) consolidate three Tehuti sensors, with a VeniceHeartbeatComponent dimension distinguishing REPORTER vs LOGGER. - New HeartbeatMonitoringOtelMetricEntity enum (2 COUNTER metrics) - New VeniceHeartbeatComponent dimension enum (REPORTER, LOGGER) - Merge duplicate catch(Exception)/catch(Throwable) blocks in HeartbeatMonitoringService into single catch(Throwable) - Add comprehensive OTel tests, enum validation tests, and NPE prevention tests
1 parent 7b6eb7c commit b68623b

16 files changed

Lines changed: 468 additions & 50 deletions

File tree

clients/da-vinci-client/src/main/java/com/linkedin/davinci/DaVinciBackend.java

Lines changed: 4 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -230,8 +230,10 @@ public DaVinciBackend(
230230
cacheBackend = cacheConfig
231231
.map(objectCacheConfig -> new ObjectCacheBackend(clientConfig, objectCacheConfig, schemaRepository));
232232

233-
HeartbeatMonitoringServiceStats heartbeatMonitoringServiceStats =
234-
new HeartbeatMonitoringServiceStats(metricsRepository, "da-vinci");
233+
HeartbeatMonitoringServiceStats heartbeatMonitoringServiceStats = new HeartbeatMonitoringServiceStats(
234+
metricsRepository,
235+
"da-vinci",
236+
configLoader.getVeniceClusterConfig().getClusterName());
235237
heartbeatMonitoringService = new HeartbeatMonitoringService(
236238
metricsRepository,
237239
readOnlyStoreRepository,
Lines changed: 42 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,42 @@
1+
package com.linkedin.davinci.stats;
2+
3+
import static com.linkedin.venice.stats.dimensions.VeniceMetricsDimensions.VENICE_CLUSTER_NAME;
4+
import static com.linkedin.venice.stats.dimensions.VeniceMetricsDimensions.VENICE_HEARTBEAT_COMPONENT;
5+
import static com.linkedin.venice.utils.Utils.setOf;
6+
7+
import com.linkedin.venice.stats.dimensions.VeniceMetricsDimensions;
8+
import com.linkedin.venice.stats.metrics.MetricEntity;
9+
import com.linkedin.venice.stats.metrics.MetricType;
10+
import com.linkedin.venice.stats.metrics.MetricUnit;
11+
import com.linkedin.venice.stats.metrics.ModuleMetricEntityInterface;
12+
import java.util.Set;
13+
14+
15+
public enum HeartbeatMonitoringOtelMetricEntity implements ModuleMetricEntityInterface {
16+
HEARTBEAT_MONITORING_EXCEPTION_COUNT(
17+
"ingestion.heartbeat_monitoring.exception_count", MetricType.COUNTER, MetricUnit.NUMBER,
18+
"Number of exceptions caught in the heartbeat monitoring service threads",
19+
setOf(VENICE_CLUSTER_NAME, VENICE_HEARTBEAT_COMPONENT)
20+
),
21+
HEARTBEAT_MONITORING_HEARTBEAT_COUNT(
22+
"ingestion.heartbeat_monitoring.heartbeat_count", MetricType.COUNTER, MetricUnit.NUMBER,
23+
"Liveness count for the heartbeat monitoring service threads",
24+
setOf(VENICE_CLUSTER_NAME, VENICE_HEARTBEAT_COMPONENT)
25+
);
26+
27+
private final MetricEntity metricEntity;
28+
29+
HeartbeatMonitoringOtelMetricEntity(
30+
String metricName,
31+
MetricType metricType,
32+
MetricUnit unit,
33+
String description,
34+
Set<VeniceMetricsDimensions> dimensions) {
35+
this.metricEntity = new MetricEntity(metricName, metricType, unit, description, dimensions);
36+
}
37+
38+
@Override
39+
public MetricEntity getMetricEntity() {
40+
return metricEntity;
41+
}
42+
}
Lines changed: 82 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -1,34 +1,104 @@
11
package com.linkedin.davinci.stats;
22

3+
import static com.linkedin.davinci.stats.HeartbeatMonitoringOtelMetricEntity.HEARTBEAT_MONITORING_EXCEPTION_COUNT;
4+
import static com.linkedin.davinci.stats.HeartbeatMonitoringOtelMetricEntity.HEARTBEAT_MONITORING_HEARTBEAT_COUNT;
5+
36
import com.linkedin.venice.stats.AbstractVeniceStats;
7+
import com.linkedin.venice.stats.OpenTelemetryMetricsSetup;
8+
import com.linkedin.venice.stats.dimensions.VeniceHeartbeatComponent;
9+
import com.linkedin.venice.stats.dimensions.VeniceMetricsDimensions;
10+
import com.linkedin.venice.stats.metrics.MetricEntityStateOneEnum;
11+
import com.linkedin.venice.stats.metrics.TehutiMetricNameEnum;
412
import io.tehuti.metrics.MetricsRepository;
5-
import io.tehuti.metrics.Sensor;
613
import io.tehuti.metrics.stats.Count;
714
import io.tehuti.metrics.stats.OccurrenceRate;
15+
import java.util.Collections;
16+
import java.util.Map;
817

918

1019
public class HeartbeatMonitoringServiceStats extends AbstractVeniceStats {
1120
private static final String HEARTBEAT_SUFFIX = "-heartbeat-monitor-service";
12-
private final Sensor heartbeatExceptionCountSensor;
13-
private final Sensor heartbeatReporterSensor;
14-
private final Sensor heartbeatLoggingSensor;
1521

16-
public HeartbeatMonitoringServiceStats(MetricsRepository metricsRepository, String heartbeatStatPrefix) {
22+
/**
23+
* Tehuti metric names. Hyphenated names are preserved for backward compatibility with existing
24+
* dashboards — the default {@link TehutiMetricNameEnum#getMetricName()} returns the lowercased
25+
* enum constant name (which uses underscores), so each constant overrides it to preserve the
26+
* original hyphenated names.
27+
*/
28+
enum TehutiMetricName implements TehutiMetricNameEnum {
29+
HEARTBEAT_MONITOR_SERVICE_EXCEPTION_COUNT("heartbeat-monitor-service-exception-count"),
30+
HEARTBEAT_REPORTER("heartbeat-reporter"), HEARTBEAT_LOGGER("heartbeat-logger");
31+
32+
private final String metricName;
33+
34+
TehutiMetricName(String metricName) {
35+
this.metricName = metricName;
36+
}
37+
38+
@Override
39+
public String getMetricName() {
40+
return metricName;
41+
}
42+
}
43+
44+
/**
45+
* 3 joint Tehuti+OTel metric states mapping to 2 OTel metrics. The reporter and logger heartbeat
46+
* states share the same OTel instrument ({@code HEARTBEAT_MONITORING_HEARTBEAT_COUNT}) but bind
47+
* different Tehuti sensors — each {@link MetricEntityStateOneEnum} wraps a distinct Tehuti sensor
48+
* while contributing to the same OTel counter differentiated by the
49+
* {@link VeniceHeartbeatComponent} dimension.
50+
*/
51+
private final MetricEntityStateOneEnum<VeniceHeartbeatComponent> exceptionCountMetrics;
52+
private final MetricEntityStateOneEnum<VeniceHeartbeatComponent> reporterHeartbeatMetrics;
53+
private final MetricEntityStateOneEnum<VeniceHeartbeatComponent> loggerHeartbeatMetrics;
54+
55+
public HeartbeatMonitoringServiceStats(
56+
MetricsRepository metricsRepository,
57+
String heartbeatStatPrefix,
58+
String clusterName) {
1759
super(metricsRepository, heartbeatStatPrefix + HEARTBEAT_SUFFIX);
18-
heartbeatExceptionCountSensor = registerSensorIfAbsent("heartbeat-monitor-service-exception-count", new Count());
19-
heartbeatReporterSensor = registerSensorIfAbsent("heartbeat-reporter", new OccurrenceRate());
20-
heartbeatLoggingSensor = registerSensorIfAbsent("heartbeat-logger", new OccurrenceRate());
60+
61+
OpenTelemetryMetricsSetup.OpenTelemetryMetricsSetupInfo otelData =
62+
OpenTelemetryMetricsSetup.builder(metricsRepository).setClusterName(clusterName).build();
63+
Map<VeniceMetricsDimensions, String> baseDimensionsMap = otelData.getBaseDimensionsMap();
64+
65+
this.exceptionCountMetrics = MetricEntityStateOneEnum.create(
66+
HEARTBEAT_MONITORING_EXCEPTION_COUNT.getMetricEntity(),
67+
otelData.getOtelRepository(),
68+
this::registerSensorIfAbsent,
69+
TehutiMetricName.HEARTBEAT_MONITOR_SERVICE_EXCEPTION_COUNT,
70+
Collections.singletonList(new Count()),
71+
baseDimensionsMap,
72+
VeniceHeartbeatComponent.class);
73+
74+
this.reporterHeartbeatMetrics = MetricEntityStateOneEnum.create(
75+
HEARTBEAT_MONITORING_HEARTBEAT_COUNT.getMetricEntity(),
76+
otelData.getOtelRepository(),
77+
this::registerSensorIfAbsent,
78+
TehutiMetricName.HEARTBEAT_REPORTER,
79+
Collections.singletonList(new OccurrenceRate()),
80+
baseDimensionsMap,
81+
VeniceHeartbeatComponent.class);
82+
83+
this.loggerHeartbeatMetrics = MetricEntityStateOneEnum.create(
84+
HEARTBEAT_MONITORING_HEARTBEAT_COUNT.getMetricEntity(),
85+
otelData.getOtelRepository(),
86+
this::registerSensorIfAbsent,
87+
TehutiMetricName.HEARTBEAT_LOGGER,
88+
Collections.singletonList(new OccurrenceRate()),
89+
baseDimensionsMap,
90+
VeniceHeartbeatComponent.class);
2191
}
2292

23-
public void recordHeartbeatExceptionCount() {
24-
heartbeatExceptionCountSensor.record();
93+
public void recordHeartbeatExceptionCount(VeniceHeartbeatComponent component) {
94+
exceptionCountMetrics.record(1, component);
2595
}
2696

2797
public void recordReporterHeartbeat() {
28-
heartbeatReporterSensor.record();
98+
reporterHeartbeatMetrics.record(1, VeniceHeartbeatComponent.REPORTER);
2999
}
30100

31101
public void recordLoggerHeartbeat() {
32-
heartbeatLoggingSensor.record();
102+
loggerHeartbeatMetrics.record(1, VeniceHeartbeatComponent.LOGGER);
33103
}
34104
}

clients/da-vinci-client/src/main/java/com/linkedin/davinci/stats/ServerMetricEntity.java

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -35,6 +35,7 @@ public static List<Class<? extends ModuleMetricEntityInterface>> getMetricEntity
3535
ServerMetadataOtelMetricEntity.class,
3636
ParticipantStoreConsumptionOtelMetricEntity.class,
3737
AdaptiveThrottlingOtelMetricEntity.class,
38+
HeartbeatMonitoringOtelMetricEntity.class,
3839
BlobTransferOtelMetricEntity.class);
3940
}
4041

clients/da-vinci-client/src/main/java/com/linkedin/davinci/stats/ingestion/heartbeat/HeartbeatMonitoringService.java

Lines changed: 7 additions & 13 deletions
Original file line numberDiff line numberDiff line change
@@ -17,6 +17,7 @@
1717
import com.linkedin.venice.meta.Version;
1818
import com.linkedin.venice.meta.VersionImpl;
1919
import com.linkedin.venice.service.AbstractVeniceService;
20+
import com.linkedin.venice.stats.dimensions.VeniceHeartbeatComponent;
2021
import com.linkedin.venice.utils.LogContext;
2122
import com.linkedin.venice.utils.Utils;
2223
import com.linkedin.venice.utils.concurrent.VeniceConcurrentHashMap;
@@ -877,13 +878,10 @@ public void run() {
877878
} catch (InterruptedException e) {
878879
// We've received an interrupt which is to be expected, so we'll just leave the loop and log
879880
break;
880-
} catch (Exception e) {
881-
exceptionThrown = true;
882-
LOGGER.error("Received exception from Ingestion-Heartbeat-Reporter-Service-Thread", e);
883-
heartbeatMonitoringServiceStats.recordHeartbeatExceptionCount();
884-
} catch (Throwable throwable) {
881+
} catch (Throwable t) {
885882
exceptionThrown = true;
886-
LOGGER.error("Received exception from Ingestion-Heartbeat-Reporter-Service-Thread", throwable);
883+
LOGGER.error("Received exception from Ingestion-Heartbeat-Reporter-Service-Thread", t);
884+
heartbeatMonitoringServiceStats.recordHeartbeatExceptionCount(VeniceHeartbeatComponent.REPORTER);
887885
}
888886
}
889887
LOGGER.info("Heartbeat lag metric reporting thread stopped. Shutting down...");
@@ -916,14 +914,10 @@ public void run() {
916914
} catch (InterruptedException e) {
917915
// We've received an interrupt which is to be expected, so we'll just leave the loop and log
918916
break;
919-
} catch (Exception e) {
920-
exceptionThrown = true;
921-
LOGGER.error("Received exception from Ingestion-Heartbeat-Lag-Logging-Service-Thread", e);
922-
heartbeatMonitoringServiceStats.recordHeartbeatExceptionCount();
923-
} catch (Throwable throwable) {
917+
} catch (Throwable t) {
924918
exceptionThrown = true;
925-
LOGGER
926-
.error("Received non-exception throwable from Ingestion-Heartbeat-Lag-Logging-Service-Thread", throwable);
919+
LOGGER.error("Received exception from Ingestion-Heartbeat-Lag-Logging-Service-Thread", t);
920+
heartbeatMonitoringServiceStats.recordHeartbeatExceptionCount(VeniceHeartbeatComponent.LOGGER);
927921
}
928922
}
929923
LOGGER.info("Heartbeat lag logging thread stopped. Shutting down...");
Lines changed: 44 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,44 @@
1+
package com.linkedin.davinci.stats;
2+
3+
import static com.linkedin.davinci.stats.HeartbeatMonitoringOtelMetricEntity.HEARTBEAT_MONITORING_EXCEPTION_COUNT;
4+
import static com.linkedin.davinci.stats.HeartbeatMonitoringOtelMetricEntity.HEARTBEAT_MONITORING_HEARTBEAT_COUNT;
5+
import static com.linkedin.venice.stats.dimensions.VeniceMetricsDimensions.VENICE_CLUSTER_NAME;
6+
import static com.linkedin.venice.stats.dimensions.VeniceMetricsDimensions.VENICE_HEARTBEAT_COMPONENT;
7+
import static com.linkedin.venice.utils.Utils.setOf;
8+
9+
import com.linkedin.venice.stats.metrics.MetricType;
10+
import com.linkedin.venice.stats.metrics.MetricUnit;
11+
import com.linkedin.venice.stats.metrics.ModuleMetricEntityTestFixture;
12+
import com.linkedin.venice.stats.metrics.ModuleMetricEntityTestFixture.MetricEntityExpectation;
13+
import java.util.HashMap;
14+
import java.util.Map;
15+
import org.testng.annotations.Test;
16+
17+
18+
public class HeartbeatMonitoringOtelMetricEntityTest {
19+
@Test
20+
public void testMetricEntities() {
21+
new ModuleMetricEntityTestFixture<>(HeartbeatMonitoringOtelMetricEntity.class, expectedDefinitions()).assertAll();
22+
}
23+
24+
private static Map<HeartbeatMonitoringOtelMetricEntity, MetricEntityExpectation> expectedDefinitions() {
25+
Map<HeartbeatMonitoringOtelMetricEntity, MetricEntityExpectation> map = new HashMap<>();
26+
map.put(
27+
HEARTBEAT_MONITORING_EXCEPTION_COUNT,
28+
new MetricEntityExpectation(
29+
"ingestion.heartbeat_monitoring.exception_count",
30+
MetricType.COUNTER,
31+
MetricUnit.NUMBER,
32+
"Number of exceptions caught in the heartbeat monitoring service threads",
33+
setOf(VENICE_CLUSTER_NAME, VENICE_HEARTBEAT_COMPONENT)));
34+
map.put(
35+
HEARTBEAT_MONITORING_HEARTBEAT_COUNT,
36+
new MetricEntityExpectation(
37+
"ingestion.heartbeat_monitoring.heartbeat_count",
38+
MetricType.COUNTER,
39+
MetricUnit.NUMBER,
40+
"Liveness count for the heartbeat monitoring service threads",
41+
setOf(VENICE_CLUSTER_NAME, VENICE_HEARTBEAT_COMPONENT)));
42+
return map;
43+
}
44+
}

0 commit comments

Comments
 (0)