From d165e04d11fc1eca473b950bd0d04cd51157d173 Mon Sep 17 00:00:00 2001 From: Jeff Kim Date: Fri, 17 May 2024 15:00:18 -0400 Subject: [PATCH 1/3] fix --- .../group/runtime/MultiThreadedEventProcessor.java | 3 --- .../runtime/MultiThreadedEventProcessorTest.java | 11 ++++++----- 2 files changed, 6 insertions(+), 8 deletions(-) diff --git a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/runtime/MultiThreadedEventProcessor.java b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/runtime/MultiThreadedEventProcessor.java index 20337a4677581..31fa52ea7d158 100644 --- a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/runtime/MultiThreadedEventProcessor.java +++ b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/runtime/MultiThreadedEventProcessor.java @@ -112,9 +112,6 @@ public MultiThreadedEventProcessor( */ private class EventProcessorThread extends Thread { private final Logger log; - private long pollStartMs; - private long timeSinceLastPollMs; - private long lastPollMs; EventProcessorThread( String name diff --git a/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/runtime/MultiThreadedEventProcessorTest.java b/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/runtime/MultiThreadedEventProcessorTest.java index 313e46dfc0473..5aaf4e5c6821a 100644 --- a/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/runtime/MultiThreadedEventProcessorTest.java +++ b/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/runtime/MultiThreadedEventProcessorTest.java @@ -61,8 +61,11 @@ public DelayEventAccumulator(Time time, long takeDelayMs) { @Override public CoordinatorEvent take() { - time.sleep(takeDelayMs); - return super.take(); + CoordinatorEvent event = super.take(); + if (event != null) { + time.sleep(takeDelayMs); + } + return event; } } @@ -475,9 +478,7 @@ public void testRecordThreadIdleRatio() throws Exception { doAnswer(invocation -> { long threadIdleTime = idleTimeCaptured.getValue(); assertEquals(100, threadIdleTime); - synchronized (recordedIdleTimesMs) { - recordedIdleTimesMs.add(threadIdleTime); - } + recordedIdleTimesMs.add(threadIdleTime); return null; }).when(mockRuntimeMetrics).recordThreadIdleTime(idleTimeCaptured.capture()); From 5a34639a3bf57cb8d306de2b5bb1827870b884e2 Mon Sep 17 00:00:00 2001 From: Jeff Kim Date: Fri, 17 May 2024 16:31:12 -0400 Subject: [PATCH 2/3] comment --- .../group/runtime/MultiThreadedEventProcessorTest.java | 2 ++ 1 file changed, 2 insertions(+) diff --git a/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/runtime/MultiThreadedEventProcessorTest.java b/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/runtime/MultiThreadedEventProcessorTest.java index 5aaf4e5c6821a..66e85ee063de1 100644 --- a/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/runtime/MultiThreadedEventProcessorTest.java +++ b/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/runtime/MultiThreadedEventProcessorTest.java @@ -478,6 +478,8 @@ public void testRecordThreadIdleRatio() throws Exception { doAnswer(invocation -> { long threadIdleTime = idleTimeCaptured.getValue(); assertEquals(100, threadIdleTime); + + // No synchronization required as the test uses a single event processor thread. recordedIdleTimesMs.add(threadIdleTime); return null; }).when(mockRuntimeMetrics).recordThreadIdleTime(idleTimeCaptured.capture()); From 044319a2968f1c18f0919c687bb875f81a6499b4 Mon Sep 17 00:00:00 2001 From: Jeff Kim Date: Mon, 20 May 2024 12:41:39 -0400 Subject: [PATCH 3/3] address comments --- .../group/runtime/MultiThreadedEventProcessorTest.java | 4 +--- 1 file changed, 1 insertion(+), 3 deletions(-) diff --git a/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/runtime/MultiThreadedEventProcessorTest.java b/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/runtime/MultiThreadedEventProcessorTest.java index 66e85ee063de1..0f2801daec3b9 100644 --- a/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/runtime/MultiThreadedEventProcessorTest.java +++ b/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/runtime/MultiThreadedEventProcessorTest.java @@ -62,9 +62,7 @@ public DelayEventAccumulator(Time time, long takeDelayMs) { @Override public CoordinatorEvent take() { CoordinatorEvent event = super.take(); - if (event != null) { - time.sleep(takeDelayMs); - } + time.sleep(takeDelayMs); return event; } }