From 6ccf8fdc263fec14cc711b1b833519bee0a3a71b Mon Sep 17 00:00:00 2001 From: lingniu Date: Wed, 1 Jul 2026 06:48:33 +0800 Subject: [PATCH] fix: stop kafka consumers without concurrent close --- .../sink/mq/KafkaEnvelopeConsumerRunner.java | 26 +++++++-- .../sink/mq/KafkaEnvelopeConsumerWorker.java | 4 ++ .../mq/KafkaEnvelopeConsumerRunnerTest.java | 53 +++++++++++++++++++ 3 files changed, 78 insertions(+), 5 deletions(-) create mode 100644 modules/sinks/sink-mq/src/test/java/com/lingniu/ingest/sink/mq/KafkaEnvelopeConsumerRunnerTest.java diff --git a/modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/mq/KafkaEnvelopeConsumerRunner.java b/modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/mq/KafkaEnvelopeConsumerRunner.java index dadd2679..3c68b380 100644 --- a/modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/mq/KafkaEnvelopeConsumerRunner.java +++ b/modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/mq/KafkaEnvelopeConsumerRunner.java @@ -95,6 +95,9 @@ public final class KafkaEnvelopeConsumerRunner implements SmartLifecycle, AutoCl try { worker.pollOnce(pollTimeout); } catch (RuntimeException ex) { + if (!running.get()) { + break; + } // Kafka/处理器异常不让 Spring 生命周期退出;退避后继续消费, // 具体坏消息由 EnvelopeConsumerProcessor 写入 DLQ。 log.warn("Kafka envelope consumer poll failed; the worker will retry after backoff", ex); @@ -116,17 +119,20 @@ public final class KafkaEnvelopeConsumerRunner implements SmartLifecycle, AutoCl if (!running.compareAndSet(true, false)) { return; } + wakeupWorkers(); if (executor != null) { - executor.shutdownNow(); - } - closeWorkers(); - if (executor != null) { + executor.shutdown(); try { - executor.awaitTermination(5, TimeUnit.SECONDS); + if (!executor.awaitTermination(5, TimeUnit.SECONDS)) { + executor.shutdownNow(); + executor.awaitTermination(5, TimeUnit.SECONDS); + } } catch (InterruptedException ex) { + executor.shutdownNow(); Thread.currentThread().interrupt(); } } + closeWorkers(); } @Override @@ -163,4 +169,14 @@ public final class KafkaEnvelopeConsumerRunner implements SmartLifecycle, AutoCl worker.close(); } } + + private void wakeupWorkers() { + List current = workers; + if (current == null) { + return; + } + for (KafkaEnvelopeConsumerWorker worker : current) { + worker.wakeup(); + } + } } diff --git a/modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/mq/KafkaEnvelopeConsumerWorker.java b/modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/mq/KafkaEnvelopeConsumerWorker.java index 9e928eb9..d640afca 100644 --- a/modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/mq/KafkaEnvelopeConsumerWorker.java +++ b/modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/mq/KafkaEnvelopeConsumerWorker.java @@ -64,6 +64,10 @@ public final class KafkaEnvelopeConsumerWorker implements AutoCloseable { return processed; } + public void wakeup() { + consumer.wakeup(); + } + @Override public void close() { consumer.close(); diff --git a/modules/sinks/sink-mq/src/test/java/com/lingniu/ingest/sink/mq/KafkaEnvelopeConsumerRunnerTest.java b/modules/sinks/sink-mq/src/test/java/com/lingniu/ingest/sink/mq/KafkaEnvelopeConsumerRunnerTest.java new file mode 100644 index 00000000..f8cd2c8a --- /dev/null +++ b/modules/sinks/sink-mq/src/test/java/com/lingniu/ingest/sink/mq/KafkaEnvelopeConsumerRunnerTest.java @@ -0,0 +1,53 @@ +package com.lingniu.ingest.sink.mq; + +import org.junit.jupiter.api.Test; + +import java.time.Duration; +import java.util.List; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.TimeUnit; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.Mockito.doAnswer; +import static org.mockito.Mockito.mock; + +class KafkaEnvelopeConsumerRunnerTest { + + @Test + void stopWakesPollLoopAndClosesWorkerAfterPollThreadExits() throws Exception { + KafkaEnvelopeConsumerWorker worker = mock(KafkaEnvelopeConsumerWorker.class); + CountDownLatch pollEntered = new CountDownLatch(1); + CountDownLatch wakeupCalled = new CountDownLatch(1); + CountDownLatch pollReturned = new CountDownLatch(1); + CountDownLatch closeCalled = new CountDownLatch(1); + + doAnswer(invocation -> { + pollEntered.countDown(); + assertThat(wakeupCalled.await(1, TimeUnit.SECONDS)).isTrue(); + pollReturned.countDown(); + return 0; + }).when(worker).pollOnce(any()); + doAnswer(invocation -> { + wakeupCalled.countDown(); + return null; + }).when(worker).wakeup(); + doAnswer(invocation -> { + assertThat(pollReturned.getCount()).isZero(); + closeCalled.countDown(); + return null; + }).when(worker).close(); + KafkaEnvelopeConsumerRunner runner = new KafkaEnvelopeConsumerRunner( + List.of(worker), + Duration.ofSeconds(30), + Duration.ofMillis(1), + false); + + runner.start(); + assertThat(pollEntered.await(1, TimeUnit.SECONDS)).isTrue(); + runner.stop(); + + assertThat(wakeupCalled.getCount()).isZero(); + assertThat(closeCalled.getCount()).isZero(); + } +}