Skip to content

Commit b979f09

Browse files
committed
[server] Add OTel metrics to DiskHealthStats
Add 1 ASYNC_GAUGE metric: disk.health.status (1=healthy, 0=unhealthy) with CLUSTER_NAME dimension. Tehuti and OTel share a single LongSupplier callback polling DiskHealthCheckService.isDiskHealthy(). clusterName added to constructor, passed from VeniceServer.
1 parent ae14805 commit b979f09

7 files changed

Lines changed: 204 additions & 19 deletions

File tree

Lines changed: 38 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,38 @@
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 com.linkedin.venice.stats.DiskHealthStats}.
16+
*/
17+
public enum DiskHealthOtelMetricEntity implements ModuleMetricEntityInterface {
18+
DISK_HEALTH_STATUS(
19+
"disk.health.status", MetricType.ASYNC_GAUGE, MetricUnit.NUMBER,
20+
"Disk health status: 1 if healthy, 0 if unhealthy", setOf(VENICE_CLUSTER_NAME)
21+
);
22+
23+
private final MetricEntity metricEntity;
24+
25+
DiskHealthOtelMetricEntity(
26+
String metricName,
27+
MetricType metricType,
28+
MetricUnit unit,
29+
String description,
30+
Set<VeniceMetricsDimensions> dimensions) {
31+
this.metricEntity = new MetricEntity(metricName, metricType, unit, description, dimensions);
32+
}
33+
34+
@Override
35+
public MetricEntity getMetricEntity() {
36+
return metricEntity;
37+
}
38+
}

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

4950
public static final Collection<MetricEntity> SERVER_METRIC_ENTITIES =
Lines changed: 34 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,34 @@
1+
package com.linkedin.davinci.stats;
2+
3+
import static com.linkedin.davinci.stats.DiskHealthOtelMetricEntity.DISK_HEALTH_STATUS;
4+
import static com.linkedin.venice.stats.dimensions.VeniceMetricsDimensions.VENICE_CLUSTER_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 DiskHealthOtelMetricEntityTest {
17+
@Test
18+
public void testMetricEntities() {
19+
new ModuleMetricEntityTestFixture<>(DiskHealthOtelMetricEntity.class, expectedDefinitions()).assertAll();
20+
}
21+
22+
private static Map<DiskHealthOtelMetricEntity, MetricEntityExpectation> expectedDefinitions() {
23+
Map<DiskHealthOtelMetricEntity, MetricEntityExpectation> map = new HashMap<>();
24+
map.put(
25+
DISK_HEALTH_STATUS,
26+
new MetricEntityExpectation(
27+
"disk.health.status",
28+
MetricType.ASYNC_GAUGE,
29+
MetricUnit.NUMBER,
30+
"Disk health status: 1 if healthy, 0 if unhealthy",
31+
setOf(VENICE_CLUSTER_NAME)));
32+
return map;
33+
}
34+
}

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(), 155, "Expected 155 unique metric entities");
25+
assertEquals(SERVER_METRIC_ENTITIES.size(), 156, "Expected 156 unique metric entities");
2626
}
2727

2828
/**

services/venice-server/src/main/java/com/linkedin/venice/server/VeniceServer.java

Lines changed: 5 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -442,7 +442,11 @@ private List<AbstractVeniceService> createServices() {
442442
serverConfig.getLogContext());
443443
services.add(diskHealthCheckService);
444444
// create stats for disk health check service
445-
new DiskHealthStats(metricsRepository, diskHealthCheckService, "disk_health_check_service");
445+
new DiskHealthStats(
446+
metricsRepository,
447+
diskHealthCheckService,
448+
"disk_health_check_service",
449+
clusterConfig.getClusterName());
446450

447451
final Optional<ResourceReadUsageTracker> resourceReadUsageTracker;
448452
if (serverConfig.isOptimizeDatabaseForBackupVersionEnabled()) {
Lines changed: 33 additions & 16 deletions
Original file line numberDiff line numberDiff line change
@@ -1,33 +1,50 @@
11
package com.linkedin.venice.stats;
22

3+
import static com.linkedin.davinci.stats.DiskHealthOtelMetricEntity.DISK_HEALTH_STATUS;
4+
35
import com.linkedin.davinci.storage.DiskHealthCheckService;
6+
import com.linkedin.venice.stats.dimensions.VeniceMetricsDimensions;
7+
import com.linkedin.venice.stats.metrics.AsyncMetricEntityStateBase;
8+
import io.opentelemetry.api.common.Attributes;
49
import io.tehuti.metrics.MetricsRepository;
5-
import io.tehuti.metrics.Sensor;
610
import io.tehuti.metrics.stats.AsyncGauge;
11+
import java.util.Map;
12+
import java.util.function.LongSupplier;
713

814

915
/**
10-
* {@code DiskHealthStats} measures the disk health conditions based on the periodic tests ran by the {@link DiskHealthCheckService}.
16+
* {@code DiskHealthStats} measures the disk health conditions based on the periodic tests ran by
17+
* the {@link DiskHealthCheckService}. Reports 1 if healthy, 0 if unhealthy.
18+
*
19+
* <p>Tehuti and OTel both poll the same {@link DiskHealthCheckService#isDiskHealthy()} method.
20+
* They cannot share a single registration because Tehuti uses {@link AsyncGauge} (polled by
21+
* Tehuti's async executor) while OTel uses {@link AsyncMetricEntityStateBase} (polled by the
22+
* OTel SDK's PeriodicMetricReader).
1123
*/
1224
public class DiskHealthStats extends AbstractVeniceStats {
13-
private DiskHealthCheckService diskHealthCheckService;
14-
15-
private Sensor diskHealthSensor;
16-
1725
public DiskHealthStats(
1826
MetricsRepository metricsRepository,
1927
DiskHealthCheckService diskHealthCheckService,
20-
String name) {
28+
String name,
29+
String clusterName) {
2130
super(metricsRepository, name);
22-
this.diskHealthCheckService = diskHealthCheckService;
2331

24-
diskHealthSensor = registerSensor(new AsyncGauge((ignored, ignored2) -> {
25-
if (this.diskHealthCheckService.isDiskHealthy()) {
26-
// report 1 if the disk in this host is healthy; otherwise, report 0.
27-
return 1;
28-
} else {
29-
return 0;
30-
}
31-
}, "disk_healthy"));
32+
// Shared callback: 1 = healthy, 0 = unhealthy
33+
LongSupplier healthCallback = () -> diskHealthCheckService.isDiskHealthy() ? 1 : 0;
34+
35+
// Tehuti: AsyncGauge
36+
registerSensor(new AsyncGauge((ignored, ignored2) -> healthCallback.getAsLong(), "disk_healthy"));
37+
38+
// OTel: ASYNC_GAUGE with CLUSTER_NAME
39+
OpenTelemetryMetricsSetup.OpenTelemetryMetricsSetupInfo otelData =
40+
OpenTelemetryMetricsSetup.builder(metricsRepository).setClusterName(clusterName).build();
41+
Map<VeniceMetricsDimensions, String> baseDimensionsMap = otelData.getBaseDimensionsMap();
42+
Attributes baseAttributes = otelData.getBaseAttributes();
43+
AsyncMetricEntityStateBase.create(
44+
DISK_HEALTH_STATUS.getMetricEntity(),
45+
otelData.getOtelRepository(),
46+
baseDimensionsMap,
47+
baseAttributes,
48+
healthCallback);
3249
}
3350
}
Lines changed: 91 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,91 @@
1+
package com.linkedin.venice.stats;
2+
3+
import static com.linkedin.davinci.stats.DiskHealthOtelMetricEntity.DISK_HEALTH_STATUS;
4+
import static com.linkedin.davinci.stats.ServerMetricEntity.SERVER_METRIC_ENTITIES;
5+
import static com.linkedin.venice.stats.dimensions.VeniceMetricsDimensions.VENICE_CLUSTER_NAME;
6+
import static org.mockito.Mockito.doReturn;
7+
import static org.mockito.Mockito.mock;
8+
9+
import com.linkedin.davinci.storage.DiskHealthCheckService;
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 DiskHealthStatsTest {
20+
private static final String TEST_METRIC_PREFIX = "server";
21+
private static final String TEST_CLUSTER_NAME = "test-cluster";
22+
private static final String METRIC_NAME = DISK_HEALTH_STATUS.getMetricEntity().getMetricName();
23+
24+
private InMemoryMetricReader inMemoryMetricReader;
25+
private VeniceMetricsRepository metricsRepository;
26+
private DiskHealthCheckService mockService;
27+
28+
@BeforeMethod
29+
public void setUp() {
30+
inMemoryMetricReader = InMemoryMetricReader.create();
31+
metricsRepository = new VeniceMetricsRepository(
32+
new VeniceMetricsConfig.Builder().setMetricPrefix(TEST_METRIC_PREFIX)
33+
.setMetricEntities(SERVER_METRIC_ENTITIES)
34+
.setEmitOtelMetrics(true)
35+
.setOtelAdditionalMetricsReader(inMemoryMetricReader)
36+
.build());
37+
mockService = mock(DiskHealthCheckService.class);
38+
doReturn(true).when(mockService).isDiskHealthy();
39+
new DiskHealthStats(metricsRepository, mockService, "disk_health", TEST_CLUSTER_NAME);
40+
}
41+
42+
@AfterMethod
43+
public void tearDown() {
44+
if (metricsRepository != null) {
45+
metricsRepository.close();
46+
}
47+
}
48+
49+
@Test
50+
public void testHealthyDiskReportsOne() {
51+
OpenTelemetryDataTestUtils
52+
.validateLongPointDataFromGauge(inMemoryMetricReader, 1, buildAttributes(), METRIC_NAME, TEST_METRIC_PREFIX);
53+
}
54+
55+
@Test
56+
public void testUnhealthyDiskReportsZero() {
57+
doReturn(false).when(mockService).isDiskHealthy();
58+
59+
OpenTelemetryDataTestUtils
60+
.validateLongPointDataFromGauge(inMemoryMetricReader, 0, buildAttributes(), METRIC_NAME, TEST_METRIC_PREFIX);
61+
}
62+
63+
@Test
64+
public void testGaugeUpdatesOnHealthChange() {
65+
OpenTelemetryDataTestUtils
66+
.validateLongPointDataFromGauge(inMemoryMetricReader, 1, buildAttributes(), METRIC_NAME, TEST_METRIC_PREFIX);
67+
68+
// Disk becomes unhealthy
69+
doReturn(false).when(mockService).isDiskHealthy();
70+
71+
OpenTelemetryDataTestUtils
72+
.validateLongPointDataFromGauge(inMemoryMetricReader, 0, buildAttributes(), METRIC_NAME, TEST_METRIC_PREFIX);
73+
}
74+
75+
@Test
76+
public void testNoNpeWhenOtelDisabled() {
77+
try (VeniceMetricsRepository disabledRepo = new VeniceMetricsRepository(
78+
new VeniceMetricsConfig.Builder().setMetricPrefix(TEST_METRIC_PREFIX).setEmitOtelMetrics(false).build())) {
79+
new DiskHealthStats(disabledRepo, mockService, "disk_health", TEST_CLUSTER_NAME);
80+
}
81+
}
82+
83+
@Test
84+
public void testNoNpeWhenPlainMetricsRepository() {
85+
new DiskHealthStats(new MetricsRepository(), mockService, "disk_health", TEST_CLUSTER_NAME);
86+
}
87+
88+
private static Attributes buildAttributes() {
89+
return Attributes.builder().put(VENICE_CLUSTER_NAME.getDimensionNameInDefaultFormat(), TEST_CLUSTER_NAME).build();
90+
}
91+
}

0 commit comments

Comments
 (0)