Skip to content

Commit 3248bf9

Browse files
authored
[server] Add OTel metrics to BackupVersionOptimizationServiceStats (linkedin#2726)
Add 1 OTel COUNTER metric version.backup.optimization.reopen_count with STORE_NAME + OPERATION_OUTCOME (SUCCESS/FAIL) dimensions. Behavioral change 1 (affects both Tehuti and OTel): Error recording moved from per-store to per-partition. This gives finer granularity and matches the success path, which already records per-partition. Behavioral change 2 (Tehuti registration timing): Tehuti sensors ( backup_version_database_optimization, backup_version_data_optimization_error) are registered lazily on the first per-store recording.
1 parent 9cd7569 commit 3248bf9

15 files changed

Lines changed: 570 additions & 21 deletions

File tree

Lines changed: 42 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,42 @@
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_OPERATION_OUTCOME;
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
18+
* {@link com.linkedin.venice.stats.BackupVersionOptimizationServiceStats}.
19+
*/
20+
public enum BackupVersionOptimizationOtelMetricEntity implements ModuleMetricEntityInterface {
21+
REOPEN_COUNT(
22+
"version.backup.optimization.reopen_count", MetricType.COUNTER, MetricUnit.NUMBER,
23+
"Count of backup version storage partition reopens by outcome (success or fail)",
24+
setOf(VENICE_CLUSTER_NAME, VENICE_STORE_NAME, VENICE_OPERATION_OUTCOME)
25+
);
26+
27+
private final MetricEntity metricEntity;
28+
29+
BackupVersionOptimizationOtelMetricEntity(
30+
String metricName,
31+
MetricType metricType,
32+
MetricUnit unit,
33+
String description,
34+
Set<VeniceMetricsDimensions> dimensions) {
35+
this.metricEntity = new MetricEntity(metricName, metricType, unit, description, dimensions);
36+
}
37+
38+
@Override
39+
public MetricEntity getMetricEntity() {
40+
return metricEntity;
41+
}
42+
}

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
@@ -47,7 +47,8 @@ public static List<Class<? extends ModuleMetricEntityInterface>> getMetricEntity
4747
DiskHealthOtelMetricEntity.class,
4848
NativeMetadataRepositoryOtelMetricEntity.class,
4949
ServerLoadOtelMetricEntity.class,
50-
RocksDBStatsOtelMetricEntity.class);
50+
RocksDBStatsOtelMetricEntity.class,
51+
BackupVersionOptimizationOtelMetricEntity.class);
5152
}
5253

5354
public static final Collection<MetricEntity> SERVER_METRIC_ENTITIES =
Lines changed: 37 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,37 @@
1+
package com.linkedin.davinci.stats;
2+
3+
import static com.linkedin.davinci.stats.BackupVersionOptimizationOtelMetricEntity.REOPEN_COUNT;
4+
import static com.linkedin.venice.stats.dimensions.VeniceMetricsDimensions.VENICE_CLUSTER_NAME;
5+
import static com.linkedin.venice.stats.dimensions.VeniceMetricsDimensions.VENICE_OPERATION_OUTCOME;
6+
import static com.linkedin.venice.stats.dimensions.VeniceMetricsDimensions.VENICE_STORE_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 BackupVersionOptimizationOtelMetricEntityTest {
19+
@Test
20+
public void testMetricEntities() {
21+
new ModuleMetricEntityTestFixture<>(BackupVersionOptimizationOtelMetricEntity.class, expectedDefinitions())
22+
.assertAll();
23+
}
24+
25+
private static Map<BackupVersionOptimizationOtelMetricEntity, MetricEntityExpectation> expectedDefinitions() {
26+
Map<BackupVersionOptimizationOtelMetricEntity, MetricEntityExpectation> map = new HashMap<>();
27+
map.put(
28+
REOPEN_COUNT,
29+
new MetricEntityExpectation(
30+
"version.backup.optimization.reopen_count",
31+
MetricType.COUNTER,
32+
MetricUnit.NUMBER,
33+
"Count of backup version storage partition reopens by outcome (success or fail)",
34+
setOf(VENICE_CLUSTER_NAME, VENICE_STORE_NAME, VENICE_OPERATION_OUTCOME)));
35+
return map;
36+
}
37+
}

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

2828
/**

internal/venice-client-common/src/main/java/com/linkedin/venice/stats/dimensions/VeniceMetricsDimensions.java

Lines changed: 4 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -168,7 +168,10 @@ public enum VeniceMetricsDimensions {
168168
VENICE_ROCKSDB_LEVEL("venice.rocksdb.level"),
169169

170170
/** {@link VeniceRocksDBBlockCacheComponent} RocksDB block cache component type. */
171-
VENICE_ROCKSDB_BLOCK_CACHE_COMPONENT("venice.rocksdb.block_cache_component");
171+
VENICE_ROCKSDB_BLOCK_CACHE_COMPONENT("venice.rocksdb.block_cache_component"),
172+
173+
/** {@link VeniceOperationOutcome} Generic operation outcome: success or fail. */
174+
VENICE_OPERATION_OUTCOME("venice.operation.outcome");
172175

173176
private final String[] dimensionName = new String[VeniceOpenTelemetryMetricNamingFormat.SIZE];
174177

Lines changed: 14 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,14 @@
1+
package com.linkedin.venice.stats.dimensions;
2+
3+
/** Generic operation outcome: success or failure. */
4+
public enum VeniceOperationOutcome implements VeniceDimensionInterface {
5+
/** Operation completed successfully. */
6+
SUCCESS,
7+
/** Operation failed. */
8+
FAIL;
9+
10+
@Override
11+
public VeniceMetricsDimensions getDimensionName() {
12+
return VeniceMetricsDimensions.VENICE_OPERATION_OUTCOME;
13+
}
14+
}

internal/venice-client-common/src/test/java/com/linkedin/venice/stats/dimensions/VeniceMetricsDimensionsTest.java

Lines changed: 9 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -172,6 +172,9 @@ public void testGetDimensionNameInSnakeCase() {
172172
case VENICE_ROCKSDB_BLOCK_CACHE_COMPONENT:
173173
assertEquals(dimension.getDimensionName(format), "venice.rocksdb.block_cache_component");
174174
break;
175+
case VENICE_OPERATION_OUTCOME:
176+
assertEquals(dimension.getDimensionName(format), "venice.operation.outcome");
177+
break;
175178
default:
176179
throw new IllegalArgumentException("Unknown dimension: " + dimension);
177180
}
@@ -342,6 +345,9 @@ public void testGetDimensionNameInCamelCase() {
342345
case VENICE_ROCKSDB_BLOCK_CACHE_COMPONENT:
343346
assertEquals(dimension.getDimensionName(format), "venice.rocksdb.blockCacheComponent");
344347
break;
348+
case VENICE_OPERATION_OUTCOME:
349+
assertEquals(dimension.getDimensionName(format), "venice.operation.outcome");
350+
break;
345351
default:
346352
throw new IllegalArgumentException("Unknown dimension: " + dimension);
347353
}
@@ -512,6 +518,9 @@ public void testGetDimensionNameInPascalCase() {
512518
case VENICE_ROCKSDB_BLOCK_CACHE_COMPONENT:
513519
assertEquals(dimension.getDimensionName(format), "Venice.Rocksdb.BlockCacheComponent");
514520
break;
521+
case VENICE_OPERATION_OUTCOME:
522+
assertEquals(dimension.getDimensionName(format), "Venice.Operation.Outcome");
523+
break;
515524
default:
516525
throw new IllegalArgumentException("Unknown dimension: " + dimension);
517526
}
Lines changed: 20 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,20 @@
1+
package com.linkedin.venice.stats.dimensions;
2+
3+
import com.linkedin.venice.utils.CollectionUtils;
4+
import java.util.Map;
5+
import org.testng.annotations.Test;
6+
7+
8+
public class VeniceOperationOutcomeTest {
9+
@Test
10+
public void testDimensionInterface() {
11+
Map<VeniceOperationOutcome, String> expectedValues = CollectionUtils.<VeniceOperationOutcome, String>mapBuilder()
12+
.put(VeniceOperationOutcome.SUCCESS, "success")
13+
.put(VeniceOperationOutcome.FAIL, "fail")
14+
.build();
15+
new VeniceDimensionTestFixture<>(
16+
VeniceOperationOutcome.class,
17+
VeniceMetricsDimensions.VENICE_OPERATION_OUTCOME,
18+
expectedValues).assertAll();
19+
}
20+
}

internal/venice-test-common/src/integrationTest/java/com/linkedin/venice/endToEnd/TestBackupVersionDatabaseOptimization.java

Lines changed: 11 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -142,18 +142,25 @@ public void execute() {
142142
// Verify whether the backup version database optimization happens or not.
143143
VeniceServerWrapper serverWrapper = venice.getVeniceServers().get(0);
144144
MetricsRepository metricsRepository = serverWrapper.getMetricsRepository();
145-
Metric optimizationMetric = metricsRepository
146-
.getMetric(".BackupVersionOptimizationService--backup_version_database_optimization.OccurrenceRate");
147-
Metric rocksdbMetric = metricsRepository.getMetric(".RocksDBMemoryStats--rocksdb.num-immutable-mem-table.Gauge");
148145

149146
// N.B.: The optimization is performed by a periodic background task, so it cannot be expected to have already
150-
// completed as soon as we get here.
147+
// completed as soon as we get here. Fetch the metric handle inside the assertion lambda — the Tehuti sensor is
148+
// registered lazily on first per-store recording, so getMetric(...) returns null until the first optimization
149+
// fires.
151150
TestUtils.waitForNonDeterministicAssertion(10, TimeUnit.SECONDS, true, () -> {
151+
Metric optimizationMetric = metricsRepository
152+
.getMetric(".BackupVersionOptimizationService--backup_version_database_optimization.OccurrenceRate");
153+
assertNotNull(
154+
optimizationMetric,
155+
"Backup version optimization Tehuti sensor should be registered after the first optimization");
152156
assertTrue(optimizationMetric.value() > 0, "Backup version database optimization should happen");
153157
/**
154158
* This assertion is used to make sure {@link com.linkedin.davinci.stats.RocksDBMemoryStats} won't crash after
155159
* reopening the backup version.
156160
*/
161+
Metric rocksdbMetric =
162+
metricsRepository.getMetric(".RocksDBMemoryStats--rocksdb.num-immutable-mem-table.Gauge");
163+
assertNotNull(rocksdbMetric, "RocksDBMemoryStats num-immutable-mem-table gauge must be registered");
157164
rocksdbMetric.value();
158165
});
159166
}

services/venice-server/src/main/java/com/linkedin/venice/cleaner/BackupVersionOptimizationService.java

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -155,17 +155,17 @@ private Runnable getOptimizationRunnable() {
155155
for (int partitionId: partitionIdSet) {
156156
try {
157157
engine.reopenStoragePartition(partitionId);
158-
stats.recordBackupVersionDatabaseOptimization();
158+
stats.recordBackupVersionDatabaseOptimization(storeName);
159159
} catch (Exception e) {
160160
LOGGER.error(
161161
"Failed to optimize database for topic-partition: {}",
162162
Utils.getReplicaId(resourceName, partitionId),
163163
e);
164+
stats.recordBackupVersionDatabaseOptimizationError(storeName);
164165
errored = true;
165166
}
166167
}
167168
if (errored) {
168-
stats.recordBackupVersionDatabaseOptimizationError();
169169
LOGGER.warn(
170170
"Encountered issue when optimizing database for resource: {}, "
171171
+ "and please check the above logs to find more details, and will retry the optimization in next iteration",

0 commit comments

Comments
 (0)