Skip to content

Commit ebe871a

Browse files
committed
[da-vinci] Add OTel metrics to AggVersionedDaVinciRecordTransformerStats
Add 2 OTel metrics for DaVinci record transformer: - record_transformer.latency (HISTOGRAM with PUT/DELETE operation dimension) - record_transformer.error_count (COUNTER with PUT/DELETE operation dimension) Separate Tehuti+OTel API: Tehuti uses the Reporter layer (AsyncGauge polling), OTel records at the call point in AggVersionedDaVinciRecordTransformerStats. Per-store VeniceConcurrentHashMap maps for latency and error count metrics, cleaned up in handleStoreDeleted.
1 parent 413d87d commit ebe871a

10 files changed

Lines changed: 404 additions & 4 deletions

File tree

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

Lines changed: 72 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,15 +1,46 @@
11
package com.linkedin.davinci.stats;
22

3+
import static com.linkedin.davinci.stats.DaVinciRecordTransformerOtelMetricEntity.RECORD_TRANSFORMER_ERROR_COUNT;
4+
import static com.linkedin.davinci.stats.DaVinciRecordTransformerOtelMetricEntity.RECORD_TRANSFORMER_LATENCY;
5+
36
import com.linkedin.davinci.config.VeniceServerConfig;
47
import com.linkedin.venice.meta.ReadOnlyStoreRepository;
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.VeniceRecordTransformerOperation;
12+
import com.linkedin.venice.stats.metrics.MetricEntityStateOneEnum;
13+
import com.linkedin.venice.utils.concurrent.VeniceConcurrentHashMap;
514
import io.tehuti.metrics.MetricsRepository;
15+
import java.util.HashMap;
16+
import java.util.Map;
617

718

819
/**
9-
* The store level stats for {@link com.linkedin.davinci.client.DaVinciRecordTransformer}
20+
* The store level stats for {@link com.linkedin.davinci.client.DaVinciRecordTransformer}.
21+
* OTel metrics are recorded directly here (separate API) because Tehuti uses the Reporter
22+
* layer ({@link DaVinciRecordTransformerStatsReporter}) with AsyncGauge polling, while OTel
23+
* records at the point of the call.
1024
*/
1125
public class AggVersionedDaVinciRecordTransformerStats
1226
extends AbstractVeniceAggVersionedStats<DaVinciRecordTransformerStats, DaVinciRecordTransformerStatsReporter> {
27+
private final VeniceOpenTelemetryMetricsRepository otelRepository;
28+
private final Map<VeniceMetricsDimensions, String> baseDimensionsMap;
29+
30+
/**
31+
* Per-store OTel metric state for latency. Bounded by the number of stores on this host.
32+
* Entries created lazily via {@link #getOrCreateLatencyMetric}, removed in
33+
* {@link #handleStoreDeleted(String)}.
34+
*/
35+
private final Map<String, MetricEntityStateOneEnum<VeniceRecordTransformerOperation>> latencyPerStore =
36+
new VeniceConcurrentHashMap<>();
37+
38+
/**
39+
* Per-store OTel metric state for error count. Same bounding and lifecycle as latencyPerStore.
40+
*/
41+
private final Map<String, MetricEntityStateOneEnum<VeniceRecordTransformerOperation>> errorCountPerStore =
42+
new VeniceConcurrentHashMap<>();
43+
1344
public AggVersionedDaVinciRecordTransformerStats(
1445
MetricsRepository metricsRepository,
1546
ReadOnlyStoreRepository metadataRepository,
@@ -20,21 +51,61 @@ public AggVersionedDaVinciRecordTransformerStats(
2051
DaVinciRecordTransformerStats::new,
2152
DaVinciRecordTransformerStatsReporter::new,
2253
serverConfig.isUnregisterMetricForDeletedStoreEnabled());
54+
55+
OpenTelemetryMetricsSetup.OpenTelemetryMetricsSetupInfo otelData =
56+
OpenTelemetryMetricsSetup.builder(metricsRepository).setClusterName(serverConfig.getClusterName()).build();
57+
this.otelRepository = otelData.getOtelRepository();
58+
this.baseDimensionsMap = otelData.getBaseDimensionsMap();
59+
}
60+
61+
@Override
62+
public void handleStoreDeleted(String storeName) {
63+
try {
64+
super.handleStoreDeleted(storeName);
65+
} finally {
66+
latencyPerStore.remove(storeName);
67+
errorCountPerStore.remove(storeName);
68+
}
2369
}
2470

2571
public void recordPutLatency(String storeName, int version, double value, long timestamp) {
2672
recordVersionedAndTotalStat(storeName, version, stat -> stat.recordPutLatency(value, timestamp));
73+
getOrCreateLatencyMetric(storeName).record(value, VeniceRecordTransformerOperation.PUT);
2774
}
2875

2976
public void recordDeleteLatency(String storeName, int version, double value, long timestamp) {
3077
recordVersionedAndTotalStat(storeName, version, stat -> stat.recordDeleteLatency(value, timestamp));
78+
getOrCreateLatencyMetric(storeName).record(value, VeniceRecordTransformerOperation.DELETE);
3179
}
3280

3381
public void recordPutError(String storeName, int version, long timestamp) {
3482
recordVersionedAndTotalStat(storeName, version, stat -> stat.recordPutError(timestamp));
83+
getOrCreateErrorCountMetric(storeName).record(1, VeniceRecordTransformerOperation.PUT);
3584
}
3685

3786
public void recordDeleteError(String storeName, int version, long timestamp) {
3887
recordVersionedAndTotalStat(storeName, version, stat -> stat.recordDeleteError(timestamp));
88+
getOrCreateErrorCountMetric(storeName).record(1, VeniceRecordTransformerOperation.DELETE);
89+
}
90+
91+
private MetricEntityStateOneEnum<VeniceRecordTransformerOperation> getOrCreateLatencyMetric(String storeName) {
92+
return latencyPerStore.computeIfAbsent(storeName, k -> createPerStoreMetric(k, RECORD_TRANSFORMER_LATENCY));
93+
}
94+
95+
private MetricEntityStateOneEnum<VeniceRecordTransformerOperation> getOrCreateErrorCountMetric(String storeName) {
96+
return errorCountPerStore.computeIfAbsent(storeName, k -> createPerStoreMetric(k, RECORD_TRANSFORMER_ERROR_COUNT));
97+
}
98+
99+
private MetricEntityStateOneEnum<VeniceRecordTransformerOperation> createPerStoreMetric(
100+
String storeName,
101+
DaVinciRecordTransformerOtelMetricEntity metricEntity) {
102+
Map<VeniceMetricsDimensions, String> storeDimensionsMap = new HashMap<>(baseDimensionsMap);
103+
storeDimensionsMap
104+
.put(VeniceMetricsDimensions.VENICE_STORE_NAME, OpenTelemetryMetricsSetup.sanitizeStoreName(storeName));
105+
return MetricEntityStateOneEnum.create(
106+
metricEntity.getMetricEntity(),
107+
otelRepository,
108+
storeDimensionsMap,
109+
VeniceRecordTransformerOperation.class);
39110
}
40111
}
Lines changed: 48 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,48 @@
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_RECORD_TRANSFORMER_OPERATION;
5+
import static com.linkedin.venice.stats.dimensions.VeniceMetricsDimensions.VENICE_STORE_NAME;
6+
import static com.linkedin.venice.utils.Utils.setOf;
7+
8+
import com.linkedin.venice.stats.dimensions.VeniceMetricsDimensions;
9+
import com.linkedin.venice.stats.metrics.MetricEntity;
10+
import com.linkedin.venice.stats.metrics.MetricType;
11+
import com.linkedin.venice.stats.metrics.MetricUnit;
12+
import com.linkedin.venice.stats.metrics.ModuleMetricEntityInterface;
13+
import java.util.Set;
14+
15+
16+
/**
17+
* OTel metric entity definitions for {@link DaVinciRecordTransformerStats}.
18+
* Tracks DaVinci record transformer latency and error counts by operation (put/delete).
19+
*/
20+
public enum DaVinciRecordTransformerOtelMetricEntity implements ModuleMetricEntityInterface {
21+
RECORD_TRANSFORMER_LATENCY(
22+
"record_transformer.latency", MetricType.HISTOGRAM, MetricUnit.MILLISECOND,
23+
"DaVinci record transformer operation latency",
24+
setOf(VENICE_CLUSTER_NAME, VENICE_STORE_NAME, VENICE_RECORD_TRANSFORMER_OPERATION)
25+
),
26+
27+
RECORD_TRANSFORMER_ERROR_COUNT(
28+
"record_transformer.error_count", MetricType.COUNTER, MetricUnit.NUMBER,
29+
"DaVinci record transformer operation error count",
30+
setOf(VENICE_CLUSTER_NAME, VENICE_STORE_NAME, VENICE_RECORD_TRANSFORMER_OPERATION)
31+
);
32+
33+
private final MetricEntity metricEntity;
34+
35+
DaVinciRecordTransformerOtelMetricEntity(
36+
String metricName,
37+
MetricType metricType,
38+
MetricUnit unit,
39+
String description,
40+
Set<VeniceMetricsDimensions> dimensions) {
41+
this.metricEntity = new MetricEntity(metricName, metricType, unit, description, dimensions);
42+
}
43+
44+
@Override
45+
public MetricEntity getMetricEntity() {
46+
return metricEntity;
47+
}
48+
}

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+
DaVinciRecordTransformerOtelMetricEntity.class);
4748
}
4849

4950
public static final Collection<MetricEntity> SERVER_METRIC_ENTITIES =
Lines changed: 187 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,187 @@
1+
package com.linkedin.davinci.stats;
2+
3+
import static com.linkedin.davinci.stats.DaVinciRecordTransformerOtelMetricEntity.RECORD_TRANSFORMER_ERROR_COUNT;
4+
import static com.linkedin.davinci.stats.DaVinciRecordTransformerOtelMetricEntity.RECORD_TRANSFORMER_LATENCY;
5+
import static com.linkedin.davinci.stats.ServerMetricEntity.SERVER_METRIC_ENTITIES;
6+
import static com.linkedin.venice.stats.dimensions.VeniceMetricsDimensions.VENICE_CLUSTER_NAME;
7+
import static com.linkedin.venice.stats.dimensions.VeniceMetricsDimensions.VENICE_RECORD_TRANSFORMER_OPERATION;
8+
import static com.linkedin.venice.stats.dimensions.VeniceMetricsDimensions.VENICE_STORE_NAME;
9+
import static org.mockito.ArgumentMatchers.anyString;
10+
import static org.mockito.Mockito.doReturn;
11+
import static org.mockito.Mockito.mock;
12+
13+
import com.linkedin.davinci.config.VeniceServerConfig;
14+
import com.linkedin.venice.meta.ReadOnlyStoreRepository;
15+
import com.linkedin.venice.meta.Store;
16+
import com.linkedin.venice.stats.VeniceMetricsConfig;
17+
import com.linkedin.venice.stats.VeniceMetricsRepository;
18+
import com.linkedin.venice.stats.dimensions.VeniceRecordTransformerOperation;
19+
import com.linkedin.venice.utils.OpenTelemetryDataTestUtils;
20+
import io.opentelemetry.api.common.Attributes;
21+
import io.opentelemetry.sdk.testing.exporter.InMemoryMetricReader;
22+
import io.tehuti.metrics.MetricsRepository;
23+
import java.util.Collections;
24+
import org.testng.annotations.AfterMethod;
25+
import org.testng.annotations.BeforeMethod;
26+
import org.testng.annotations.Test;
27+
28+
29+
public class AggVersionedDaVinciRecordTransformerStatsOtelTest {
30+
private static final String TEST_METRIC_PREFIX = "davinci_client";
31+
private static final String TEST_CLUSTER_NAME = "test-cluster";
32+
private static final String TEST_STORE_NAME = "test-store";
33+
34+
private InMemoryMetricReader inMemoryMetricReader;
35+
private VeniceMetricsRepository metricsRepository;
36+
private AggVersionedDaVinciRecordTransformerStats aggStats;
37+
38+
@BeforeMethod
39+
public void setUp() {
40+
inMemoryMetricReader = InMemoryMetricReader.create();
41+
metricsRepository = new VeniceMetricsRepository(
42+
new VeniceMetricsConfig.Builder().setMetricPrefix(TEST_METRIC_PREFIX)
43+
.setMetricEntities(SERVER_METRIC_ENTITIES)
44+
.setEmitOtelMetrics(true)
45+
.setOtelAdditionalMetricsReader(inMemoryMetricReader)
46+
.build());
47+
aggStats = createAggStats(metricsRepository);
48+
}
49+
50+
@AfterMethod
51+
public void tearDown() {
52+
if (metricsRepository != null) {
53+
metricsRepository.close();
54+
}
55+
}
56+
57+
@Test
58+
public void testRecordPutLatency() {
59+
long timestamp = System.currentTimeMillis();
60+
aggStats.recordPutLatency(TEST_STORE_NAME, 1, 50.0, timestamp);
61+
aggStats.recordPutLatency(TEST_STORE_NAME, 1, 100.0, timestamp);
62+
63+
OpenTelemetryDataTestUtils.validateExponentialHistogramPointData(
64+
inMemoryMetricReader,
65+
50.0,
66+
100.0,
67+
2,
68+
150.0,
69+
buildAttributes(VeniceRecordTransformerOperation.PUT),
70+
RECORD_TRANSFORMER_LATENCY.getMetricEntity().getMetricName(),
71+
TEST_METRIC_PREFIX);
72+
}
73+
74+
@Test
75+
public void testRecordDeleteLatency() {
76+
long timestamp = System.currentTimeMillis();
77+
aggStats.recordDeleteLatency(TEST_STORE_NAME, 1, 25.0, timestamp);
78+
79+
OpenTelemetryDataTestUtils.validateExponentialHistogramPointData(
80+
inMemoryMetricReader,
81+
25.0,
82+
25.0,
83+
1,
84+
25.0,
85+
buildAttributes(VeniceRecordTransformerOperation.DELETE),
86+
RECORD_TRANSFORMER_LATENCY.getMetricEntity().getMetricName(),
87+
TEST_METRIC_PREFIX);
88+
}
89+
90+
@Test
91+
public void testRecordPutError() {
92+
long timestamp = System.currentTimeMillis();
93+
aggStats.recordPutError(TEST_STORE_NAME, 1, timestamp);
94+
aggStats.recordPutError(TEST_STORE_NAME, 1, timestamp);
95+
96+
OpenTelemetryDataTestUtils.validateLongPointDataFromCounter(
97+
inMemoryMetricReader,
98+
2,
99+
buildAttributes(VeniceRecordTransformerOperation.PUT),
100+
RECORD_TRANSFORMER_ERROR_COUNT.getMetricEntity().getMetricName(),
101+
TEST_METRIC_PREFIX);
102+
}
103+
104+
@Test
105+
public void testRecordDeleteError() {
106+
long timestamp = System.currentTimeMillis();
107+
aggStats.recordDeleteError(TEST_STORE_NAME, 1, timestamp);
108+
109+
OpenTelemetryDataTestUtils.validateLongPointDataFromCounter(
110+
inMemoryMetricReader,
111+
1,
112+
buildAttributes(VeniceRecordTransformerOperation.DELETE),
113+
RECORD_TRANSFORMER_ERROR_COUNT.getMetricEntity().getMetricName(),
114+
TEST_METRIC_PREFIX);
115+
}
116+
117+
@Test
118+
public void testOperationDimensionIsolation() {
119+
long timestamp = System.currentTimeMillis();
120+
aggStats.recordPutError(TEST_STORE_NAME, 1, timestamp);
121+
aggStats.recordPutError(TEST_STORE_NAME, 1, timestamp);
122+
aggStats.recordDeleteError(TEST_STORE_NAME, 1, timestamp);
123+
124+
OpenTelemetryDataTestUtils.validateLongPointDataFromCounter(
125+
inMemoryMetricReader,
126+
2,
127+
buildAttributes(VeniceRecordTransformerOperation.PUT),
128+
RECORD_TRANSFORMER_ERROR_COUNT.getMetricEntity().getMetricName(),
129+
TEST_METRIC_PREFIX);
130+
131+
OpenTelemetryDataTestUtils.validateLongPointDataFromCounter(
132+
inMemoryMetricReader,
133+
1,
134+
buildAttributes(VeniceRecordTransformerOperation.DELETE),
135+
RECORD_TRANSFORMER_ERROR_COUNT.getMetricEntity().getMetricName(),
136+
TEST_METRIC_PREFIX);
137+
}
138+
139+
// --- NPE prevention tests ---
140+
141+
@Test
142+
public void testNoNpeWhenOtelDisabled() {
143+
try (VeniceMetricsRepository disabledRepo = new VeniceMetricsRepository(
144+
new VeniceMetricsConfig.Builder().setMetricPrefix(TEST_METRIC_PREFIX).setEmitOtelMetrics(false).build())) {
145+
exerciseAllRecordingPaths(disabledRepo);
146+
}
147+
}
148+
149+
@Test
150+
public void testNoNpeWhenPlainMetricsRepository() {
151+
exerciseAllRecordingPaths(new MetricsRepository());
152+
}
153+
154+
private void exerciseAllRecordingPaths(MetricsRepository repo) {
155+
AggVersionedDaVinciRecordTransformerStats safeStats = createAggStats(repo);
156+
long ts = System.currentTimeMillis();
157+
safeStats.recordPutLatency(TEST_STORE_NAME, 1, 10.0, ts);
158+
safeStats.recordDeleteLatency(TEST_STORE_NAME, 1, 10.0, ts);
159+
safeStats.recordPutError(TEST_STORE_NAME, 1, ts);
160+
safeStats.recordDeleteError(TEST_STORE_NAME, 1, ts);
161+
}
162+
163+
// --- Helpers ---
164+
165+
private static AggVersionedDaVinciRecordTransformerStats createAggStats(MetricsRepository repo) {
166+
ReadOnlyStoreRepository metadataRepository = mock(ReadOnlyStoreRepository.class);
167+
Store mockStore = mock(Store.class);
168+
doReturn(TEST_STORE_NAME).when(mockStore).getName();
169+
doReturn(Collections.emptyList()).when(mockStore).getVersions();
170+
doReturn(0).when(mockStore).getCurrentVersion();
171+
doReturn(mockStore).when(metadataRepository).getStoreOrThrow(anyString());
172+
173+
VeniceServerConfig serverConfig = mock(VeniceServerConfig.class);
174+
doReturn(false).when(serverConfig).isUnregisterMetricForDeletedStoreEnabled();
175+
doReturn(TEST_CLUSTER_NAME).when(serverConfig).getClusterName();
176+
177+
return new AggVersionedDaVinciRecordTransformerStats(repo, metadataRepository, serverConfig);
178+
}
179+
180+
private Attributes buildAttributes(VeniceRecordTransformerOperation operation) {
181+
return Attributes.builder()
182+
.put(VENICE_CLUSTER_NAME.getDimensionNameInDefaultFormat(), TEST_CLUSTER_NAME)
183+
.put(VENICE_STORE_NAME.getDimensionNameInDefaultFormat(), TEST_STORE_NAME)
184+
.put(VENICE_RECORD_TRANSFORMER_OPERATION.getDimensionNameInDefaultFormat(), operation.getDimensionValue())
185+
.build();
186+
}
187+
}

0 commit comments

Comments
 (0)