11package 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+
6+ import com .google .common .annotations .VisibleForTesting ;
37import com .linkedin .davinci .config .VeniceServerConfig ;
48import com .linkedin .venice .meta .ReadOnlyStoreRepository ;
9+ import com .linkedin .venice .stats .OpenTelemetryMetricsSetup ;
10+ import com .linkedin .venice .stats .VeniceOpenTelemetryMetricsRepository ;
11+ import com .linkedin .venice .stats .dimensions .VeniceMetricsDimensions ;
12+ import com .linkedin .venice .stats .dimensions .VeniceRecordTransformerOperation ;
13+ import com .linkedin .venice .stats .metrics .MetricEntityStateOneEnum ;
14+ import com .linkedin .venice .utils .concurrent .VeniceConcurrentHashMap ;
515import io .tehuti .metrics .MetricsRepository ;
16+ import java .util .HashMap ;
17+ import java .util .Map ;
618
719
820/**
9- * The store level stats for {@link com.linkedin.davinci.client.DaVinciRecordTransformer}
21+ * The store level stats for {@link com.linkedin.davinci.client.DaVinciRecordTransformer}.
22+ * OTel metrics are recorded directly here (separate API) because Tehuti uses the Reporter
23+ * layer ({@link DaVinciRecordTransformerStatsReporter}) with AsyncGauge polling, while OTel
24+ * records at the point of the call.
1025 */
1126public class AggVersionedDaVinciRecordTransformerStats
1227 extends AbstractVeniceAggVersionedStats <DaVinciRecordTransformerStats , DaVinciRecordTransformerStatsReporter > {
28+ private final VeniceOpenTelemetryMetricsRepository otelRepository ;
29+ private final Map <VeniceMetricsDimensions , String > baseDimensionsMap ;
30+ private final boolean emitOtelMetrics ;
31+
32+ /**
33+ * Per-store OTel metric state for latency. Bounded by the number of stores on this host.
34+ * Entries created lazily inside {@link #recordOtelLatency}, removed in
35+ * {@link #handleStoreDeleted(String)}.
36+ */
37+ private final Map <String , MetricEntityStateOneEnum <VeniceRecordTransformerOperation >> latencyPerStore =
38+ new VeniceConcurrentHashMap <>();
39+
40+ /**
41+ * Per-store OTel metric state for error count. Same bounding and lifecycle as latencyPerStore.
42+ */
43+ private final Map <String , MetricEntityStateOneEnum <VeniceRecordTransformerOperation >> errorCountPerStore =
44+ new VeniceConcurrentHashMap <>();
45+
1346 public AggVersionedDaVinciRecordTransformerStats (
1447 MetricsRepository metricsRepository ,
1548 ReadOnlyStoreRepository metadataRepository ,
@@ -20,21 +53,90 @@ public AggVersionedDaVinciRecordTransformerStats(
2053 DaVinciRecordTransformerStats ::new ,
2154 DaVinciRecordTransformerStatsReporter ::new ,
2255 serverConfig .isUnregisterMetricForDeletedStoreEnabled ());
56+
57+ OpenTelemetryMetricsSetup .OpenTelemetryMetricsSetupInfo otelData =
58+ OpenTelemetryMetricsSetup .builder (metricsRepository ).setClusterName (serverConfig .getClusterName ()).build ();
59+ this .otelRepository = otelData .getOtelRepository ();
60+ this .baseDimensionsMap = otelData .getBaseDimensionsMap ();
61+ this .emitOtelMetrics = otelData .emitOpenTelemetryMetrics ();
62+ }
63+
64+ @ Override
65+ public void handleStoreDeleted (String storeName ) {
66+ try {
67+ super .handleStoreDeleted (storeName );
68+ } finally {
69+ latencyPerStore .remove (storeName );
70+ errorCountPerStore .remove (storeName );
71+ }
2372 }
2473
2574 public void recordPutLatency (String storeName , int version , double value , long timestamp ) {
2675 recordVersionedAndTotalStat (storeName , version , stat -> stat .recordPutLatency (value , timestamp ));
76+ recordOtelLatency (storeName , value , VeniceRecordTransformerOperation .PUT );
2777 }
2878
2979 public void recordDeleteLatency (String storeName , int version , double value , long timestamp ) {
3080 recordVersionedAndTotalStat (storeName , version , stat -> stat .recordDeleteLatency (value , timestamp ));
81+ recordOtelLatency (storeName , value , VeniceRecordTransformerOperation .DELETE );
3182 }
3283
3384 public void recordPutError (String storeName , int version , long timestamp ) {
3485 recordVersionedAndTotalStat (storeName , version , stat -> stat .recordPutError (timestamp ));
86+ recordOtelErrorCount (storeName , VeniceRecordTransformerOperation .PUT );
3587 }
3688
3789 public void recordDeleteError (String storeName , int version , long timestamp ) {
3890 recordVersionedAndTotalStat (storeName , version , stat -> stat .recordDeleteError (timestamp ));
91+ recordOtelErrorCount (storeName , VeniceRecordTransformerOperation .DELETE );
92+ }
93+
94+ private void recordOtelLatency (String storeName , double value , VeniceRecordTransformerOperation operation ) {
95+ if (!emitOtelMetrics ) {
96+ return ;
97+ }
98+ latencyPerStore .computeIfAbsent (storeName , k -> createPerStoreMetric (k , RECORD_TRANSFORMER_LATENCY ))
99+ .record (value , operation );
100+ }
101+
102+ private void recordOtelErrorCount (String storeName , VeniceRecordTransformerOperation operation ) {
103+ if (!emitOtelMetrics ) {
104+ return ;
105+ }
106+ errorCountPerStore .computeIfAbsent (storeName , k -> createPerStoreMetric (k , RECORD_TRANSFORMER_ERROR_COUNT ))
107+ .record (1 , operation );
108+ }
109+
110+ @ VisibleForTesting
111+ boolean hasLatencyMetricFor (String storeName ) {
112+ return latencyPerStore .containsKey (storeName );
113+ }
114+
115+ @ VisibleForTesting
116+ boolean hasErrorCountMetricFor (String storeName ) {
117+ return errorCountPerStore .containsKey (storeName );
118+ }
119+
120+ @ VisibleForTesting
121+ int latencyStoreCount () {
122+ return latencyPerStore .size ();
123+ }
124+
125+ @ VisibleForTesting
126+ int errorCountStoreCount () {
127+ return errorCountPerStore .size ();
128+ }
129+
130+ private MetricEntityStateOneEnum <VeniceRecordTransformerOperation > createPerStoreMetric (
131+ String storeName ,
132+ DaVinciRecordTransformerOtelMetricEntity metricEntity ) {
133+ Map <VeniceMetricsDimensions , String > storeDimensionsMap = new HashMap <>(baseDimensionsMap );
134+ storeDimensionsMap
135+ .put (VeniceMetricsDimensions .VENICE_STORE_NAME , OpenTelemetryMetricsSetup .sanitizeStoreName (storeName ));
136+ return MetricEntityStateOneEnum .create (
137+ metricEntity .getMetricEntity (),
138+ otelRepository ,
139+ storeDimensionsMap ,
140+ VeniceRecordTransformerOperation .class );
39141 }
40142}
0 commit comments