From 8c2b3256003e90c375e81ceffc6b1dc250ab4ab9 Mon Sep 17 00:00:00 2001 From: Gary Russell Date: Thu, 3 Nov 2022 09:43:53 -0400 Subject: [PATCH] Fix Missing Re-Interrupts --- .../java/org/springframework/kafka/core/KafkaAdmin.java | 9 ++++++--- .../kafka/listener/AbstractMessageListenerContainer.java | 6 +++++- 2 files changed, 11 insertions(+), 4 deletions(-) diff --git a/spring-kafka/src/main/java/org/springframework/kafka/core/KafkaAdmin.java b/spring-kafka/src/main/java/org/springframework/kafka/core/KafkaAdmin.java index b0f7914a..b3164cfb 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/core/KafkaAdmin.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/core/KafkaAdmin.java @@ -201,12 +201,15 @@ public class KafkaAdmin extends KafkaResourceFactory addOrModifyTopicsIfNeeded(adminClient, newTopics); return true; } - catch (Exception e) { + catch (InterruptedException ex) { + Thread.currentThread().interrupt(); + } + catch (Exception ex) { if (!this.initializingContext || this.fatalIfBrokerNotAvailable) { - throw new IllegalStateException("Could not configure topics", e); + throw new IllegalStateException("Could not configure topics", ex); } else { - LOGGER.error(e, "Could not configure topics"); + LOGGER.error(ex, "Could not configure topics"); } } finally { diff --git a/spring-kafka/src/main/java/org/springframework/kafka/listener/AbstractMessageListenerContainer.java b/spring-kafka/src/main/java/org/springframework/kafka/listener/AbstractMessageListenerContainer.java index 242548f0..09156fb2 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/listener/AbstractMessageListenerContainer.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/listener/AbstractMessageListenerContainer.java @@ -488,7 +488,11 @@ public abstract class AbstractMessageListenerContainer entry.getValue().get(this.topicCheckTimeout, TimeUnit.SECONDS); return false; } - catch (@SuppressWarnings("unused") Exception e) { + catch (InterruptedException ex) { + Thread.currentThread().interrupt(); + return true; + } + catch (@SuppressWarnings("unused") Exception ex) { return true; } })