cadonna commented on a change in pull request #9614:
URL: https://github.com/apache/kafka/pull/9614#discussion_r529313022



##########
File path: 
streams/src/main/java/org/apache/kafka/streams/processor/internals/metrics/StreamsMetricsImpl.java
##########
@@ -214,6 +215,20 @@ public RocksDBMetricsRecordingTrigger 
rocksDBMetricsRecordingTrigger() {
         }
     }
 
+    public final Sensor clientLevelSensor(final String sensorName,
+                                          final RecordingLevel recordingLevel,
+                                          final Sensor... parents) {
+        final String key = CLIENT_LEVEL_GROUP;
+        synchronized (clientLevelSensors) {
+            final String fullSensorName = key + SENSOR_NAME_DELIMITER + 
sensorName;
+            return Optional.ofNullable(metrics.getSensor(fullSensorName))
+                .orElseGet(() -> {
+                    clientLevelSensors.push(fullSensorName);
+                    return metrics.sensor(fullSensorName, recordingLevel, 
parents);
+                });

Review comment:
       Although, we use this in other methods, I think the following is a bit 
simpler to read:
   ```suggestion
           final Sensor sensor = metrics.getSensor(fullSensorName);
           if (sensor == null) {
               clientLevelSensors.push(fullSensorName);
               return metrics.sensor(fullSensorName, recordingLevel, parents);
           }
           return sensor;
   ```

##########
File path: 
streams/src/main/java/org/apache/kafka/streams/processor/internals/metrics/StreamsMetricsImpl.java
##########
@@ -214,6 +215,20 @@ public RocksDBMetricsRecordingTrigger 
rocksDBMetricsRecordingTrigger() {
         }
     }
 
+    public final Sensor clientLevelSensor(final String sensorName,
+                                          final RecordingLevel recordingLevel,
+                                          final Sensor... parents) {
+        final String key = CLIENT_LEVEL_GROUP;
+        synchronized (clientLevelSensors) {
+            final String fullSensorName = key + SENSOR_NAME_DELIMITER + 
sensorName;

Review comment:
       You can inline the value of variable `key` here and remove `key`.
   
   ```suggestion
               final String fullSensorName = CLIENT_LEVEL_GROUP + 
SENSOR_NAME_DELIMITER + sensorName;
   ```

##########
File path: 
streams/src/test/java/org/apache/kafka/streams/processor/internals/metrics/StreamsMetricsImplTest.java
##########
@@ -577,6 +579,38 @@ public void shouldGetExistingCacheLevelSensor() {
         assertThat(actualSensor, is(equalToObject(sensor)));
     }
 
+    @Test
+    public void shouldGetNewClientLevelSensor() {
+        final Metrics metrics = mock(Metrics.class);
+        final RecordingLevel recordingLevel = RecordingLevel.INFO;
+        setupGetNewSensorTest(metrics, recordingLevel);
+        final StreamsMetricsImpl streamsMetrics = new 
StreamsMetricsImpl(metrics, CLIENT_ID, VERSION, time);
+
+        final Sensor actualSensor = streamsMetrics.clientLevelSensor(
+            SENSOR_NAME_1,
+            recordingLevel
+        );

Review comment:
       ```suggestion
           final Sensor actualSensor = 
streamsMetrics.clientLevelSensor(SENSOR_NAME_1, recordingLevel);
   ```

##########
File path: 
streams/src/test/java/org/apache/kafka/streams/processor/internals/StreamThreadTest.java
##########
@@ -2438,6 +2440,60 @@ public void shouldConstructAdminMetrics() {
         assertEquals(testMetricName, 
adminClientMetrics.get(testMetricName).metricName());
     }
 
+    @Test
+    public void shouldCountFailedStreamThread() {
+        verifyFailedStreamThread(false);
+        verifyFailedStreamThread(true);
+    }
+
+    public void verifyFailedStreamThread(final boolean shouldFail) {
+        final Consumer<byte[], byte[]> consumer = 
EasyMock.createNiceMock(Consumer.class);
+        final TaskManager taskManager = 
EasyMock.createNiceMock(TaskManager.class);
+        final StreamsMetricsImpl streamsMetrics =
+            new StreamsMetricsImpl(metrics, CLIENT_ID, 
StreamsConfig.METRICS_LATEST, mockTime);
+        final StreamThread thread = new StreamThread(
+            mockTime,
+            config,
+            null,
+            consumer,
+            consumer,
+            null,
+            null,
+            taskManager,
+            streamsMetrics,
+            internalTopologyBuilder,
+            CLIENT_ID,
+            new LogContext(""),
+            new AtomicInteger(),
+            new AtomicLong(Long.MAX_VALUE),
+            null,
+            e -> { }
+        ) {
+            @Override
+            void runOnce() {
+                setState(StreamThread.State.PENDING_SHUTDOWN);
+                if (shouldFail) {
+                    throw new 
StreamsException(Thread.currentThread().getName());
+                }
+            }
+        };
+        expect(taskManager.activeTaskMap()).andReturn(Collections.emptyMap());
+        expect(taskManager.standbyTaskMap()).andReturn(Collections.emptyMap());
+
+        taskManager.process(anyInt(), anyObject());
+        EasyMock.expectLastCall().andThrow(new 
StreamsException(Thread.currentThread().getName()));
+

Review comment:
       You can remove these lines. They are dead code.

##########
File path: 
streams/src/test/java/org/apache/kafka/streams/processor/internals/StreamThreadTest.java
##########
@@ -2438,6 +2440,60 @@ public void shouldConstructAdminMetrics() {
         assertEquals(testMetricName, 
adminClientMetrics.get(testMetricName).metricName());
     }
 
+    @Test
+    public void shouldCountFailedStreamThread() {
+        verifyFailedStreamThread(false);
+        verifyFailedStreamThread(true);
+    }
+
+    public void verifyFailedStreamThread(final boolean shouldFail) {
+        final Consumer<byte[], byte[]> consumer = 
EasyMock.createNiceMock(Consumer.class);
+        final TaskManager taskManager = 
EasyMock.createNiceMock(TaskManager.class);
+        final StreamsMetricsImpl streamsMetrics =
+            new StreamsMetricsImpl(metrics, CLIENT_ID, 
StreamsConfig.METRICS_LATEST, mockTime);
+        final StreamThread thread = new StreamThread(
+            mockTime,
+            config,
+            null,
+            consumer,
+            consumer,
+            null,
+            null,
+            taskManager,
+            streamsMetrics,
+            internalTopologyBuilder,
+            CLIENT_ID,
+            new LogContext(""),
+            new AtomicInteger(),
+            new AtomicLong(Long.MAX_VALUE),
+            null,
+            e -> { }
+        ) {
+            @Override
+            void runOnce() {
+                setState(StreamThread.State.PENDING_SHUTDOWN);
+                if (shouldFail) {
+                    throw new 
StreamsException(Thread.currentThread().getName());
+                }
+            }
+        };
+        expect(taskManager.activeTaskMap()).andReturn(Collections.emptyMap());
+        expect(taskManager.standbyTaskMap()).andReturn(Collections.emptyMap());
+
+        taskManager.process(anyInt(), anyObject());
+        EasyMock.expectLastCall().andThrow(new 
StreamsException(Thread.currentThread().getName()));
+
+        EasyMock.replay(taskManager);
+
+        thread.updateThreadMetadata("metadata");
+        thread.setState(StreamThread.State.STARTING);
+        thread.runLoop();
+
+        final Metric failedThreads = 
StreamsTestUtils.getMetricByName(metrics.metrics(), "failed-stream-threads", 
"stream-metrics");
+
+        assertEquals(shouldFail ? 1.0 : 0.0, failedThreads.metricValue());

Review comment:
       ```suggestion
           EasyMock.replay(taskManager);
           thread.updateThreadMetadata("metadata");
           thread.setState(StreamThread.State.STARTING);
           
           thread.runLoop();
   
           final Metric failedThreads = 
StreamsTestUtils.getMetricByName(metrics.metrics(), "failed-stream-threads", 
"stream-metrics");
           assertThat(failedThreads.metricValue(), is(shouldFail ? 1.0 : 0.0));
   ```

##########
File path: 
streams/src/test/java/org/apache/kafka/streams/processor/internals/StreamThreadTest.java
##########
@@ -2438,6 +2440,60 @@ public void shouldConstructAdminMetrics() {
         assertEquals(testMetricName, 
adminClientMetrics.get(testMetricName).metricName());
     }
 
+    @Test
+    public void shouldCountFailedStreamThread() {
+        verifyFailedStreamThread(false);
+        verifyFailedStreamThread(true);
+    }
+
+    public void verifyFailedStreamThread(final boolean shouldFail) {

Review comment:
       I would rename this method to 
`runAndVerifyFailedStreamThreadRecording()`.

##########
File path: 
streams/src/test/java/org/apache/kafka/streams/processor/internals/metrics/StreamsMetricsImplTest.java
##########
@@ -577,6 +579,38 @@ public void shouldGetExistingCacheLevelSensor() {
         assertThat(actualSensor, is(equalToObject(sensor)));
     }
 
+    @Test
+    public void shouldGetNewClientLevelSensor() {
+        final Metrics metrics = mock(Metrics.class);
+        final RecordingLevel recordingLevel = RecordingLevel.INFO;
+        setupGetNewSensorTest(metrics, recordingLevel);
+        final StreamsMetricsImpl streamsMetrics = new 
StreamsMetricsImpl(metrics, CLIENT_ID, VERSION, time);
+
+        final Sensor actualSensor = streamsMetrics.clientLevelSensor(
+            SENSOR_NAME_1,
+            recordingLevel
+        );
+
+        verify(metrics);
+        assertThat(actualSensor, is(equalToObject(sensor)));
+    }
+
+    @Test
+    public void shouldGetExistingClientLevelSensor() {
+        final Metrics metrics = mock(Metrics.class);
+        final RecordingLevel recordingLevel = RecordingLevel.INFO;
+        setupGetExistingSensorTest(metrics);
+        final StreamsMetricsImpl streamsMetrics = new 
StreamsMetricsImpl(metrics, CLIENT_ID, VERSION, time);
+
+        final Sensor actualSensor = streamsMetrics.clientLevelSensor(
+            SENSOR_NAME_1,
+            recordingLevel
+        );

Review comment:
       ```suggestion
           final Sensor actualSensor = 
streamsMetrics.clientLevelSensor(SENSOR_NAME_1, recordingLevel);
   ```

##########
File path: 
streams/src/test/java/org/apache/kafka/streams/processor/internals/StreamThreadTest.java
##########
@@ -2438,6 +2440,60 @@ public void shouldConstructAdminMetrics() {
         assertEquals(testMetricName, 
adminClientMetrics.get(testMetricName).metricName());
     }
 
+    @Test
+    public void shouldCountFailedStreamThread() {
+        verifyFailedStreamThread(false);
+        verifyFailedStreamThread(true);
+    }

Review comment:
       Could you please specify two separate unit tests? One could be named 
`shouldRecordFailedStreamThread()` and the other 
`shouldNotRecordFailedStreamThread()`.

##########
File path: 
streams/src/test/java/org/apache/kafka/streams/internals/metrics/ClientMetricsTest.java
##########
@@ -99,6 +121,29 @@ public void shouldAddAliveStreamThreadsMetric() {
         );
     }
 
+    @Test
+    public void shouldGetFailedStreamThreadsSensor() {
+        final String name = "failed-stream-threads";
+        final String description = "The number of failed stream threads since 
the start of the Kafka Streams client";
+        expect(streamsMetrics.clientLevelSensor(name, 
RecordingLevel.INFO)).andReturn(expectedSensor);
+        expect(streamsMetrics.clientLevelTagMap()).andReturn(tagMap);
+        StreamsMetricsImpl.addSumMetricToSensor(
+            expectedSensor,
+            CLIENT_LEVEL_GROUP,
+            tagMap,
+            name,
+            false,
+            description
+        );
+

Review comment:
       Please remove empty line.

##########
File path: 
streams/src/test/java/org/apache/kafka/streams/processor/internals/metrics/StreamsMetricsImplTest.java
##########
@@ -629,17 +663,20 @@ private void setupRemoveSensorsTest(final Metrics metrics,
     }
 
     @Test
-    public void shouldRemoveClientLevelMetrics() {
+    public void shouldRemoveClientLevelMetricsAndSensors() {
         final Metrics metrics = niceMock(Metrics.class);
         final StreamsMetricsImpl streamsMetrics = new 
StreamsMetricsImpl(metrics, CLIENT_ID, VERSION, time);
-        addSensorsOnAllLevels(metrics, streamsMetrics);
+        final Capture<String> sensorKeys = addSensorsOnAllLevels(metrics, 
streamsMetrics);
         resetToDefault(metrics);
-        expect(metrics.removeMetric(metricName1)).andStubReturn(null);
-        expect(metrics.removeMetric(metricName2)).andStubReturn(null);
-        replay(metrics);
 

Review comment:
       Please remove this empty line.

##########
File path: 
streams/src/test/java/org/apache/kafka/streams/processor/internals/metrics/StreamsMetricsImplTest.java
##########
@@ -629,17 +663,20 @@ private void setupRemoveSensorsTest(final Metrics metrics,
     }
 
     @Test
-    public void shouldRemoveClientLevelMetrics() {
+    public void shouldRemoveClientLevelMetricsAndSensors() {
         final Metrics metrics = niceMock(Metrics.class);
         final StreamsMetricsImpl streamsMetrics = new 
StreamsMetricsImpl(metrics, CLIENT_ID, VERSION, time);
-        addSensorsOnAllLevels(metrics, streamsMetrics);
+        final Capture<String> sensorKeys = addSensorsOnAllLevels(metrics, 
streamsMetrics);
         resetToDefault(metrics);
-        expect(metrics.removeMetric(metricName1)).andStubReturn(null);
-        expect(metrics.removeMetric(metricName2)).andStubReturn(null);
-        replay(metrics);
 
-        streamsMetrics.removeAllClientLevelMetrics();
+        metrics.removeSensor(sensorKeys.getValues().get(0));
+        metrics.removeSensor(sensorKeys.getValues().get(1));
+
+        
expect(metrics.removeMetric(metricName1)).andReturn(mock(KafkaMetric.class));
+        
expect(metrics.removeMetric(metricName2)).andReturn(mock(KafkaMetric.class));
+        replay(metrics);
 
+        streamsMetrics.removeAllClientLevelSensorsAndMetrics();
         verify(metrics);

Review comment:
       ```suggestion
           metrics.removeSensor(sensorKeys.getValues().get(0));
           metrics.removeSensor(sensorKeys.getValues().get(1));
           
expect(metrics.removeMetric(metricName1)).andReturn(mock(KafkaMetric.class));
           
expect(metrics.removeMetric(metricName2)).andReturn(mock(KafkaMetric.class));
           replay(metrics);
   
           streamsMetrics.removeAllClientLevelSensorsAndMetrics();
   
           verify(metrics);
   ```




----------------------------------------------------------------
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

For queries about this service, please contact Infrastructure at:
[email protected]


Reply via email to