Skip to content

Commit 07bb5d5

Browse files
authored
[da-vinci] Add OTel metrics to StuckConsumerRepairStats (linkedin#2723)
Add 3 OTel COUNTER metrics under the ingestion.pubsub.consumer.stuck.* namespace: - ingestion.pubsub.consumer.stuck.detected_count: scans that detected a stuck PubSub consumer - ingestion.pubsub.consumer.stuck.task_repaired_count: ingestion tasks killed to unblock a stuck consumer - ingestion.pubsub.consumer.stuck.unresolved_count: stuck consumers found with no fixable task identified This is a singleton stats class with CLUSTER_NAME as the only dimension. clusterName was added to the constructor, passed from AggKafkaConsumerService via serverConfig.getClusterName().
1 parent 55a09c7 commit 07bb5d5

8 files changed

Lines changed: 332 additions & 14 deletions

File tree

clients/da-vinci-client/src/main/java/com/linkedin/davinci/kafka/consumer/AggKafkaConsumerService.java

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -119,7 +119,7 @@ public AggKafkaConsumerService(
119119
this.kafkaClusterUrlResolver = serverConfig.getKafkaClusterUrlResolver();
120120
this.metadataRepository = metadataRepository;
121121
if (serverConfig.isStuckConsumerRepairEnabled()) {
122-
this.stuckConsumerStats = new StuckConsumerRepairStats(metricsRepository);
122+
this.stuckConsumerStats = new StuckConsumerRepairStats(metricsRepository, serverConfig.getClusterName());
123123
this.stuckConsumerRepairExecutorService = Executors.newSingleThreadScheduledExecutor(
124124
new DaemonThreadFactory(this.getClass().getName() + "-StuckConsumerRepair", serverConfig.getLogContext()));
125125
int intervalInSeconds = serverConfig.getStuckConsumerRepairIntervalSecond();

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

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -51,7 +51,8 @@ public static List<Class<? extends ModuleMetricEntityInterface>> getMetricEntity
5151
RocksDBStatsOtelMetricEntity.class,
5252
BackupVersionOptimizationOtelMetricEntity.class,
5353
ParticipantStateTransitionOtelMetricEntity.class,
54-
DaVinciRecordTransformerOtelMetricEntity.class);
54+
DaVinciRecordTransformerOtelMetricEntity.class,
55+
StuckConsumerRepairOtelMetricEntity.class);
5556
}
5657

5758
public static final Collection<MetricEntity> SERVER_METRIC_ENTITIES =
Lines changed: 50 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,50 @@
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.utils.Utils.setOf;
5+
6+
import com.linkedin.venice.stats.dimensions.VeniceMetricsDimensions;
7+
import com.linkedin.venice.stats.metrics.MetricEntity;
8+
import com.linkedin.venice.stats.metrics.MetricType;
9+
import com.linkedin.venice.stats.metrics.MetricUnit;
10+
import com.linkedin.venice.stats.metrics.ModuleMetricEntityInterface;
11+
import java.util.Set;
12+
13+
14+
/**
15+
* OTel metric entity definitions for {@link StuckConsumerRepairStats}.
16+
* Metrics track the stuck PubSub consumer detection and repair lifecycle.
17+
*/
18+
public enum StuckConsumerRepairOtelMetricEntity implements ModuleMetricEntityInterface {
19+
STUCK_CONSUMER_DETECTED_COUNT(
20+
"ingestion.pubsub.consumer.stuck.detected_count", MetricType.COUNTER, MetricUnit.NUMBER,
21+
"Count of scans that detected a stuck PubSub consumer", setOf(VENICE_CLUSTER_NAME)
22+
),
23+
24+
STUCK_CONSUMER_TASK_REPAIRED_COUNT(
25+
"ingestion.pubsub.consumer.stuck.task_repaired_count", MetricType.COUNTER, MetricUnit.NUMBER,
26+
"Count of ingestion tasks killed to unblock a stuck PubSub consumer", setOf(VENICE_CLUSTER_NAME)
27+
),
28+
29+
STUCK_CONSUMER_UNRESOLVED_COUNT(
30+
"ingestion.pubsub.consumer.stuck.unresolved_count", MetricType.COUNTER, MetricUnit.NUMBER,
31+
"Count of scans where a stuck PubSub consumer was found but no fixable task identified",
32+
setOf(VENICE_CLUSTER_NAME)
33+
);
34+
35+
private final MetricEntity metricEntity;
36+
37+
StuckConsumerRepairOtelMetricEntity(
38+
String metricName,
39+
MetricType metricType,
40+
MetricUnit unit,
41+
String description,
42+
Set<VeniceMetricsDimensions> dimensions) {
43+
this.metricEntity = new MetricEntity(metricName, metricType, unit, description, dimensions);
44+
}
45+
46+
@Override
47+
public MetricEntity getMetricEntity() {
48+
return metricEntity;
49+
}
50+
}
Lines changed: 54 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -1,33 +1,76 @@
11
package com.linkedin.davinci.stats;
22

3+
import static com.linkedin.davinci.stats.StuckConsumerRepairOtelMetricEntity.STUCK_CONSUMER_DETECTED_COUNT;
4+
import static com.linkedin.davinci.stats.StuckConsumerRepairOtelMetricEntity.STUCK_CONSUMER_TASK_REPAIRED_COUNT;
5+
import static com.linkedin.davinci.stats.StuckConsumerRepairOtelMetricEntity.STUCK_CONSUMER_UNRESOLVED_COUNT;
6+
37
import com.linkedin.venice.stats.AbstractVeniceStats;
8+
import com.linkedin.venice.stats.OpenTelemetryMetricsSetup;
9+
import com.linkedin.venice.stats.dimensions.VeniceMetricsDimensions;
10+
import com.linkedin.venice.stats.metrics.MetricEntityStateBase;
11+
import com.linkedin.venice.stats.metrics.TehutiMetricNameEnum;
12+
import io.opentelemetry.api.common.Attributes;
413
import io.tehuti.metrics.MetricsRepository;
5-
import io.tehuti.metrics.Sensor;
614
import io.tehuti.metrics.stats.OccurrenceRate;
15+
import java.util.Collections;
16+
import java.util.Map;
717

818

919
public class StuckConsumerRepairStats extends AbstractVeniceStats {
10-
private Sensor stuckConsumerFound;
11-
private Sensor ingestionTaskRepair;
12-
private Sensor repairFailure;
20+
/** Tehuti metric names for StuckConsumerRepairStats sensors. */
21+
enum TehutiMetricName implements TehutiMetricNameEnum {
22+
STUCK_CONSUMER_FOUND, INGESTION_TASK_REPAIR, REPAIR_FAILURE
23+
}
24+
25+
private final MetricEntityStateBase stuckConsumerFoundOtel;
26+
private final MetricEntityStateBase ingestionTaskRepairOtel;
27+
private final MetricEntityStateBase repairFailureOtel;
1328

14-
public StuckConsumerRepairStats(MetricsRepository metricsRepository) {
29+
public StuckConsumerRepairStats(MetricsRepository metricsRepository, String clusterName) {
1530
super(metricsRepository, "StuckConsumerRepair");
1631

17-
this.stuckConsumerFound = registerSensor("stuck_consumer_found", new OccurrenceRate());
18-
this.ingestionTaskRepair = registerSensor("ingestion_task_repair", new OccurrenceRate());
19-
this.repairFailure = registerSensor("repair_failure", new OccurrenceRate());
32+
OpenTelemetryMetricsSetup.OpenTelemetryMetricsSetupInfo otelData =
33+
OpenTelemetryMetricsSetup.builder(metricsRepository).setClusterName(clusterName).build();
34+
Map<VeniceMetricsDimensions, String> baseDimensionsMap = otelData.getBaseDimensionsMap();
35+
Attributes baseAttributes = otelData.getBaseAttributes();
36+
37+
stuckConsumerFoundOtel = MetricEntityStateBase.create(
38+
STUCK_CONSUMER_DETECTED_COUNT.getMetricEntity(),
39+
otelData.getOtelRepository(),
40+
this::registerSensorIfAbsent,
41+
TehutiMetricName.STUCK_CONSUMER_FOUND,
42+
Collections.singletonList(new OccurrenceRate()),
43+
baseDimensionsMap,
44+
baseAttributes);
45+
46+
ingestionTaskRepairOtel = MetricEntityStateBase.create(
47+
STUCK_CONSUMER_TASK_REPAIRED_COUNT.getMetricEntity(),
48+
otelData.getOtelRepository(),
49+
this::registerSensorIfAbsent,
50+
TehutiMetricName.INGESTION_TASK_REPAIR,
51+
Collections.singletonList(new OccurrenceRate()),
52+
baseDimensionsMap,
53+
baseAttributes);
54+
55+
repairFailureOtel = MetricEntityStateBase.create(
56+
STUCK_CONSUMER_UNRESOLVED_COUNT.getMetricEntity(),
57+
otelData.getOtelRepository(),
58+
this::registerSensorIfAbsent,
59+
TehutiMetricName.REPAIR_FAILURE,
60+
Collections.singletonList(new OccurrenceRate()),
61+
baseDimensionsMap,
62+
baseAttributes);
2063
}
2164

2265
public void recordStuckConsumerFound() {
23-
stuckConsumerFound.record();
66+
stuckConsumerFoundOtel.record(1);
2467
}
2568

2669
public void recordIngestionTaskRepair() {
27-
ingestionTaskRepair.record();
70+
ingestionTaskRepairOtel.record(1);
2871
}
2972

3073
public void recordRepairFailure() {
31-
repairFailure.record();
74+
repairFailureOtel.record(1);
3275
}
3376
}

clients/da-vinci-client/src/test/java/com/linkedin/davinci/stats/ServerMetricEntityTest.java

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -22,7 +22,7 @@
2222
public class ServerMetricEntityTest {
2323
@Test
2424
public void testServerMetricEntitiesCount() {
25-
assertEquals(SERVER_METRIC_ENTITIES.size(), 179, "Expected 179 unique metric entities");
25+
assertEquals(SERVER_METRIC_ENTITIES.size(), 182, "Expected 182 unique metric entities");
2626
}
2727

2828
/**
Lines changed: 52 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,52 @@
1+
package com.linkedin.davinci.stats;
2+
3+
import static com.linkedin.davinci.stats.StuckConsumerRepairOtelMetricEntity.STUCK_CONSUMER_DETECTED_COUNT;
4+
import static com.linkedin.davinci.stats.StuckConsumerRepairOtelMetricEntity.STUCK_CONSUMER_TASK_REPAIRED_COUNT;
5+
import static com.linkedin.davinci.stats.StuckConsumerRepairOtelMetricEntity.STUCK_CONSUMER_UNRESOLVED_COUNT;
6+
import static com.linkedin.venice.stats.dimensions.VeniceMetricsDimensions.VENICE_CLUSTER_NAME;
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 StuckConsumerRepairOtelMetricEntityTest {
19+
@Test
20+
public void testMetricEntities() {
21+
new ModuleMetricEntityTestFixture<>(StuckConsumerRepairOtelMetricEntity.class, expectedDefinitions()).assertAll();
22+
}
23+
24+
private static Map<StuckConsumerRepairOtelMetricEntity, MetricEntityExpectation> expectedDefinitions() {
25+
Map<StuckConsumerRepairOtelMetricEntity, MetricEntityExpectation> map = new HashMap<>();
26+
map.put(
27+
STUCK_CONSUMER_DETECTED_COUNT,
28+
new MetricEntityExpectation(
29+
"ingestion.pubsub.consumer.stuck.detected_count",
30+
MetricType.COUNTER,
31+
MetricUnit.NUMBER,
32+
"Count of scans that detected a stuck PubSub consumer",
33+
setOf(VENICE_CLUSTER_NAME)));
34+
map.put(
35+
STUCK_CONSUMER_TASK_REPAIRED_COUNT,
36+
new MetricEntityExpectation(
37+
"ingestion.pubsub.consumer.stuck.task_repaired_count",
38+
MetricType.COUNTER,
39+
MetricUnit.NUMBER,
40+
"Count of ingestion tasks killed to unblock a stuck PubSub consumer",
41+
setOf(VENICE_CLUSTER_NAME)));
42+
map.put(
43+
STUCK_CONSUMER_UNRESOLVED_COUNT,
44+
new MetricEntityExpectation(
45+
"ingestion.pubsub.consumer.stuck.unresolved_count",
46+
MetricType.COUNTER,
47+
MetricUnit.NUMBER,
48+
"Count of scans where a stuck PubSub consumer was found but no fixable task identified",
49+
setOf(VENICE_CLUSTER_NAME)));
50+
return map;
51+
}
52+
}
Lines changed: 149 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,149 @@
1+
package com.linkedin.davinci.stats;
2+
3+
import static com.linkedin.davinci.stats.ServerMetricEntity.SERVER_METRIC_ENTITIES;
4+
import static com.linkedin.venice.stats.dimensions.VeniceMetricsDimensions.VENICE_CLUSTER_NAME;
5+
import static org.testng.Assert.assertNotNull;
6+
import static org.testng.Assert.assertTrue;
7+
8+
import com.linkedin.venice.stats.VeniceMetricsConfig;
9+
import com.linkedin.venice.stats.VeniceMetricsRepository;
10+
import com.linkedin.venice.utils.OpenTelemetryDataTestUtils;
11+
import io.opentelemetry.api.common.Attributes;
12+
import io.opentelemetry.sdk.testing.exporter.InMemoryMetricReader;
13+
import io.tehuti.metrics.MetricConfig;
14+
import io.tehuti.metrics.MetricsRepository;
15+
import io.tehuti.metrics.stats.AsyncGauge;
16+
import org.testng.annotations.AfterMethod;
17+
import org.testng.annotations.BeforeMethod;
18+
import org.testng.annotations.Test;
19+
20+
21+
public class StuckConsumerRepairStatsTest {
22+
private static final String TEST_METRIC_PREFIX = "server";
23+
private static final String TEST_CLUSTER_NAME = "test-cluster";
24+
25+
private InMemoryMetricReader inMemoryMetricReader;
26+
private VeniceMetricsRepository metricsRepository;
27+
private StuckConsumerRepairStats stats;
28+
// Dedicated executor avoids shutting down Tehuti's static DEFAULT_ASYNC_GAUGE_EXECUTOR singleton when this
29+
// repository is closed in tearDown(). Closing the static executor would break AsyncGauge measurements
30+
// for any subsequent test running in the same JVM.
31+
private AsyncGauge.AsyncGaugeExecutor asyncGaugeExecutor;
32+
33+
@BeforeMethod
34+
public void setUp() {
35+
inMemoryMetricReader = InMemoryMetricReader.create();
36+
asyncGaugeExecutor = new AsyncGauge.AsyncGaugeExecutor.Builder().build();
37+
metricsRepository = new VeniceMetricsRepository(
38+
new VeniceMetricsConfig.Builder().setMetricPrefix(TEST_METRIC_PREFIX)
39+
.setMetricEntities(SERVER_METRIC_ENTITIES)
40+
.setEmitOtelMetrics(true)
41+
.setOtelAdditionalMetricsReader(inMemoryMetricReader)
42+
.setTehutiMetricConfig(new MetricConfig(asyncGaugeExecutor))
43+
.build());
44+
stats = new StuckConsumerRepairStats(metricsRepository, TEST_CLUSTER_NAME);
45+
}
46+
47+
@AfterMethod
48+
public void tearDown() {
49+
// Closes only the dedicated executor; the JVM-wide static executor stays alive for other tests.
50+
if (metricsRepository != null) {
51+
metricsRepository.close();
52+
}
53+
}
54+
55+
// --- Positive tests: OTel counter accumulation ---
56+
57+
@Test
58+
public void testRecordStuckConsumerFound() {
59+
stats.recordStuckConsumerFound();
60+
stats.recordStuckConsumerFound();
61+
62+
OpenTelemetryDataTestUtils.validateLongPointDataFromCounter(
63+
inMemoryMetricReader,
64+
2,
65+
buildAttributes(),
66+
StuckConsumerRepairOtelMetricEntity.STUCK_CONSUMER_DETECTED_COUNT.getMetricEntity().getMetricName(),
67+
TEST_METRIC_PREFIX);
68+
}
69+
70+
@Test
71+
public void testRecordIngestionTaskRepair() {
72+
stats.recordIngestionTaskRepair();
73+
stats.recordIngestionTaskRepair();
74+
stats.recordIngestionTaskRepair();
75+
76+
OpenTelemetryDataTestUtils.validateLongPointDataFromCounter(
77+
inMemoryMetricReader,
78+
3,
79+
buildAttributes(),
80+
StuckConsumerRepairOtelMetricEntity.STUCK_CONSUMER_TASK_REPAIRED_COUNT.getMetricEntity().getMetricName(),
81+
TEST_METRIC_PREFIX);
82+
}
83+
84+
@Test
85+
public void testRecordRepairFailure() {
86+
stats.recordRepairFailure();
87+
88+
OpenTelemetryDataTestUtils.validateLongPointDataFromCounter(
89+
inMemoryMetricReader,
90+
1,
91+
buildAttributes(),
92+
StuckConsumerRepairOtelMetricEntity.STUCK_CONSUMER_UNRESOLVED_COUNT.getMetricEntity().getMetricName(),
93+
TEST_METRIC_PREFIX);
94+
}
95+
96+
// --- Positive tests: Tehuti sensor existence and recording ---
97+
98+
@Test
99+
public void testTehutiSensorsRegisteredAndRecorded() {
100+
// OccurrenceRate sensors use the ".OccurrenceRate" suffix
101+
String foundSensor = ".StuckConsumerRepair--stuck_consumer_found.OccurrenceRate";
102+
String repairSensor = ".StuckConsumerRepair--ingestion_task_repair.OccurrenceRate";
103+
String failureSensor = ".StuckConsumerRepair--repair_failure.OccurrenceRate";
104+
105+
assertNotNull(metricsRepository.getMetric(foundSensor), "Tehuti sensor should exist for stuck_consumer_found");
106+
assertNotNull(metricsRepository.getMetric(repairSensor), "Tehuti sensor should exist for ingestion_task_repair");
107+
assertNotNull(metricsRepository.getMetric(failureSensor), "Tehuti sensor should exist for repair_failure");
108+
109+
// Record and verify values are non-zero (OccurrenceRate is time-dependent, so just assert > 0)
110+
stats.recordStuckConsumerFound();
111+
stats.recordIngestionTaskRepair();
112+
stats.recordRepairFailure();
113+
assertTrue(metricsRepository.getMetric(foundSensor).value() > 0, "stuck_consumer_found rate should be > 0");
114+
assertTrue(metricsRepository.getMetric(repairSensor).value() > 0, "ingestion_task_repair rate should be > 0");
115+
assertTrue(metricsRepository.getMetric(failureSensor).value() > 0, "repair_failure rate should be > 0");
116+
}
117+
118+
// --- NPE prevention tests ---
119+
120+
@Test
121+
public void testNoNpeWhenOtelDisabled() {
122+
AsyncGauge.AsyncGaugeExecutor localExecutor = new AsyncGauge.AsyncGaugeExecutor.Builder().build();
123+
try (VeniceMetricsRepository disabledRepo = new VeniceMetricsRepository(
124+
new VeniceMetricsConfig.Builder().setMetricPrefix(TEST_METRIC_PREFIX)
125+
.setEmitOtelMetrics(false)
126+
.setTehutiMetricConfig(new MetricConfig(localExecutor))
127+
.build())) {
128+
exerciseAllRecordingPaths(disabledRepo);
129+
}
130+
}
131+
132+
@Test
133+
public void testNoNpeWhenPlainMetricsRepository() {
134+
exerciseAllRecordingPaths(new MetricsRepository());
135+
}
136+
137+
private void exerciseAllRecordingPaths(MetricsRepository repo) {
138+
StuckConsumerRepairStats safeStats = new StuckConsumerRepairStats(repo, TEST_CLUSTER_NAME);
139+
safeStats.recordStuckConsumerFound();
140+
safeStats.recordIngestionTaskRepair();
141+
safeStats.recordRepairFailure();
142+
}
143+
144+
// --- Helper ---
145+
146+
private Attributes buildAttributes() {
147+
return Attributes.builder().put(VENICE_CLUSTER_NAME.getDimensionNameInDefaultFormat(), TEST_CLUSTER_NAME).build();
148+
}
149+
}
Lines changed: 23 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,23 @@
1+
package com.linkedin.davinci.stats;
2+
3+
import com.linkedin.venice.stats.metrics.TehutiMetricNameEnumTestFixture;
4+
import java.util.HashMap;
5+
import java.util.Map;
6+
import org.testng.annotations.Test;
7+
8+
9+
public class StuckConsumerRepairTehutiMetricNameTest {
10+
@Test
11+
public void testTehutiMetricNames() {
12+
new TehutiMetricNameEnumTestFixture<>(StuckConsumerRepairStats.TehutiMetricName.class, expectedMetricNames())
13+
.assertAll();
14+
}
15+
16+
private static Map<StuckConsumerRepairStats.TehutiMetricName, String> expectedMetricNames() {
17+
Map<StuckConsumerRepairStats.TehutiMetricName, String> map = new HashMap<>();
18+
map.put(StuckConsumerRepairStats.TehutiMetricName.STUCK_CONSUMER_FOUND, "stuck_consumer_found");
19+
map.put(StuckConsumerRepairStats.TehutiMetricName.INGESTION_TASK_REPAIR, "ingestion_task_repair");
20+
map.put(StuckConsumerRepairStats.TehutiMetricName.REPAIR_FAILURE, "repair_failure");
21+
return map;
22+
}
23+
}

0 commit comments

Comments
 (0)