Skip to content

Commit 2730809

Browse files
committed
[da-vinci] Add OTel metrics to StuckConsumerRepairStats
Add 3 COUNTER metrics under ingestion.pubsub.consumer.stuck.* namespace: - detected_count: scans that detected a stuck consumer - task_repaired_count: ingestion tasks killed to unblock - unresolved_count: stuck consumers found with no fixable task Joint Tehuti+OTel API via MetricEntityStateBase. Singleton class with CLUSTER_NAME dimension only. clusterName added to constructor, passed from AggKafkaConsumerService via serverConfig.
1 parent 413d87d commit 2730809

8 files changed

Lines changed: 319 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
@@ -43,7 +43,8 @@ public static List<Class<? extends ModuleMetricEntityInterface>> getMetricEntity
4343
ServerReadQuotaOtelMetricEntity.class,
4444
ServerConnectionOtelMetricEntity.class,
4545
StoreBufferServiceOtelMetricEntity.class,
46-
StorageEngineOtelMetricEntity.class);
46+
StorageEngineOtelMetricEntity.class,
47+
StuckConsumerRepairOtelMetricEntity.class);
4748
}
4849

4950
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(), 153, "Expected 153 unique metric entities");
25+
assertEquals(SERVER_METRIC_ENTITIES.size(), 156, "Expected 156 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: 136 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,136 @@
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.MetricsRepository;
14+
import org.testng.annotations.AfterMethod;
15+
import org.testng.annotations.BeforeMethod;
16+
import org.testng.annotations.Test;
17+
18+
19+
public class StuckConsumerRepairStatsTest {
20+
private static final String TEST_METRIC_PREFIX = "server";
21+
private static final String TEST_CLUSTER_NAME = "test-cluster";
22+
23+
private InMemoryMetricReader inMemoryMetricReader;
24+
private VeniceMetricsRepository metricsRepository;
25+
private StuckConsumerRepairStats stats;
26+
27+
@BeforeMethod
28+
public void setUp() {
29+
inMemoryMetricReader = InMemoryMetricReader.create();
30+
metricsRepository = new VeniceMetricsRepository(
31+
new VeniceMetricsConfig.Builder().setMetricPrefix(TEST_METRIC_PREFIX)
32+
.setMetricEntities(SERVER_METRIC_ENTITIES)
33+
.setEmitOtelMetrics(true)
34+
.setOtelAdditionalMetricsReader(inMemoryMetricReader)
35+
.build());
36+
stats = new StuckConsumerRepairStats(metricsRepository, TEST_CLUSTER_NAME);
37+
}
38+
39+
@AfterMethod
40+
public void tearDown() {
41+
if (metricsRepository != null) {
42+
metricsRepository.close();
43+
}
44+
}
45+
46+
// --- Positive tests: OTel counter accumulation ---
47+
48+
@Test
49+
public void testRecordStuckConsumerFound() {
50+
stats.recordStuckConsumerFound();
51+
stats.recordStuckConsumerFound();
52+
53+
OpenTelemetryDataTestUtils.validateLongPointDataFromCounter(
54+
inMemoryMetricReader,
55+
2,
56+
buildAttributes(),
57+
StuckConsumerRepairOtelMetricEntity.STUCK_CONSUMER_DETECTED_COUNT.getMetricEntity().getMetricName(),
58+
TEST_METRIC_PREFIX);
59+
}
60+
61+
@Test
62+
public void testRecordIngestionTaskRepair() {
63+
stats.recordIngestionTaskRepair();
64+
stats.recordIngestionTaskRepair();
65+
stats.recordIngestionTaskRepair();
66+
67+
OpenTelemetryDataTestUtils.validateLongPointDataFromCounter(
68+
inMemoryMetricReader,
69+
3,
70+
buildAttributes(),
71+
StuckConsumerRepairOtelMetricEntity.STUCK_CONSUMER_TASK_REPAIRED_COUNT.getMetricEntity().getMetricName(),
72+
TEST_METRIC_PREFIX);
73+
}
74+
75+
@Test
76+
public void testRecordRepairFailure() {
77+
stats.recordRepairFailure();
78+
79+
OpenTelemetryDataTestUtils.validateLongPointDataFromCounter(
80+
inMemoryMetricReader,
81+
1,
82+
buildAttributes(),
83+
StuckConsumerRepairOtelMetricEntity.STUCK_CONSUMER_UNRESOLVED_COUNT.getMetricEntity().getMetricName(),
84+
TEST_METRIC_PREFIX);
85+
}
86+
87+
// --- Positive tests: Tehuti sensor existence and recording ---
88+
89+
@Test
90+
public void testTehutiSensorsRegisteredAndRecorded() {
91+
// OccurrenceRate sensors use the ".OccurrenceRate" suffix
92+
String foundSensor = ".StuckConsumerRepair--stuck_consumer_found.OccurrenceRate";
93+
String repairSensor = ".StuckConsumerRepair--ingestion_task_repair.OccurrenceRate";
94+
String failureSensor = ".StuckConsumerRepair--repair_failure.OccurrenceRate";
95+
96+
assertNotNull(metricsRepository.getMetric(foundSensor), "Tehuti sensor should exist for stuck_consumer_found");
97+
assertNotNull(metricsRepository.getMetric(repairSensor), "Tehuti sensor should exist for ingestion_task_repair");
98+
assertNotNull(metricsRepository.getMetric(failureSensor), "Tehuti sensor should exist for repair_failure");
99+
100+
// Record and verify values are non-zero (OccurrenceRate is time-dependent, so just assert > 0)
101+
stats.recordStuckConsumerFound();
102+
stats.recordIngestionTaskRepair();
103+
stats.recordRepairFailure();
104+
assertTrue(metricsRepository.getMetric(foundSensor).value() > 0, "stuck_consumer_found rate should be > 0");
105+
assertTrue(metricsRepository.getMetric(repairSensor).value() > 0, "ingestion_task_repair rate should be > 0");
106+
assertTrue(metricsRepository.getMetric(failureSensor).value() > 0, "repair_failure rate should be > 0");
107+
}
108+
109+
// --- NPE prevention tests ---
110+
111+
@Test
112+
public void testNoNpeWhenOtelDisabled() {
113+
try (VeniceMetricsRepository disabledRepo = new VeniceMetricsRepository(
114+
new VeniceMetricsConfig.Builder().setMetricPrefix(TEST_METRIC_PREFIX).setEmitOtelMetrics(false).build())) {
115+
exerciseAllRecordingPaths(disabledRepo);
116+
}
117+
}
118+
119+
@Test
120+
public void testNoNpeWhenPlainMetricsRepository() {
121+
exerciseAllRecordingPaths(new MetricsRepository());
122+
}
123+
124+
private void exerciseAllRecordingPaths(MetricsRepository repo) {
125+
StuckConsumerRepairStats safeStats = new StuckConsumerRepairStats(repo, TEST_CLUSTER_NAME);
126+
safeStats.recordStuckConsumerFound();
127+
safeStats.recordIngestionTaskRepair();
128+
safeStats.recordRepairFailure();
129+
}
130+
131+
// --- Helper ---
132+
133+
private Attributes buildAttributes() {
134+
return Attributes.builder().put(VENICE_CLUSTER_NAME.getDimensionNameInDefaultFormat(), TEST_CLUSTER_NAME).build();
135+
}
136+
}
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)