diff --git a/spring-kafka/src/main/java/org/springframework/kafka/listener/KafkaMessageListenerContainer.java b/spring-kafka/src/main/java/org/springframework/kafka/listener/KafkaMessageListenerContainer.java index db5a34d9..57b3a963 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/listener/KafkaMessageListenerContainer.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/listener/KafkaMessageListenerContainer.java @@ -772,9 +772,6 @@ public class KafkaMessageListenerContainer extends AbstractMessageListener if (!this.isAnyManualAck && !this.autoCommit) { this.acks.add(record); } - if (this.isRecordAck) { - this.consumer.wakeup(); - } } catch (Exception e) { if (this.containerProperties.isAckOnError() && !this.autoCommit) { @@ -792,9 +789,6 @@ public class KafkaMessageListenerContainer extends AbstractMessageListener } } } - if (this.isManualAck || this.isBatchAck) { - this.consumer.wakeup(); - } } private void processCommits() { @@ -995,9 +989,6 @@ public class KafkaMessageListenerContainer extends AbstractMessageListener ListenerConsumer.this.logger.debug("Interrupt ignored"); } } - if (!ListenerConsumer.this.isManualImmediateAck && this.active) { - ListenerConsumer.this.consumer.wakeup(); - } } } finally { diff --git a/spring-kafka/src/test/java/org/springframework/kafka/listener/KafkaMessageListenerContainerTests.java b/spring-kafka/src/test/java/org/springframework/kafka/listener/KafkaMessageListenerContainerTests.java index e7f757b2..3ce46ab7 100644 --- a/spring-kafka/src/test/java/org/springframework/kafka/listener/KafkaMessageListenerContainerTests.java +++ b/spring-kafka/src/test/java/org/springframework/kafka/listener/KafkaMessageListenerContainerTests.java @@ -571,7 +571,7 @@ public class KafkaMessageListenerContainerTests { }); containerProps.setSyncCommits(true); containerProps.setAckMode(AckMode.BATCH); - containerProps.setPollTimeout(10000); + containerProps.setPollTimeout(100); containerProps.setAckOnError(false); KafkaMessageListenerContainer container = new KafkaMessageListenerContainer<>(cf, @@ -640,7 +640,7 @@ public class KafkaMessageListenerContainerTests { }); containerProps.setSyncCommits(true); containerProps.setAckMode(AckMode.BATCH); - containerProps.setPollTimeout(10000); + containerProps.setPollTimeout(100); containerProps.setAckOnError(false); KafkaMessageListenerContainer container = new KafkaMessageListenerContainer<>(cf, @@ -714,7 +714,7 @@ public class KafkaMessageListenerContainerTests { }); containerProps.setSyncCommits(true); containerProps.setAckMode(AckMode.MANUAL_IMMEDIATE); - containerProps.setPollTimeout(10000); + containerProps.setPollTimeout(100); containerProps.setAckOnError(false); KafkaMessageListenerContainer container = new KafkaMessageListenerContainer<>(cf, @@ -777,7 +777,7 @@ public class KafkaMessageListenerContainerTests { }); containerProps.setSyncCommits(true); containerProps.setAckMode(AckMode.BATCH); - containerProps.setPollTimeout(10000); + containerProps.setPollTimeout(100); containerProps.setAckOnError(true); final CountDownLatch latch = new CountDownLatch(4); containerProps.setGenericErrorHandler((BatchErrorHandler) (t, messages) -> {