diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/metrics/StreamsMetricsImpl.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/metrics/StreamsMetricsImpl.java index a05dbd9d2f00e..e4e370888a08b 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/metrics/StreamsMetricsImpl.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/metrics/StreamsMetricsImpl.java @@ -44,7 +44,6 @@ import java.util.LinkedList; import java.util.Map; import java.util.Objects; -import java.util.Optional; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentMap; import java.util.concurrent.TimeUnit; @@ -220,11 +219,12 @@ public final Sensor clientLevelSensor(final String sensorName, final Sensor... parents) { synchronized (clientLevelSensors) { final String fullSensorName = CLIENT_LEVEL_GROUP + SENSOR_NAME_DELIMITER + sensorName; - return Optional.ofNullable(metrics.getSensor(fullSensorName)) - .orElseGet(() -> { - clientLevelSensors.push(fullSensorName); - return metrics.sensor(fullSensorName, recordingLevel, parents); - }); + final Sensor sensor = metrics.getSensor(fullSensorName); + if (sensor == null) { + clientLevelSensors.push(fullSensorName); + return metrics.sensor(fullSensorName, recordingLevel, parents); + } + return sensor; } } @@ -234,12 +234,7 @@ public final Sensor threadLevelSensor(final String threadId, final Sensor... parents) { final String key = threadSensorPrefix(threadId); synchronized (threadLevelSensors) { - final String fullSensorName = key + SENSOR_NAME_DELIMITER + sensorName; - return Optional.ofNullable(metrics.getSensor(fullSensorName)) - .orElseGet(() -> { - threadLevelSensors.computeIfAbsent(key, ignored -> new LinkedList<>()).push(fullSensorName); - return metrics.sensor(fullSensorName, recordingLevel, parents); - }); + return getSensors(threadLevelSensors, sensorName, key, recordingLevel, parents); } } @@ -327,12 +322,7 @@ public final Sensor taskLevelSensor(final String threadId, final Sensor... parents) { final String key = taskSensorPrefix(threadId, taskId); synchronized (taskLevelSensors) { - final String fullSensorName = key + SENSOR_NAME_DELIMITER + sensorName; - return Optional.ofNullable(metrics.getSensor(fullSensorName)) - .orElseGet(() -> { - taskLevelSensors.computeIfAbsent(key, ignored -> new LinkedList<>()).push(fullSensorName); - return metrics.sensor(fullSensorName, recordingLevel, parents); - }); + return getSensors(taskLevelSensors, sensorName, key, recordingLevel, parents); } } @@ -359,12 +349,7 @@ public Sensor nodeLevelSensor(final String threadId, final Sensor... parents) { final String key = nodeSensorPrefix(threadId, taskId, processorNodeName); synchronized (nodeLevelSensors) { - final String fullSensorName = key + SENSOR_NAME_DELIMITER + sensorName; - return Optional.ofNullable(metrics.getSensor(fullSensorName)) - .orElseGet(() -> { - nodeLevelSensors.computeIfAbsent(key, ignored -> new LinkedList<>()).push(fullSensorName); - return metrics.sensor(fullSensorName, recordingLevel, parents); - }); + return getSensors(nodeLevelSensors, sensorName, key, recordingLevel, parents); } } @@ -393,12 +378,7 @@ public Sensor cacheLevelSensor(final String threadId, final Sensor... parents) { final String key = cacheSensorPrefix(threadId, taskName, storeName); synchronized (cacheLevelSensors) { - final String fullSensorName = key + SENSOR_NAME_DELIMITER + sensorName; - return Optional.ofNullable(metrics.getSensor(fullSensorName)) - .orElseGet(() -> { - cacheLevelSensors.computeIfAbsent(key, ignored -> new LinkedList<>()).push(fullSensorName); - return metrics.sensor(fullSensorName, recordingLevel, parents); - }); + return getSensors(cacheLevelSensors, sensorName, key, recordingLevel, parents); } } @@ -437,9 +417,6 @@ public final Sensor storeLevelSensor(final String taskId, final RecordingLevel recordingLevel, final Sensor... parents) { final String key = storeSensorPrefix(Thread.currentThread().getName(), taskId, storeName); - final String fullSensorName = key + SENSOR_NAME_DELIMITER + sensorName; - final Sensor sensor = metrics.getSensor(fullSensorName); - if (sensor == null) { // since the keys in the map storeLevelSensors contain the name of the current thread and threads only // access keys in which their name is contained, the value in the maps do not need to be thread safe // and we can use a LinkedList here. @@ -447,10 +424,7 @@ public final Sensor storeLevelSensor(final String taskId, // that contain its name. Similar is true for the other metric levels. Thread-level metrics need some // special attention, since they are created before the thread is constructed. The creation of those // metrics could be moved into the run() method of the thread. - storeLevelSensors.computeIfAbsent(key, ignored -> new LinkedList<>()).push(fullSensorName); - return metrics.sensor(fullSensorName, recordingLevel, parents); - } - return sensor; + return getSensors(storeLevelSensors, sensorName, key, recordingLevel, parents); } public void addStoreLevelMutableMetric(final String taskId, @@ -926,6 +900,20 @@ public static T maybeMeasureLatency(final Supplier actionToMeasure, } } + private Sensor getSensors(final Map> sensors, + final String sensorName, + final String key, + final RecordingLevel recordingLevel, + final Sensor... parents) { + final String fullSensorName = key + SENSOR_NAME_DELIMITER + sensorName; + final Sensor sensor = metrics.getSensor(fullSensorName); + if (sensor == null) { + sensors.computeIfAbsent(key, ignored -> new LinkedList<>()).push(fullSensorName); + return metrics.sensor(fullSensorName, recordingLevel, parents); + } + return sensor; + } + /** * Deletes a sensor and its parents, if any */