From d745cf450daf46273004d4b6de5f600239e000be Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Daniel=20F=C4=85derski?= Date: Mon, 18 Aug 2025 17:47:14 +0200 Subject: [PATCH] stop subscription immediately even when there are pending retries --- .../consumers/consumer/BatchConsumer.java | 3 +- .../supervisor/process/ConsumerProcess.java | 44 ++++++++++--------- 2 files changed, 25 insertions(+), 22 deletions(-) diff --git a/hermes-consumers/src/main/java/pl/allegro/tech/hermes/consumers/consumer/BatchConsumer.java b/hermes-consumers/src/main/java/pl/allegro/tech/hermes/consumers/consumer/BatchConsumer.java index 04047ed47d..baf1e31b65 100644 --- a/hermes-consumers/src/main/java/pl/allegro/tech/hermes/consumers/consumer/BatchConsumer.java +++ b/hermes-consumers/src/main/java/pl/allegro/tech/hermes/consumers/consumer/BatchConsumer.java @@ -282,7 +282,8 @@ && shouldRetryOnClientError(retryClientErrors, result)) .withStopStrategy( attempt -> attempt.getDelaySinceFirstAttempt() > messageTtlMillis - || Thread.currentThread().isInterrupted()) + || Thread.currentThread().isInterrupted() + || !consuming) .withRetryListener( getRetryListener( result -> { diff --git a/hermes-consumers/src/main/java/pl/allegro/tech/hermes/consumers/supervisor/process/ConsumerProcess.java b/hermes-consumers/src/main/java/pl/allegro/tech/hermes/consumers/supervisor/process/ConsumerProcess.java index 821fbf7340..cbf6010296 100644 --- a/hermes-consumers/src/main/java/pl/allegro/tech/hermes/consumers/supervisor/process/ConsumerProcess.java +++ b/hermes-consumers/src/main/java/pl/allegro/tech/hermes/consumers/supervisor/process/ConsumerProcess.java @@ -67,17 +67,7 @@ public void run() { } finally { logger.info("Releasing consumer process thread of subscription {}", getSubscriptionName()); refreshHealthcheck(); - try { - stop(); - } catch (Exception exceptionWhileStopping) { - logger.error( - "An error occurred while stopping consumer process of subscription {}", - getSubscriptionName(), - exceptionWhileStopping); - } finally { - onConsumerStopped.accept(getSubscriptionName()); - Thread.currentThread().setName("consumer-released-thread"); - } + stop(); } } @@ -119,7 +109,7 @@ private void process(Signal signal) { "Stopping main loop for consumer {}. {}", signal.getTarget(), signal.getLogWithIdAndType()); - this.running = false; + stop(); break; case RETRANSMIT: retransmit(signal); @@ -159,15 +149,27 @@ private void start(Signal signal) { } private void stop() { - long startTime = clock.millis(); - logger.info("Stopping consumer for subscription {}", getSubscriptionName()); - - consumer.tearDown(); - - logger.info( - "Stopped consumer for subscription {} in {}ms", - getSubscriptionName(), - clock.millis() - startTime); + if (!running) { + return; + } + this.running = false; + try { + long startTime = clock.millis(); + logger.info("Stopping consumer for subscription {}", getSubscriptionName()); + consumer.tearDown(); + logger.info( + "Stopped consumer for subscription {} in {}ms", + getSubscriptionName(), + clock.millis() - startTime); + } catch (Exception exceptionWhileStopping) { + logger.error( + "An error occurred while stopping consumer process of subscription {}", + getSubscriptionName(), + exceptionWhileStopping); + } finally { + onConsumerStopped.accept(getSubscriptionName()); + Thread.currentThread().setName("consumer-released-thread"); + } } private void retransmit(Signal signal) {