diff --git a/spring-kafka-test/src/main/java/org/springframework/kafka/test/utils/KafkaTestUtils.java b/spring-kafka-test/src/main/java/org/springframework/kafka/test/utils/KafkaTestUtils.java index 577247bb..961a845a 100644 --- a/spring-kafka-test/src/main/java/org/springframework/kafka/test/utils/KafkaTestUtils.java +++ b/spring-kafka-test/src/main/java/org/springframework/kafka/test/utils/KafkaTestUtils.java @@ -150,12 +150,13 @@ public final class KafkaTestUtils { * @since 2.0 */ public static ConsumerRecord getSingleRecord(Consumer consumer, String topic, long timeout) { - ConsumerRecords received = getRecords(consumer, timeout); - Iterator> iterator = received.records(topic).iterator(); - assertThat(iterator.hasNext()).as("No records found for topic").isTrue(); - iterator.next(); - assertThat(iterator.hasNext()).as("More than one record for topic found").isFalse(); - if (received.count() > 1) { + long expire = System.currentTimeMillis() + timeout; + ConsumerRecords received; + Iterator> iterator; + long remaining = timeout; + do { + received = getRecords(consumer, remaining); + iterator = received.records(topic).iterator(); Map reset = new HashMap<>(); received.forEach(rec -> { if (!rec.topic().equals(topic)) { @@ -163,7 +164,18 @@ public final class KafkaTestUtils { } }); reset.forEach((tp, off) -> consumer.seek(tp, off)); + try { + Thread.sleep(50); + } + catch (InterruptedException e) { + Thread.currentThread().interrupt(); + } + remaining = expire - System.currentTimeMillis(); } + while (!iterator.hasNext() && remaining > 0); + assertThat(iterator.hasNext()).as("No records found for topic").isTrue(); + iterator.next(); + assertThat(iterator.hasNext()).as("More than one record for topic found").isFalse(); return received.records(topic).iterator().next(); } diff --git a/spring-kafka-test/src/test/java/org/springframework/kafka/test/utils/KafkaTestUtilsTests.java b/spring-kafka-test/src/test/java/org/springframework/kafka/test/utils/KafkaTestUtilsTests.java index f35dd0f3..4281ff20 100644 --- a/spring-kafka-test/src/test/java/org/springframework/kafka/test/utils/KafkaTestUtilsTests.java +++ b/spring-kafka-test/src/test/java/org/springframework/kafka/test/utils/KafkaTestUtilsTests.java @@ -17,6 +17,7 @@ package org.springframework.kafka.test.utils; import static org.assertj.core.api.Assertions.assertThat; +import static org.assertj.core.api.Assertions.assertThatExceptionOfType; import java.util.Map; @@ -26,6 +27,7 @@ import org.apache.kafka.clients.consumer.KafkaConsumer; import org.apache.kafka.clients.producer.KafkaProducer; import org.apache.kafka.clients.producer.ProducerRecord; import org.junit.jupiter.api.Test; +import org.opentest4j.AssertionFailedError; import org.springframework.kafka.test.EmbeddedKafkaBroker; import org.springframework.kafka.test.context.EmbeddedKafka; @@ -35,7 +37,7 @@ import org.springframework.kafka.test.context.EmbeddedKafka; * @since 2.2.7 * */ -@EmbeddedKafka(topics = { "singleTopic1", "singleTopic2", "singleTopic3" }) +@EmbeddedKafka(topics = { "singleTopic1", "singleTopic2", "singleTopic3", "singleTopic4", "singleTopic5" }) public class KafkaTestUtilsTests { @Test @@ -54,6 +56,26 @@ public class KafkaTestUtilsTests { consumer.close(); } + @Test + void testGetSingleWithMoreThatOneTopicRecordNotThereYet(EmbeddedKafkaBroker broker) { + Map producerProps = KafkaTestUtils.producerProps(broker); + KafkaProducer producer = new KafkaProducer<>(producerProps); + producer.send(new ProducerRecord<>("singleTopic4", 1, "foo")); + Map consumerProps = KafkaTestUtils.consumerProps("ktuTests", "false", broker); + consumerProps.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest"); + KafkaConsumer consumer = new KafkaConsumer<>(consumerProps); + broker.consumeFromEmbeddedTopics(consumer, "singleTopic4", "singleTopic5"); + long t1 = System.currentTimeMillis(); + assertThatExceptionOfType(AssertionFailedError.class).isThrownBy(() -> + KafkaTestUtils.getSingleRecord(consumer, "singleTopic5", 2000L)); + assertThat(System.currentTimeMillis() - t1).isGreaterThanOrEqualTo(2000L); + producer.send(new ProducerRecord<>("singleTopic5", 1, "foo")); + producer.close(); + KafkaTestUtils.getSingleRecord(consumer, "singleTopic4"); + KafkaTestUtils.getSingleRecord(consumer, "singleTopic5"); + consumer.close(); + } + @Test public void testGetOneRecord(EmbeddedKafkaBroker broker) throws Exception { Map producerProps = KafkaTestUtils.producerProps(broker);