Skip to content
Merged
Changes from 1 commit
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
33 changes: 16 additions & 17 deletions clients/src/main/java/org/apache/kafka/common/network/Selector.java
Original file line number Diff line number Diff line change
Expand Up @@ -1027,9 +1027,10 @@ public int numStagedReceives(KafkaChannel channel) {

private class SelectorMetrics implements AutoCloseable {
private final Metrics metrics;
private final String metricGrpPrefix;
Comment thread
navina marked this conversation as resolved.
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 @@ -1051,10 +1052,10 @@ private 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 @@ -1143,10 +1144,10 @@ public SelectorMetrics(Metrics metrics, String metricGrpPrefix, Map<String, Stri

private Meter createMeter(Metrics metrics, String groupName, Map<String, String> metricTags,
SampledStat stat, String baseName, String descriptiveName) {
MetricName rateMetricName = metrics.metricName(baseName + "-rate", groupName,
String.format("The number of %s per second", descriptiveName), metricTags);
MetricName totalMetricName = metrics.metricName(baseName + "-total", groupName,
String.format("The total number of %s", descriptiveName), metricTags);
MetricName rateMetricName = metrics.metricName((baseName + "-rate").intern(), groupName,
Comment thread
smccauliff marked this conversation as resolved.
String.format("The number of %s per second", descriptiveName).intern(), metricTags);
MetricName totalMetricName = metrics.metricName((baseName + "-total").intern(), groupName,
String.format("The total number of %s", descriptiveName).intern(), metricTags);
if (stat == null)
return new Meter(rateMetricName, totalMetricName);
else
Expand Down Expand Up @@ -1180,29 +1181,27 @@ public void maybeRegisterConnectionMetrics(String connectionId) {
String nodeRequestName = "node-" + connectionId + ".bytes-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, "outgoing-byte", "outgoing bytes"));
nodeRequest.add(createMeter(metrics, metricGrpName, tags, new Count(), "request", "requests sent"));
MetricName metricName = metrics.metricName("request-size-avg", metricGrpName, "The average size of requests sent.", tags);
nodeRequest.add(createMeter(metrics, perConnectionMetricGrpName, tags, "outgoing-byte", "outgoing bytes"));
nodeRequest.add(createMeter(metrics, perConnectionMetricGrpName, tags, new Count(), "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 nodeResponseName = "node-" + connectionId + ".bytes-received";
Sensor nodeResponse = sensor(nodeResponseName);
nodeResponse.add(createMeter(metrics, metricGrpName, tags, "incoming-byte", "incoming bytes"));
nodeResponse.add(createMeter(metrics, metricGrpName, tags, new Count(), "response", "responses received"));
nodeResponse.add(createMeter(metrics, perConnectionMetricGrpName, tags, "incoming-byte", "incoming bytes"));
nodeResponse.add(createMeter(metrics, perConnectionMetricGrpName, tags, new Count(), "response", "responses received"));

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