Skip to content

Commit ae32e8f

Browse files
authored
[controller] Add OTel metrics to SystemStoreHealthCheckStats and ProtocolVersionAutoDetectionStats (linkedin#2529)
## SystemStoreHealthCheckStats (2 OTel metrics from 3 Tehuti sensors) - `SYSTEM_STORE_UNHEALTHY_COUNT` (`ASYNC_GAUGE`) — consolidates `bad_meta_system_store_count` and `bad_push_status_system_store_count` into a single OTel metric, differentiated by `VENICE_SYSTEM_STORE_TYPE` dimension. Uses `AsyncMetricEntityStateOneEnum<VeniceSystemStoreType>` with callbacks reading existing `AtomicLong` counters. - `SYSTEM_STORE_UNREPAIRABLE_COUNT` (`ASYNC_GAUGE`) — maps 1:1 to `not_repairable_system_store_count`. Uses `AsyncMetricEntityStateBase`. - Tehuti and OTel are registered separately because: (1) multiple Tehuti sensors map to a single OTel metric differentiated by dimension, and (2) `AsyncMetricEntityStateOneEnum` only supports OTel registration, not Tehuti. - OTel switch throws `IllegalArgumentException` for unmapped `VeniceSystemStoreType` values to fail fast if the enum is extended. --- ## ProtocolVersionAutoDetectionStats (2 OTel metrics from 2 Tehuti sensors) - `PROTOCOL_VERSION_AUTO_DETECTION_FAILURE_COUNT` (`GAUGE`) — combined Tehuti Gauge + OTel via `MetricEntityStateBase`. - `PROTOCOL_VERSION_AUTO_DETECTION_TIME` (`MIN_MAX_COUNT_SUM_AGGREGATIONS`) — combined Tehuti Avg + OTel via `MetricEntityStateBase`. --- ## Existing Test Updates — Tehuti + OTel Verification - `TestProtocolVersionAutoDetectionService`: replaced mocked stats with real stats backed by `VeniceMetricsRepository` + `InMemoryMetricReader` in two async tests. Added Tehuti metric existence/value assertions and OTel gauge/histogram validation. - `SystemStoreRepairTaskTest.testRepairSystemStore`: added real `SystemStoreHealthCheckStats` with `doCallRealMethod()` on stats methods. New `verifySystemStoreMetrics` helper validates both Tehuti and OTel at two state transitions.
1 parent 4c79424 commit ae32e8f

8 files changed

Lines changed: 933 additions & 31 deletions

File tree

services/venice-controller/src/main/java/com/linkedin/venice/controller/HelixVeniceClusterResources.java

Lines changed: 1 addition & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -253,9 +253,7 @@ public HelixVeniceClusterResources(
253253
this.protocolVersionAutoDetectionService = new ProtocolVersionAutoDetectionService(
254254
clusterName,
255255
admin,
256-
new ProtocolVersionAutoDetectionStats(
257-
metricsRepository,
258-
"admin_operation_protocol_version_auto_detection_service_" + clusterName),
256+
new ProtocolVersionAutoDetectionStats(metricsRepository, clusterName),
259257
config.getProtocolVersionAutoDetectionSleepMS());
260258
} else {
261259
this.protocolVersionAutoDetectionService = null;

services/venice-controller/src/main/java/com/linkedin/venice/controller/VeniceController.java

Lines changed: 5 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -25,9 +25,11 @@
2525
import com.linkedin.venice.controller.stats.DisabledPartitionStats;
2626
import com.linkedin.venice.controller.stats.ErrorPartitionStats;
2727
import com.linkedin.venice.controller.stats.LogCompactionStats;
28+
import com.linkedin.venice.controller.stats.ProtocolVersionAutoDetectionStats;
2829
import com.linkedin.venice.controller.stats.PushJobStatusStats;
2930
import com.linkedin.venice.controller.stats.SparkServerStats;
3031
import com.linkedin.venice.controller.stats.StoreBackupVersionCleanupServiceStats;
32+
import com.linkedin.venice.controller.stats.SystemStoreHealthCheckStats;
3133
import com.linkedin.venice.controller.stats.TopicCleanupServiceStats;
3234
import com.linkedin.venice.controller.stats.VeniceAdminStats;
3335
import com.linkedin.venice.controller.supersetschema.SupersetSchemaGenerator;
@@ -89,7 +91,9 @@ public class VeniceController {
8991
AddVersionLatencyStats.AddVersionLatencyOtelMetricEntity.class,
9092
DeferredVersionSwapStats.DeferredVersionSwapOtelMetricEntity.class,
9193
DisabledPartitionStats.DisabledPartitionOtelMetricEntity.class,
92-
ErrorPartitionStats.ErrorPartitionOtelMetricEntity.class);
94+
ErrorPartitionStats.ErrorPartitionOtelMetricEntity.class,
95+
SystemStoreHealthCheckStats.SystemStoreHealthCheckOtelMetricEntity.class,
96+
ProtocolVersionAutoDetectionStats.ProtocolVersionAutoDetectionOtelMetricEntity.class);
9397

9498
// services
9599
private final VeniceControllerService controllerService;
Lines changed: 80 additions & 14 deletions
Original file line numberDiff line numberDiff line change
@@ -1,31 +1,97 @@
11
package com.linkedin.venice.controller.stats;
22

3+
import static com.linkedin.venice.stats.dimensions.VeniceMetricsDimensions.VENICE_CLUSTER_NAME;
4+
import static com.linkedin.venice.utils.Utils.setOf;
5+
36
import com.linkedin.venice.stats.AbstractVeniceStats;
7+
import com.linkedin.venice.stats.OpenTelemetryMetricsSetup;
8+
import com.linkedin.venice.stats.VeniceOpenTelemetryMetricsRepository;
9+
import com.linkedin.venice.stats.dimensions.VeniceMetricsDimensions;
10+
import com.linkedin.venice.stats.metrics.MetricEntity;
11+
import com.linkedin.venice.stats.metrics.MetricEntityStateBase;
12+
import com.linkedin.venice.stats.metrics.MetricType;
13+
import com.linkedin.venice.stats.metrics.MetricUnit;
14+
import com.linkedin.venice.stats.metrics.ModuleMetricEntityInterface;
15+
import com.linkedin.venice.stats.metrics.TehutiMetricNameEnum;
16+
import io.opentelemetry.api.common.Attributes;
417
import io.tehuti.metrics.MetricsRepository;
5-
import io.tehuti.metrics.Sensor;
618
import io.tehuti.metrics.stats.Avg;
719
import io.tehuti.metrics.stats.Gauge;
20+
import java.util.Collections;
21+
import java.util.Map;
22+
import java.util.Set;
823

924

1025
public class ProtocolVersionAutoDetectionStats extends AbstractVeniceStats {
11-
private final Sensor protocolVersionAutoDetectionErrorSensor;
12-
private final Sensor protocolVersionAutoDetectionLatencySensor;
13-
private final static String PROTOCOL_VERSION_AUTO_DETECTION_ERROR = "protocol_version_auto_detection_error";
14-
private final static String PROTOCOL_VERSION_AUTO_DETECTION_LATENCY = "protocol_version_auto_detection_latency";
15-
16-
public ProtocolVersionAutoDetectionStats(MetricsRepository metricsRepository, String name) {
17-
super(metricsRepository, name);
18-
protocolVersionAutoDetectionErrorSensor =
19-
registerSensorIfAbsent(PROTOCOL_VERSION_AUTO_DETECTION_ERROR, new Gauge());
20-
protocolVersionAutoDetectionLatencySensor =
21-
registerSensorIfAbsent(PROTOCOL_VERSION_AUTO_DETECTION_LATENCY, new Avg());
26+
private static final String TEHUTI_PREFIX = "admin_operation_protocol_version_auto_detection_service_";
27+
28+
private final MetricEntityStateBase consecutiveFailureMetric;
29+
private final MetricEntityStateBase detectionTimeMetric;
30+
31+
public ProtocolVersionAutoDetectionStats(MetricsRepository metricsRepository, String clusterName) {
32+
super(metricsRepository, TEHUTI_PREFIX + clusterName);
33+
34+
OpenTelemetryMetricsSetup.OpenTelemetryMetricsSetupInfo otelData =
35+
OpenTelemetryMetricsSetup.builder(metricsRepository).setClusterName(clusterName).build();
36+
VeniceOpenTelemetryMetricsRepository otelRepository = otelData.getOtelRepository();
37+
Map<VeniceMetricsDimensions, String> baseDimensionsMap = otelData.getBaseDimensionsMap();
38+
Attributes baseAttributes = otelData.getBaseAttributes();
39+
40+
consecutiveFailureMetric = MetricEntityStateBase.create(
41+
ProtocolVersionAutoDetectionOtelMetricEntity.PROTOCOL_VERSION_AUTO_DETECTION_FAILURE_COUNT.getMetricEntity(),
42+
otelRepository,
43+
this::registerSensorIfAbsent,
44+
ProtocolVersionAutoDetectionTehutiMetricNameEnum.PROTOCOL_VERSION_AUTO_DETECTION_ERROR,
45+
Collections.singletonList(new Gauge()),
46+
baseDimensionsMap,
47+
baseAttributes);
48+
49+
detectionTimeMetric = MetricEntityStateBase.create(
50+
ProtocolVersionAutoDetectionOtelMetricEntity.PROTOCOL_VERSION_AUTO_DETECTION_TIME.getMetricEntity(),
51+
otelRepository,
52+
this::registerSensorIfAbsent,
53+
ProtocolVersionAutoDetectionTehutiMetricNameEnum.PROTOCOL_VERSION_AUTO_DETECTION_LATENCY,
54+
Collections.singletonList(new Avg()),
55+
baseDimensionsMap,
56+
baseAttributes);
2257
}
2358

2459
public void recordProtocolVersionAutoDetectionErrorSensor(int count) {
25-
protocolVersionAutoDetectionErrorSensor.record(count);
60+
consecutiveFailureMetric.record(count);
2661
}
2762

2863
public void recordProtocolVersionAutoDetectionLatencySensor(double latencyInMs) {
29-
protocolVersionAutoDetectionLatencySensor.record(latencyInMs);
64+
detectionTimeMetric.record(latencyInMs);
65+
}
66+
67+
enum ProtocolVersionAutoDetectionTehutiMetricNameEnum implements TehutiMetricNameEnum {
68+
PROTOCOL_VERSION_AUTO_DETECTION_ERROR, PROTOCOL_VERSION_AUTO_DETECTION_LATENCY
69+
}
70+
71+
public enum ProtocolVersionAutoDetectionOtelMetricEntity implements ModuleMetricEntityInterface {
72+
PROTOCOL_VERSION_AUTO_DETECTION_FAILURE_COUNT(
73+
"protocol_version_auto_detection.consecutive_failure_count", MetricType.GAUGE, MetricUnit.NUMBER,
74+
"Consecutive failures in protocol version auto-detection", setOf(VENICE_CLUSTER_NAME)
75+
),
76+
PROTOCOL_VERSION_AUTO_DETECTION_TIME(
77+
"protocol_version_auto_detection.time", MetricType.MIN_MAX_COUNT_SUM_AGGREGATIONS, MetricUnit.MILLISECOND,
78+
"Latency of protocol version auto-detection", setOf(VENICE_CLUSTER_NAME)
79+
);
80+
81+
private final MetricEntity metricEntity;
82+
83+
ProtocolVersionAutoDetectionOtelMetricEntity(
84+
String metricName,
85+
MetricType metricType,
86+
MetricUnit unit,
87+
String description,
88+
Set<VeniceMetricsDimensions> dimensionsList) {
89+
this.metricEntity = new MetricEntity(metricName, metricType, unit, description, dimensionsList);
90+
}
91+
92+
@Override
93+
public MetricEntity getMetricEntity() {
94+
return metricEntity;
95+
}
3096
}
3197
}

services/venice-controller/src/main/java/com/linkedin/venice/controller/stats/SystemStoreHealthCheckStats.java

Lines changed: 90 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -1,9 +1,27 @@
11
package com.linkedin.venice.controller.stats;
22

3+
import static com.linkedin.venice.stats.dimensions.VeniceMetricsDimensions.VENICE_CLUSTER_NAME;
4+
import static com.linkedin.venice.stats.dimensions.VeniceMetricsDimensions.VENICE_SYSTEM_STORE_TYPE;
5+
import static com.linkedin.venice.utils.Utils.setOf;
6+
37
import com.linkedin.venice.stats.AbstractVeniceStats;
8+
import com.linkedin.venice.stats.OpenTelemetryMetricsSetup;
9+
import com.linkedin.venice.stats.VeniceOpenTelemetryMetricsRepository;
10+
import com.linkedin.venice.stats.dimensions.VeniceMetricsDimensions;
11+
import com.linkedin.venice.stats.dimensions.VeniceSystemStoreType;
12+
import com.linkedin.venice.stats.metrics.AsyncMetricEntityStateBase;
13+
import com.linkedin.venice.stats.metrics.AsyncMetricEntityStateOneEnum;
14+
import com.linkedin.venice.stats.metrics.MetricEntity;
15+
import com.linkedin.venice.stats.metrics.MetricType;
16+
import com.linkedin.venice.stats.metrics.MetricUnit;
17+
import com.linkedin.venice.stats.metrics.ModuleMetricEntityInterface;
18+
import com.linkedin.venice.stats.metrics.TehutiMetricNameEnum;
19+
import io.opentelemetry.api.common.Attributes;
420
import io.tehuti.metrics.MetricsRepository;
521
import io.tehuti.metrics.Sensor;
622
import io.tehuti.metrics.stats.AsyncGauge;
23+
import java.util.Map;
24+
import java.util.Set;
725
import java.util.concurrent.atomic.AtomicLong;
826

927

@@ -20,16 +38,53 @@ public class SystemStoreHealthCheckStats extends AbstractVeniceStats {
2038

2139
public SystemStoreHealthCheckStats(MetricsRepository metricsRepository, String name) {
2240
super(metricsRepository, name);
41+
42+
// Tehuti and OTel are registered separately because: (1) multiple Tehuti sensors (bad_meta + bad_push_status)
43+
// map to a single OTel metric differentiated by dimension, and (2) AsyncMetricEntityStateOneEnum only supports
44+
// OTel registration, not Tehuti.
2345
badMetaSystemStoreCountSensor = registerSensorIfAbsent(
24-
new AsyncGauge((ignored, ignored2) -> badMetaSystemStoreCounter.get(), "bad_meta_system_store_count"));
46+
new AsyncGauge(
47+
(ignored, ignored2) -> badMetaSystemStoreCounter.get(),
48+
SystemStoreHealthCheckTehutiMetricNameEnum.BAD_META_SYSTEM_STORE_COUNT.getMetricName()));
2549
badPushStatusSystemStoreCountSensor = registerSensorIfAbsent(
2650
new AsyncGauge(
2751
(ignored, ignored2) -> badPushStatusSystemStoreCounter.get(),
28-
"bad_push_status_system_store_count"));
52+
SystemStoreHealthCheckTehutiMetricNameEnum.BAD_PUSH_STATUS_SYSTEM_STORE_COUNT.getMetricName()));
2953
notRepairableSystemStoreCountSensor = registerSensorIfAbsent(
3054
new AsyncGauge(
3155
(ignored, ignored2) -> notRepairableSystemStoreCounter.get(),
32-
"not_repairable_system_store_count"));
56+
SystemStoreHealthCheckTehutiMetricNameEnum.NOT_REPAIRABLE_SYSTEM_STORE_COUNT.getMetricName()));
57+
58+
// OTel setup
59+
OpenTelemetryMetricsSetup.OpenTelemetryMetricsSetupInfo otelData =
60+
OpenTelemetryMetricsSetup.builder(metricsRepository).setClusterName(name).build();
61+
VeniceOpenTelemetryMetricsRepository otelRepository = otelData.getOtelRepository();
62+
Map<VeniceMetricsDimensions, String> baseDimensionsMap = otelData.getBaseDimensionsMap();
63+
Attributes baseAttributes = otelData.getBaseAttributes();
64+
65+
// OTel async gauges
66+
AsyncMetricEntityStateOneEnum.create(
67+
SystemStoreHealthCheckOtelMetricEntity.SYSTEM_STORE_UNHEALTHY_COUNT.getMetricEntity(),
68+
otelRepository,
69+
baseDimensionsMap,
70+
VeniceSystemStoreType.class,
71+
type -> {
72+
switch (type) {
73+
case META_STORE:
74+
return badMetaSystemStoreCounter::get;
75+
case DAVINCI_PUSH_STATUS_STORE:
76+
return badPushStatusSystemStoreCounter::get;
77+
default:
78+
throw new IllegalArgumentException("Unmapped VeniceSystemStoreType for unhealthy count metric: " + type);
79+
}
80+
});
81+
82+
AsyncMetricEntityStateBase.create(
83+
SystemStoreHealthCheckOtelMetricEntity.SYSTEM_STORE_UNREPAIRABLE_COUNT.getMetricEntity(),
84+
otelRepository,
85+
baseDimensionsMap,
86+
baseAttributes,
87+
notRepairableSystemStoreCounter::get);
3388
}
3489

3590
public AtomicLong getBadMetaSystemStoreCounter() {
@@ -43,4 +98,36 @@ public AtomicLong getBadPushStatusSystemStoreCounter() {
4398
public AtomicLong getNotRepairableSystemStoreCounter() {
4499
return notRepairableSystemStoreCounter;
45100
}
101+
102+
enum SystemStoreHealthCheckTehutiMetricNameEnum implements TehutiMetricNameEnum {
103+
BAD_META_SYSTEM_STORE_COUNT, BAD_PUSH_STATUS_SYSTEM_STORE_COUNT, NOT_REPAIRABLE_SYSTEM_STORE_COUNT
104+
}
105+
106+
public enum SystemStoreHealthCheckOtelMetricEntity implements ModuleMetricEntityInterface {
107+
SYSTEM_STORE_UNHEALTHY_COUNT(
108+
"system_store.health_check.unhealthy_count", MetricType.ASYNC_GAUGE, MetricUnit.NUMBER,
109+
"Unhealthy system stores, differentiated by system store type",
110+
setOf(VENICE_CLUSTER_NAME, VENICE_SYSTEM_STORE_TYPE)
111+
),
112+
SYSTEM_STORE_UNREPAIRABLE_COUNT(
113+
"system_store.health_check.unrepairable_count", MetricType.ASYNC_GAUGE, MetricUnit.NUMBER,
114+
"System stores that cannot be repaired", setOf(VENICE_CLUSTER_NAME)
115+
);
116+
117+
private final MetricEntity metricEntity;
118+
119+
SystemStoreHealthCheckOtelMetricEntity(
120+
String metricName,
121+
MetricType metricType,
122+
MetricUnit unit,
123+
String description,
124+
Set<VeniceMetricsDimensions> dimensionsList) {
125+
this.metricEntity = new MetricEntity(metricName, metricType, unit, description, dimensionsList);
126+
}
127+
128+
@Override
129+
public MetricEntity getMetricEntity() {
130+
return metricEntity;
131+
}
132+
}
46133
}

0 commit comments

Comments
 (0)