Skip to content
Merged
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
24 changes: 12 additions & 12 deletions clients/src/main/java/org/apache/kafka/common/network/Selector.java
Original file line number Diff line number Diff line change
Expand Up @@ -1114,9 +1114,10 @@ public void close() {

class SelectorMetrics implements AutoCloseable {
private final Metrics metrics;
private final String metricGrpPrefix;
private final Map<String, String> metricTags;
private final boolean metricsPerConnection;
private final String metricGrpName;
private final String perConnectionMetricGrpName;

public final Sensor connectionClosed;
public final Sensor connectionCreated;
Expand All @@ -1142,10 +1143,10 @@ class SelectorMetrics implements AutoCloseable {

public SelectorMetrics(Metrics metrics, String metricGrpPrefix, Map<String, String> metricTags, boolean metricsPerConnection) {
this.metrics = metrics;
this.metricGrpPrefix = metricGrpPrefix;
this.metricTags = metricTags;
this.metricsPerConnection = metricsPerConnection;
String metricGrpName = metricGrpPrefix + "-metrics";
this.metricGrpName = metricGrpPrefix + "-metrics";
this.perConnectionMetricGrpName = metricGrpPrefix + "-node-metrics";
StringBuilder tagsSuffix = new StringBuilder();

for (Map.Entry<String, String> tag: metricTags.entrySet()) {
Expand Down Expand Up @@ -1293,34 +1294,33 @@ public void maybeRegisterConnectionMetrics(String connectionId) {
String nodeRequestName = "node-" + connectionId + ".requests-sent";
Sensor nodeRequest = this.metrics.getSensor(nodeRequestName);
if (nodeRequest == null) {
String metricGrpName = metricGrpPrefix + "-node-metrics";
Map<String, String> tags = new LinkedHashMap<>(metricTags);
tags.put("node-id", "node-" + connectionId);

nodeRequest = sensor(nodeRequestName);
nodeRequest.add(createMeter(metrics, metricGrpName, tags, new WindowedCount(), "request", "requests sent"));
MetricName metricName = metrics.metricName("request-size-avg", metricGrpName, "The average size of requests sent.", tags);
nodeRequest.add(createMeter(metrics, perConnectionMetricGrpName, tags, new WindowedCount(), "request", "requests sent"));
MetricName metricName = metrics.metricName("request-size-avg", perConnectionMetricGrpName, "The average size of requests sent.", tags);
nodeRequest.add(metricName, new Avg());
metricName = metrics.metricName("request-size-max", metricGrpName, "The maximum size of any request sent.", tags);
metricName = metrics.metricName("request-size-max", perConnectionMetricGrpName, "The maximum size of any request sent.", tags);
nodeRequest.add(metricName, new Max());

String bytesSentName = "node-" + connectionId + ".bytes-sent";
Sensor bytesSent = sensor(bytesSentName);
bytesSent.add(createMeter(metrics, metricGrpName, tags, "outgoing-byte", "outgoing bytes"));
bytesSent.add(createMeter(metrics, perConnectionMetricGrpName, tags, "outgoing-byte", "outgoing bytes"));

String nodeResponseName = "node-" + connectionId + ".responses-received";
Sensor nodeResponse = sensor(nodeResponseName);
nodeResponse.add(createMeter(metrics, metricGrpName, tags, new WindowedCount(), "response", "responses received"));
nodeResponse.add(createMeter(metrics, perConnectionMetricGrpName, tags, new WindowedCount(), "response", "responses received"));

String bytesReceivedName = "node-" + connectionId + ".bytes-received";
Sensor bytesReceive = sensor(bytesReceivedName);
bytesReceive.add(createMeter(metrics, metricGrpName, tags, "incoming-byte", "incoming bytes"));
bytesReceive.add(createMeter(metrics, perConnectionMetricGrpName, tags, "incoming-byte", "incoming bytes"));

String nodeTimeName = "node-" + connectionId + ".latency";
Sensor nodeRequestTime = sensor(nodeTimeName);
metricName = metrics.metricName("request-latency-avg", metricGrpName, tags);
metricName = metrics.metricName("request-latency-avg", perConnectionMetricGrpName, tags);
nodeRequestTime.add(metricName, new Avg());
metricName = metrics.metricName("request-latency-max", metricGrpName, tags);
metricName = metrics.metricName("request-latency-max", perConnectionMetricGrpName, tags);
nodeRequestTime.add(metricName, new Max());
}
}
Expand Down