Skip to content

Commit 502671f

Browse files
committed
[da-vinci] Add OTel metrics to NativeMetadataRepositoryStats
Add per-store ASYNC_DOUBLE_GAUGE metric: metadata.staleness_duration with STORE_NAME dimension. Returns time in ms since the store metadata was last fetched from the meta system store. Returns NaN for removed stores (OTel SDK drops the data point — store disappears from dashboards). Tehuti unchanged: single high-watermark gauge across all stores. OTel: per-store staleness — backends compute max at query time.
1 parent 413d87d commit 502671f

6 files changed

Lines changed: 282 additions & 5 deletions

File tree

Lines changed: 39 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,39 @@
1+
package com.linkedin.davinci.stats;
2+
3+
import static com.linkedin.venice.stats.dimensions.VeniceMetricsDimensions.VENICE_STORE_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 NativeMetadataRepositoryStats}.
16+
*/
17+
public enum NativeMetadataRepositoryOtelMetricEntity implements ModuleMetricEntityInterface {
18+
METADATA_CACHE_STALENESS(
19+
"metadata.staleness_duration", MetricType.ASYNC_DOUBLE_GAUGE, MetricUnit.MILLISECOND,
20+
"Per-store metadata staleness in ms since the store metadata was last fetched from the meta system store",
21+
setOf(VENICE_STORE_NAME)
22+
);
23+
24+
private final MetricEntity metricEntity;
25+
26+
NativeMetadataRepositoryOtelMetricEntity(
27+
String metricName,
28+
MetricType metricType,
29+
MetricUnit unit,
30+
String description,
31+
Set<VeniceMetricsDimensions> dimensions) {
32+
this.metricEntity = new MetricEntity(metricName, metricType, unit, description, dimensions);
33+
}
34+
35+
@Override
36+
public MetricEntity getMetricEntity() {
37+
return metricEntity;
38+
}
39+
}

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

Lines changed: 55 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -1,26 +1,58 @@
11
package com.linkedin.davinci.stats;
22

3+
import static com.linkedin.davinci.stats.NativeMetadataRepositoryOtelMetricEntity.METADATA_CACHE_STALENESS;
4+
35
import com.linkedin.venice.stats.AbstractVeniceStats;
6+
import com.linkedin.venice.stats.OpenTelemetryMetricsSetup;
7+
import com.linkedin.venice.stats.VeniceOpenTelemetryMetricsRepository;
8+
import com.linkedin.venice.stats.dimensions.VeniceMetricsDimensions;
9+
import com.linkedin.venice.stats.metrics.AsyncMetricEntityStateBase;
410
import com.linkedin.venice.utils.concurrent.VeniceConcurrentHashMap;
11+
import io.opentelemetry.api.common.Attributes;
512
import io.tehuti.metrics.MetricsRepository;
6-
import io.tehuti.metrics.Sensor;
713
import io.tehuti.metrics.stats.AsyncGauge;
814
import java.time.Clock;
15+
import java.util.HashMap;
916
import java.util.Map;
17+
import java.util.function.DoubleSupplier;
1018

1119

20+
/**
21+
* Tracks metadata cache staleness for {@link com.linkedin.davinci.repository.NativeMetadataRepository}.
22+
*
23+
* <p>Tehuti emits a single high-watermark gauge (oldest store's staleness across all stores).
24+
* OTel emits per-store ASYNC_DOUBLE_GAUGE with STORE_NAME dimension — backends can compute the
25+
* high watermark at query time via max aggregation.
26+
*
27+
* <p>Per-store OTel callbacks are registered lazily on first {@link #updateCacheTimestamp} call
28+
* and read from the shared {@link #metadataCacheTimestampMapInMs}. When a store is removed,
29+
* the callback returns {@code NaN} (timestamp absent from the map, store no longer tracked). OTel callbacks cannot be deregistered
30+
* (SDK limitation), so the per-store entry stays registered until the process exits.
31+
*/
1232
public class NativeMetadataRepositoryStats extends AbstractVeniceStats {
13-
private final Sensor storeMetadataStalenessSensor;
1433
private final Map<String, Long> metadataCacheTimestampMapInMs = new VeniceConcurrentHashMap<>();
1534
private final Clock clock;
1635

36+
// OTel: per-store ASYNC_GAUGE for staleness. Bounded by number of subscribed stores.
37+
private final VeniceOpenTelemetryMetricsRepository otelRepository;
38+
private final Map<VeniceMetricsDimensions, String> baseDimensionsMap;
39+
private final Map<String, AsyncMetricEntityStateBase> otelPerStore = new VeniceConcurrentHashMap<>();
40+
1741
public NativeMetadataRepositoryStats(MetricsRepository metricsRepository, String name, Clock clock) {
1842
super(metricsRepository, name);
1943
this.clock = clock;
20-
this.storeMetadataStalenessSensor = registerSensor(
44+
45+
// Tehuti: single high-watermark gauge across all stores
46+
registerSensor(
2147
new AsyncGauge(
2248
(ignored1, ignored2) -> getMetadataStalenessHighWatermarkMs(),
2349
"store_metadata_staleness_high_watermark_ms"));
50+
51+
// OTel setup
52+
OpenTelemetryMetricsSetup.OpenTelemetryMetricsSetupInfo otelData =
53+
OpenTelemetryMetricsSetup.builder(metricsRepository).build();
54+
this.otelRepository = otelData.getOtelRepository();
55+
this.baseDimensionsMap = otelData.getBaseDimensionsMap();
2456
}
2557

2658
public final double getMetadataStalenessHighWatermarkMs() {
@@ -38,9 +70,29 @@ public final double getMetadataStalenessHighWatermarkMs() {
3870

3971
public void updateCacheTimestamp(String storeName, long cacheTimeStampInMs) {
4072
metadataCacheTimestampMapInMs.put(storeName, cacheTimeStampInMs);
73+
registerOtelGaugeIfAbsent(storeName);
4174
}
4275

4376
public void removeCacheTimestamp(String storeName) {
4477
metadataCacheTimestampMapInMs.remove(storeName);
78+
// OTel callback stays registered but returns NaN (timestamp absent from map, store no longer tracked)
79+
}
80+
81+
private void registerOtelGaugeIfAbsent(String storeName) {
82+
if (otelRepository == null) {
83+
return;
84+
}
85+
otelPerStore.computeIfAbsent(storeName, k -> {
86+
Map<VeniceMetricsDimensions, String> dims = new HashMap<>(baseDimensionsMap);
87+
dims.put(VeniceMetricsDimensions.VENICE_STORE_NAME, OpenTelemetryMetricsSetup.sanitizeStoreName(k));
88+
Attributes attrs = otelRepository.createAttributes(METADATA_CACHE_STALENESS.getMetricEntity(), dims);
89+
// DoubleSupplier callback: returns NaN when store is removed (no timestamp in map),
90+
// consistent with the Tehuti high-watermark gauge behavior.
91+
return AsyncMetricEntityStateBase
92+
.create(METADATA_CACHE_STALENESS.getMetricEntity(), otelRepository, dims, attrs, (DoubleSupplier) () -> {
93+
Long ts = metadataCacheTimestampMapInMs.get(k);
94+
return ts == null ? Double.NaN : (double) (clock.millis() - ts);
95+
});
96+
});
4597
}
4698
}

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

4950
public static final Collection<MetricEntity> SERVER_METRIC_ENTITIES =
Lines changed: 35 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,35 @@
1+
package com.linkedin.davinci.stats;
2+
3+
import static com.linkedin.davinci.stats.NativeMetadataRepositoryOtelMetricEntity.METADATA_CACHE_STALENESS;
4+
import static com.linkedin.venice.stats.dimensions.VeniceMetricsDimensions.VENICE_STORE_NAME;
5+
import static com.linkedin.venice.utils.Utils.setOf;
6+
7+
import com.linkedin.venice.stats.metrics.MetricType;
8+
import com.linkedin.venice.stats.metrics.MetricUnit;
9+
import com.linkedin.venice.stats.metrics.ModuleMetricEntityTestFixture;
10+
import com.linkedin.venice.stats.metrics.ModuleMetricEntityTestFixture.MetricEntityExpectation;
11+
import java.util.HashMap;
12+
import java.util.Map;
13+
import org.testng.annotations.Test;
14+
15+
16+
public class NativeMetadataRepositoryOtelMetricEntityTest {
17+
@Test
18+
public void testMetricEntities() {
19+
new ModuleMetricEntityTestFixture<>(NativeMetadataRepositoryOtelMetricEntity.class, expectedDefinitions())
20+
.assertAll();
21+
}
22+
23+
private static Map<NativeMetadataRepositoryOtelMetricEntity, MetricEntityExpectation> expectedDefinitions() {
24+
Map<NativeMetadataRepositoryOtelMetricEntity, MetricEntityExpectation> map = new HashMap<>();
25+
map.put(
26+
METADATA_CACHE_STALENESS,
27+
new MetricEntityExpectation(
28+
"metadata.staleness_duration",
29+
MetricType.ASYNC_DOUBLE_GAUGE,
30+
MetricUnit.MILLISECOND,
31+
"Per-store metadata staleness in ms since the store metadata was last fetched from the meta system store",
32+
setOf(VENICE_STORE_NAME)));
33+
return map;
34+
}
35+
}
Lines changed: 150 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,150 @@
1+
package com.linkedin.davinci.stats;
2+
3+
import static com.linkedin.davinci.stats.NativeMetadataRepositoryOtelMetricEntity.METADATA_CACHE_STALENESS;
4+
import static com.linkedin.davinci.stats.ServerMetricEntity.SERVER_METRIC_ENTITIES;
5+
import static com.linkedin.venice.stats.dimensions.VeniceMetricsDimensions.VENICE_STORE_NAME;
6+
import static org.mockito.Mockito.doReturn;
7+
import static org.mockito.Mockito.mock;
8+
import static org.testng.Assert.assertFalse;
9+
10+
import com.linkedin.venice.stats.VeniceMetricsConfig;
11+
import com.linkedin.venice.stats.VeniceMetricsRepository;
12+
import com.linkedin.venice.utils.OpenTelemetryDataTestUtils;
13+
import io.opentelemetry.api.common.Attributes;
14+
import io.opentelemetry.sdk.metrics.data.MetricData;
15+
import io.opentelemetry.sdk.testing.exporter.InMemoryMetricReader;
16+
import io.tehuti.metrics.MetricsRepository;
17+
import java.time.Clock;
18+
import java.util.Collection;
19+
import org.testng.annotations.AfterMethod;
20+
import org.testng.annotations.BeforeMethod;
21+
import org.testng.annotations.Test;
22+
23+
24+
public class NativeMetadataRepositoryStatsOtelTest {
25+
private static final String TEST_METRIC_PREFIX = "server";
26+
private static final String METRIC_NAME = METADATA_CACHE_STALENESS.getMetricEntity().getMetricName();
27+
28+
private InMemoryMetricReader inMemoryMetricReader;
29+
private VeniceMetricsRepository metricsRepository;
30+
private Clock mockClock;
31+
private NativeMetadataRepositoryStats stats;
32+
33+
@BeforeMethod
34+
public void setUp() {
35+
inMemoryMetricReader = InMemoryMetricReader.create();
36+
metricsRepository = new VeniceMetricsRepository(
37+
new VeniceMetricsConfig.Builder().setMetricPrefix(TEST_METRIC_PREFIX)
38+
.setMetricEntities(SERVER_METRIC_ENTITIES)
39+
.setEmitOtelMetrics(true)
40+
.setOtelAdditionalMetricsReader(inMemoryMetricReader)
41+
.build());
42+
mockClock = mock(Clock.class);
43+
doReturn(1000L).when(mockClock).millis();
44+
stats = new NativeMetadataRepositoryStats(metricsRepository, "test", mockClock);
45+
}
46+
47+
@AfterMethod
48+
public void tearDown() {
49+
if (metricsRepository != null) {
50+
metricsRepository.close();
51+
}
52+
}
53+
54+
@Test
55+
public void testPerStoreStaleness() {
56+
stats.updateCacheTimestamp("store-a", 500);
57+
stats.updateCacheTimestamp("store-b", 800);
58+
59+
// store-a staleness = 1000 - 500 = 500ms
60+
validateGauge(500, "store-a");
61+
// store-b staleness = 1000 - 800 = 200ms
62+
validateGauge(200, "store-b");
63+
}
64+
65+
@Test
66+
public void testStalenessUpdatesOnClockAdvance() {
67+
stats.updateCacheTimestamp("store-a", 900);
68+
69+
validateGauge(100, "store-a");
70+
71+
// Clock advances
72+
doReturn(2000L).when(mockClock).millis();
73+
validateGauge(1100, "store-a");
74+
}
75+
76+
@Test
77+
public void testStalenessUpdatesOnCacheRefresh() {
78+
stats.updateCacheTimestamp("store-a", 500);
79+
validateGauge(500, "store-a");
80+
81+
// Store metadata refreshed — staleness drops
82+
stats.updateCacheTimestamp("store-a", 900);
83+
validateGauge(100, "store-a");
84+
}
85+
86+
@Test
87+
public void testRemovedStoreReportsNaN() {
88+
stats.updateCacheTimestamp("store-a", 500);
89+
validateGauge(500, "store-a");
90+
91+
stats.removeCacheTimestamp("store-a");
92+
// After removal, callback returns NaN (no data). Validate directly since the
93+
// tolerance-based helper doesn't support NaN comparison (NaN != NaN in IEEE 754).
94+
validateGaugeAbsent("store-a");
95+
}
96+
97+
@Test
98+
public void testMultiStoreIsolation() {
99+
stats.updateCacheTimestamp("store-a", 200);
100+
stats.updateCacheTimestamp("store-b", 900);
101+
102+
validateGauge(800, "store-a");
103+
validateGauge(100, "store-b");
104+
}
105+
106+
@Test
107+
public void testNoNpeWhenOtelDisabled() {
108+
try (VeniceMetricsRepository disabledRepo = new VeniceMetricsRepository(
109+
new VeniceMetricsConfig.Builder().setMetricPrefix(TEST_METRIC_PREFIX).setEmitOtelMetrics(false).build())) {
110+
NativeMetadataRepositoryStats stats = new NativeMetadataRepositoryStats(disabledRepo, "test", mockClock);
111+
stats.updateCacheTimestamp("store-a", 500);
112+
stats.removeCacheTimestamp("store-a");
113+
}
114+
}
115+
116+
@Test
117+
public void testNoNpeWhenPlainMetricsRepository() {
118+
NativeMetadataRepositoryStats stats = new NativeMetadataRepositoryStats(new MetricsRepository(), "test", mockClock);
119+
stats.updateCacheTimestamp("store-a", 500);
120+
stats.removeCacheTimestamp("store-a");
121+
}
122+
123+
/**
124+
* Verifies that a removed store emits no OTel data point. The callback returns NaN
125+
* which the OTel SDK drops entirely — the metric data point disappears from collection.
126+
*/
127+
private void validateGaugeAbsent(String storeName) {
128+
Collection<MetricData> metricsData = inMemoryMetricReader.collectAllMetrics();
129+
String fullMetricName = "venice." + TEST_METRIC_PREFIX + "." + METRIC_NAME;
130+
boolean hasDataPoint = metricsData.stream()
131+
.filter(m -> m.getName().equals(fullMetricName))
132+
.flatMap(m -> m.getDoubleGaugeData().getPoints().stream())
133+
.anyMatch(p -> p.getAttributes().equals(buildAttributes(storeName)));
134+
assertFalse(hasDataPoint, "Expected no data point for removed store: " + storeName);
135+
}
136+
137+
private void validateGauge(double expectedValue, String storeName) {
138+
OpenTelemetryDataTestUtils.validateDoublePointDataFromGauge(
139+
inMemoryMetricReader,
140+
expectedValue,
141+
0.01,
142+
buildAttributes(storeName),
143+
METRIC_NAME,
144+
TEST_METRIC_PREFIX);
145+
}
146+
147+
private static Attributes buildAttributes(String storeName) {
148+
return Attributes.builder().put(VENICE_STORE_NAME.getDimensionNameInDefaultFormat(), storeName).build();
149+
}
150+
}

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(), 154, "Expected 154 unique metric entities");
2626
}
2727

2828
/**

0 commit comments

Comments
 (0)