From 5dbbea5e430fe52e3b5dcbc52d85dcd0d934a5c8 Mon Sep 17 00:00:00 2001 From: Luke Chen Date: Thu, 27 Apr 2023 18:09:37 +0800 Subject: [PATCH 1/2] improve log --- .../src/main/java/org/apache/kafka/queue/KafkaEventQueue.java | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/server-common/src/main/java/org/apache/kafka/queue/KafkaEventQueue.java b/server-common/src/main/java/org/apache/kafka/queue/KafkaEventQueue.java index 70f859d9f893f..76d52e2e63596 100644 --- a/server-common/src/main/java/org/apache/kafka/queue/KafkaEventQueue.java +++ b/server-common/src/main/java/org/apache/kafka/queue/KafkaEventQueue.java @@ -250,7 +250,7 @@ private void handleEvents() throws InterruptedException { continue; } else if (shuttingDown) { remove(eventContext); - toDeliver = new RejectedExecutionException(); + toDeliver = new RejectedExecutionException("The event queue is shutting down"); toRun = eventContext; continue; } @@ -300,7 +300,7 @@ Exception enqueue(EventContext eventContext, lock.lock(); try { if (shuttingDown) { - return new RejectedExecutionException(); + return new RejectedExecutionException("The event queue is shutting down"); } if (interrupted) { return new InterruptedException(); From 0a911a23c9648f26aacd85753edf05eb8d3e77da Mon Sep 17 00:00:00 2001 From: Luke Chen Date: Fri, 28 Apr 2023 10:33:00 +0800 Subject: [PATCH 2/2] add reason for InterruptedException --- .../main/java/org/apache/kafka/queue/KafkaEventQueue.java | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/server-common/src/main/java/org/apache/kafka/queue/KafkaEventQueue.java b/server-common/src/main/java/org/apache/kafka/queue/KafkaEventQueue.java index 76d52e2e63596..2fde8285dab4b 100644 --- a/server-common/src/main/java/org/apache/kafka/queue/KafkaEventQueue.java +++ b/server-common/src/main/java/org/apache/kafka/queue/KafkaEventQueue.java @@ -245,7 +245,7 @@ private void handleEvents() throws InterruptedException { continue; } else if (interrupted) { remove(eventContext); - toDeliver = new InterruptedException(); + toDeliver = new InterruptedException("The event handler thread is interrupted"); toRun = eventContext; continue; } else if (shuttingDown) { @@ -264,7 +264,7 @@ private void handleEvents() throws InterruptedException { } } else { if (interrupted) { - toDeliver = new InterruptedException(); + toDeliver = new InterruptedException("The event handler thread is interrupted"); } else { toDeliver = null; } @@ -303,7 +303,7 @@ Exception enqueue(EventContext eventContext, return new RejectedExecutionException("The event queue is shutting down"); } if (interrupted) { - return new InterruptedException(); + return new InterruptedException("The event handler thread is interrupted"); } OptionalLong existingDeadlineNs = OptionalLong.empty(); if (eventContext.tag != null) {