GH-1293: Fix KafkaTestUtils.getSingleRecord()

Resolves: https://github.com/spring-projects/spring-kafka/issues/1293

`getSingleRecord()` would fail prematurely if records were present in other topics

- Add a loop to prevent early exit
- Always reset the offsets for other topics before assertions
This commit is contained in:
Gary Russell
2019-11-07 12:42:20 -05:00
committed by Artem Bilan
parent 2376699772
commit 2643f5f22e
2 changed files with 41 additions and 7 deletions

View File

@@ -150,12 +150,13 @@ public final class KafkaTestUtils {
* @since 2.0
*/
public static <K, V> ConsumerRecord<K, V> getSingleRecord(Consumer<K, V> consumer, String topic, long timeout) {
ConsumerRecords<K, V> received = getRecords(consumer, timeout);
Iterator<ConsumerRecord<K, V>> 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<K, V> received;
Iterator<ConsumerRecord<K, V>> iterator;
long remaining = timeout;
do {
received = getRecords(consumer, remaining);
iterator = received.records(topic).iterator();
Map<TopicPartition, Long> 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();
}

View File

@@ -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<String, Object> producerProps = KafkaTestUtils.producerProps(broker);
KafkaProducer<Integer, String> producer = new KafkaProducer<>(producerProps);
producer.send(new ProducerRecord<>("singleTopic4", 1, "foo"));
Map<String, Object> consumerProps = KafkaTestUtils.consumerProps("ktuTests", "false", broker);
consumerProps.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest");
KafkaConsumer<Integer, String> 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<String, Object> producerProps = KafkaTestUtils.producerProps(broker);