diff --git a/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderTests.java b/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderTests.java index d7d3d5b56..d8d9dba93 100644 --- a/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderTests.java +++ b/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderTests.java @@ -544,6 +544,8 @@ public class KafkaBinderTests extends Message receivedMessage = receive(dlqChannel, 5); assertThat(receivedMessage).isNotNull(); assertThat(receivedMessage.getPayload()).isEqualTo("foo".getBytes()); + //Adding a 1 second sleep to give the retrying enough time to complete. + Thread.sleep(1000); assertThat(handler.getInvocationCount()) .isEqualTo(consumerProperties.getMaxAttempts()); assertThat(receivedMessage.getHeaders()