Skip to content

Commit 055d73a

Browse files
committed
[da-vinci][server] Add OTel metrics to ServerConnectionStats
3 OTel metrics added: - connection.active_count (UP_DOWN_COUNTER, by source) - connection.request_count (COUNTER, by source) - connection.setup_time (HISTOGRAM) New dimension: VeniceConnectionSource (ROUTER, CLIENT) ServerMetricEntity count: 132 -> 135
1 parent b961f65 commit 055d73a

13 files changed

Lines changed: 584 additions & 17 deletions

File tree

Lines changed: 49 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,49 @@
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_CONNECTION_SOURCE;
5+
import static com.linkedin.venice.utils.Utils.setOf;
6+
7+
import com.linkedin.venice.stats.dimensions.VeniceMetricsDimensions;
8+
import com.linkedin.venice.stats.metrics.MetricEntity;
9+
import com.linkedin.venice.stats.metrics.MetricType;
10+
import com.linkedin.venice.stats.metrics.MetricUnit;
11+
import com.linkedin.venice.stats.metrics.ModuleMetricEntityInterface;
12+
import java.util.Set;
13+
14+
15+
/**
16+
* OTel metric entity definitions for {@link com.linkedin.venice.stats.ServerConnectionStats}.
17+
*/
18+
public enum ServerConnectionOtelMetricEntity implements ModuleMetricEntityInterface {
19+
CONNECTION_ACTIVE_COUNT(
20+
"connection.active_count", MetricType.UP_DOWN_COUNTER, MetricUnit.NUMBER,
21+
"Active connection count by source (router or client)", setOf(VENICE_CLUSTER_NAME, VENICE_CONNECTION_SOURCE)
22+
),
23+
24+
CONNECTION_REQUEST_COUNT(
25+
"connection.request_count", MetricType.COUNTER, MetricUnit.NUMBER, "Connection establishment requests by source",
26+
setOf(VENICE_CLUSTER_NAME, VENICE_CONNECTION_SOURCE)
27+
),
28+
29+
CONNECTION_SETUP_TIME(
30+
"connection.setup_time", MetricType.HISTOGRAM, MetricUnit.MILLISECOND,
31+
"SSL handshake setup latency from channel init to completion", setOf(VENICE_CLUSTER_NAME)
32+
);
33+
34+
private final MetricEntity metricEntity;
35+
36+
ServerConnectionOtelMetricEntity(
37+
String metricName,
38+
MetricType metricType,
39+
MetricUnit unit,
40+
String description,
41+
Set<VeniceMetricsDimensions> dimensions) {
42+
this.metricEntity = new MetricEntity(metricName, metricType, unit, description, dimensions);
43+
}
44+
45+
@Override
46+
public MetricEntity getMetricEntity() {
47+
return metricEntity;
48+
}
49+
}

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
@@ -38,7 +38,8 @@ public static List<Class<? extends ModuleMetricEntityInterface>> getMetricEntity
3838
HeartbeatMonitoringOtelMetricEntity.class,
3939
BlobTransferOtelMetricEntity.class,
4040
KafkaConsumerServiceOtelMetricEntity.class,
41-
RocksDBMemoryOtelMetricEntity.class);
41+
RocksDBMemoryOtelMetricEntity.class,
42+
ServerConnectionOtelMetricEntity.class);
4243
}
4344

4445
public static final Collection<MetricEntity> SERVER_METRIC_ENTITIES =
Lines changed: 53 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,53 @@
1+
package com.linkedin.davinci.stats;
2+
3+
import static com.linkedin.davinci.stats.ServerConnectionOtelMetricEntity.CONNECTION_ACTIVE_COUNT;
4+
import static com.linkedin.davinci.stats.ServerConnectionOtelMetricEntity.CONNECTION_REQUEST_COUNT;
5+
import static com.linkedin.davinci.stats.ServerConnectionOtelMetricEntity.CONNECTION_SETUP_TIME;
6+
import static com.linkedin.venice.stats.dimensions.VeniceMetricsDimensions.VENICE_CLUSTER_NAME;
7+
import static com.linkedin.venice.stats.dimensions.VeniceMetricsDimensions.VENICE_CONNECTION_SOURCE;
8+
import static com.linkedin.venice.utils.Utils.setOf;
9+
10+
import com.linkedin.venice.stats.metrics.MetricType;
11+
import com.linkedin.venice.stats.metrics.MetricUnit;
12+
import com.linkedin.venice.stats.metrics.ModuleMetricEntityTestFixture;
13+
import com.linkedin.venice.stats.metrics.ModuleMetricEntityTestFixture.MetricEntityExpectation;
14+
import java.util.HashMap;
15+
import java.util.Map;
16+
import org.testng.annotations.Test;
17+
18+
19+
public class ServerConnectionOtelMetricEntityTest {
20+
@Test
21+
public void testMetricEntities() {
22+
new ModuleMetricEntityTestFixture<>(ServerConnectionOtelMetricEntity.class, expectedDefinitions()).assertAll();
23+
}
24+
25+
private static Map<ServerConnectionOtelMetricEntity, MetricEntityExpectation> expectedDefinitions() {
26+
Map<ServerConnectionOtelMetricEntity, MetricEntityExpectation> map = new HashMap<>();
27+
map.put(
28+
CONNECTION_ACTIVE_COUNT,
29+
new MetricEntityExpectation(
30+
"connection.active_count",
31+
MetricType.UP_DOWN_COUNTER,
32+
MetricUnit.NUMBER,
33+
"Active connection count by source (router or client)",
34+
setOf(VENICE_CLUSTER_NAME, VENICE_CONNECTION_SOURCE)));
35+
map.put(
36+
CONNECTION_REQUEST_COUNT,
37+
new MetricEntityExpectation(
38+
"connection.request_count",
39+
MetricType.COUNTER,
40+
MetricUnit.NUMBER,
41+
"Connection establishment requests by source",
42+
setOf(VENICE_CLUSTER_NAME, VENICE_CONNECTION_SOURCE)));
43+
map.put(
44+
CONNECTION_SETUP_TIME,
45+
new MetricEntityExpectation(
46+
"connection.setup_time",
47+
MetricType.HISTOGRAM,
48+
MetricUnit.MILLISECOND,
49+
"SSL handshake setup latency from channel init to completion",
50+
setOf(VENICE_CLUSTER_NAME)));
51+
return map;
52+
}
53+
}

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

2828
/**
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+
/**
4+
* Dimension enum representing the source type of a connection to the server.
5+
* Maps to {@link VeniceMetricsDimensions#VENICE_CONNECTION_SOURCE}.
6+
*/
7+
public enum VeniceConnectionSource implements VeniceDimensionInterface {
8+
ROUTER, CLIENT;
9+
10+
@Override
11+
public VeniceMetricsDimensions getDimensionName() {
12+
return VeniceMetricsDimensions.VENICE_CONNECTION_SOURCE;
13+
}
14+
}

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
@@ -135,7 +135,10 @@ public enum VeniceMetricsDimensions {
135135
VENICE_HEARTBEAT_COMPONENT("venice.heartbeat.component"),
136136

137137
/** {@link VeniceConsumerPoolAction} Consumer pool action (subscribe, update_assignment). */
138-
VENICE_CONSUMER_POOL_ACTION("venice.consumer_pool.action");
138+
VENICE_CONSUMER_POOL_ACTION("venice.consumer_pool.action"),
139+
140+
/** {@link VeniceConnectionSource} Connection source type: router or client. */
141+
VENICE_CONNECTION_SOURCE("venice.connection.source");
139142

140143
private final String[] dimensionName = new String[VeniceOpenTelemetryMetricNamingFormat.SIZE];
141144

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 VeniceConnectionSourceTest {
9+
@Test
10+
public void testDimensionInterface() {
11+
Map<VeniceConnectionSource, String> expectedValues = CollectionUtils.<VeniceConnectionSource, String>mapBuilder()
12+
.put(VeniceConnectionSource.ROUTER, "router")
13+
.put(VeniceConnectionSource.CLIENT, "client")
14+
.build();
15+
new VeniceDimensionTestFixture<>(
16+
VeniceConnectionSource.class,
17+
VeniceMetricsDimensions.VENICE_CONNECTION_SOURCE,
18+
expectedValues).assertAll();
19+
}
20+
}

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
@@ -139,6 +139,9 @@ public void testGetDimensionNameInSnakeCase() {
139139
case VENICE_CONSUMER_POOL_ACTION:
140140
assertEquals(dimension.getDimensionName(format), "venice.consumer_pool.action");
141141
break;
142+
case VENICE_CONNECTION_SOURCE:
143+
assertEquals(dimension.getDimensionName(format), "venice.connection.source");
144+
break;
142145
default:
143146
throw new IllegalArgumentException("Unknown dimension: " + dimension);
144147
}
@@ -276,6 +279,9 @@ public void testGetDimensionNameInCamelCase() {
276279
case VENICE_CONSUMER_POOL_ACTION:
277280
assertEquals(dimension.getDimensionName(format), "venice.consumerPool.action");
278281
break;
282+
case VENICE_CONNECTION_SOURCE:
283+
assertEquals(dimension.getDimensionName(format), "venice.connection.source");
284+
break;
279285
default:
280286
throw new IllegalArgumentException("Unknown dimension: " + dimension);
281287
}
@@ -413,6 +419,9 @@ public void testGetDimensionNameInPascalCase() {
413419
case VENICE_CONSUMER_POOL_ACTION:
414420
assertEquals(dimension.getDimensionName(format), "Venice.ConsumerPool.Action");
415421
break;
422+
case VENICE_CONNECTION_SOURCE:
423+
assertEquals(dimension.getDimensionName(format), "Venice.Connection.Source");
424+
break;
416425
default:
417426
throw new IllegalArgumentException("Unknown dimension: " + dimension);
418427
}

services/venice-server/src/main/java/com/linkedin/venice/listener/HttpChannelInitializer.java

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -175,7 +175,7 @@ public HttpChannelInitializer(
175175
if (sslFactory.isPresent()) {
176176
this.serverConnectionStatsHandler = new ServerConnectionStatsHandler(
177177
this.identityParser,
178-
new ServerConnectionStats(metricsRepository, "server_connection_stats"),
178+
new ServerConnectionStats(metricsRepository, "server_connection_stats", serverConfig.getClusterName()),
179179
serverConfig.getRouterPrincipalName());
180180
} else {
181181
this.serverConnectionStatsHandler = null;
Lines changed: 76 additions & 13 deletions
Original file line numberDiff line numberDiff line change
@@ -1,9 +1,21 @@
11
package com.linkedin.venice.stats;
22

3+
import static com.linkedin.davinci.stats.ServerConnectionOtelMetricEntity.CONNECTION_ACTIVE_COUNT;
4+
import static com.linkedin.davinci.stats.ServerConnectionOtelMetricEntity.CONNECTION_REQUEST_COUNT;
5+
import static com.linkedin.davinci.stats.ServerConnectionOtelMetricEntity.CONNECTION_SETUP_TIME;
6+
7+
import com.linkedin.venice.stats.dimensions.VeniceConnectionSource;
8+
import com.linkedin.venice.stats.dimensions.VeniceMetricsDimensions;
9+
import com.linkedin.venice.stats.metrics.MetricEntityStateBase;
10+
import com.linkedin.venice.stats.metrics.MetricEntityStateOneEnum;
11+
import com.linkedin.venice.stats.metrics.TehutiMetricNameEnum;
312
import io.tehuti.metrics.MetricsRepository;
413
import io.tehuti.metrics.Sensor;
514
import io.tehuti.metrics.stats.AsyncGauge;
615
import io.tehuti.metrics.stats.OccurrenceRate;
16+
import java.util.Arrays;
17+
import java.util.Collections;
18+
import java.util.Map;
719
import java.util.concurrent.atomic.AtomicLong;
820

921

@@ -12,46 +24,97 @@ public class ServerConnectionStats extends AbstractVeniceStats {
1224
public static final String ROUTER_CONNECTION_COUNT_GAUGE = "router_connection_count";
1325
public static final String CLIENT_CONNECTION_REQUEST = "client_connection_request";
1426
public static final String CLIENT_CONNECTION_COUNT_GAUGE = "client_connection_count";
27+
public static final String CONNECTION_REQUEST = "connection_request";
1528
public static final String NEW_CONNECTION_SETUP_LATENCY = "new_connection_setup_latency";
1629

17-
private final Sensor routerConnectionRequestSensor;
18-
private final Sensor clientConnectionRequestSensor;
1930
private final Sensor connectionRequestSensor;
20-
private final Sensor newConnectionSetupLatencySensor;
2131

2232
private final AtomicLong routerConnectionCount = new AtomicLong();
2333
private final AtomicLong clientConnectionCount = new AtomicLong();
2434

25-
public ServerConnectionStats(MetricsRepository metricsRepository, String name) {
35+
// OTel: UP_DOWN_COUNTER, no Tehuti binding (Tehuti uses AsyncGauge on AtomicLong)
36+
private final MetricEntityStateOneEnum<VeniceConnectionSource> activeCountOtel;
37+
// OTel: shared instrument, each bound to its respective Tehuti OccurrenceRate sensor
38+
private final MetricEntityStateOneEnum<VeniceConnectionSource> routerRequestCountOtel;
39+
private final MetricEntityStateOneEnum<VeniceConnectionSource> clientRequestCountOtel;
40+
// OTel: joint Tehuti+OTel histogram
41+
private final MetricEntityStateBase setupTimeOtel;
42+
43+
enum TehutiMetricName implements TehutiMetricNameEnum {
44+
ROUTER_CONNECTION_REQUEST, CLIENT_CONNECTION_REQUEST, NEW_CONNECTION_SETUP_LATENCY
45+
}
46+
47+
public ServerConnectionStats(MetricsRepository metricsRepository, String name, String clusterName) {
2648
super(metricsRepository, name);
49+
50+
OpenTelemetryMetricsSetup.OpenTelemetryMetricsSetupInfo otelData =
51+
OpenTelemetryMetricsSetup.builder(metricsRepository).setClusterName(clusterName).build();
52+
VeniceOpenTelemetryMetricsRepository otelRepository = otelData.getOtelRepository();
53+
Map<VeniceMetricsDimensions, String> baseDimensionsMap = otelData.getBaseDimensionsMap();
54+
55+
// Tehuti AsyncGauges for active connection counts
2756
registerSensorIfAbsent(
2857
new AsyncGauge((ignored, ignored2) -> routerConnectionCount.get(), ROUTER_CONNECTION_COUNT_GAUGE));
29-
routerConnectionRequestSensor = registerSensorIfAbsent(ROUTER_CONNECTION_REQUEST, new OccurrenceRate());
3058
registerSensorIfAbsent(
3159
new AsyncGauge((ignored, ignored2) -> clientConnectionCount.get(), CLIENT_CONNECTION_COUNT_GAUGE));
32-
clientConnectionRequestSensor = registerSensorIfAbsent(CLIENT_CONNECTION_REQUEST, new OccurrenceRate());
33-
connectionRequestSensor = registerSensorIfAbsent("connection_request", new OccurrenceRate());
34-
newConnectionSetupLatencySensor = registerSensorIfAbsent(
35-
NEW_CONNECTION_SETUP_LATENCY,
36-
TehutiUtils.getPercentileStatWithAvgAndMax(getName(), NEW_CONNECTION_SETUP_LATENCY));
60+
61+
activeCountOtel = MetricEntityStateOneEnum.create(
62+
CONNECTION_ACTIVE_COUNT.getMetricEntity(),
63+
otelRepository,
64+
baseDimensionsMap,
65+
VeniceConnectionSource.class);
66+
67+
routerRequestCountOtel = MetricEntityStateOneEnum.create(
68+
CONNECTION_REQUEST_COUNT.getMetricEntity(),
69+
otelRepository,
70+
this::registerSensorIfAbsent,
71+
TehutiMetricName.ROUTER_CONNECTION_REQUEST,
72+
Collections.singletonList(new OccurrenceRate()),
73+
baseDimensionsMap,
74+
VeniceConnectionSource.class);
75+
76+
clientRequestCountOtel = MetricEntityStateOneEnum.create(
77+
CONNECTION_REQUEST_COUNT.getMetricEntity(),
78+
otelRepository,
79+
this::registerSensorIfAbsent,
80+
TehutiMetricName.CLIENT_CONNECTION_REQUEST,
81+
Collections.singletonList(new OccurrenceRate()),
82+
baseDimensionsMap,
83+
VeniceConnectionSource.class);
84+
85+
// Tehuti only — OTel total derived at query time
86+
connectionRequestSensor = registerSensorIfAbsent(CONNECTION_REQUEST, new OccurrenceRate());
87+
88+
setupTimeOtel = MetricEntityStateBase.create(
89+
CONNECTION_SETUP_TIME.getMetricEntity(),
90+
otelRepository,
91+
this::registerSensorIfAbsent,
92+
TehutiMetricName.NEW_CONNECTION_SETUP_LATENCY,
93+
Arrays.asList(TehutiUtils.getPercentileStatWithAvgAndMax(getName(), NEW_CONNECTION_SETUP_LATENCY)),
94+
baseDimensionsMap,
95+
otelData.getBaseAttributes());
3796
}
3897

3998
public void incrementRouterConnectionCount() {
4099
routerConnectionCount.incrementAndGet();
41-
routerConnectionRequestSensor.record(1);
100+
routerRequestCountOtel.record(1, VeniceConnectionSource.ROUTER);
101+
activeCountOtel.record(1, VeniceConnectionSource.ROUTER);
42102
}
43103

44104
public void decrementRouterConnectionCount() {
45105
routerConnectionCount.decrementAndGet();
106+
activeCountOtel.record(-1, VeniceConnectionSource.ROUTER);
46107
}
47108

48109
public void incrementClientConnectionCount() {
49110
clientConnectionCount.incrementAndGet();
50-
clientConnectionRequestSensor.record(1);
111+
clientRequestCountOtel.record(1, VeniceConnectionSource.CLIENT);
112+
activeCountOtel.record(1, VeniceConnectionSource.CLIENT);
51113
}
52114

53115
public void decrementClientConnectionCount() {
54116
clientConnectionCount.decrementAndGet();
117+
activeCountOtel.record(-1, VeniceConnectionSource.CLIENT);
55118
}
56119

57120
public void newConnectionRequest() {
@@ -62,6 +125,6 @@ public void newConnectionRequest() {
62125
* Record the latency from the start of channel initialization to SSL handshake completion.
63126
*/
64127
public void recordNewConnectionSetupLatency(double latencyMs) {
65-
newConnectionSetupLatencySensor.record(latencyMs);
128+
setupTimeOtel.record(latencyMs);
66129
}
67130
}

0 commit comments

Comments
 (0)