Fix KafkaItemReaderTests

Tests in this class fail intermittently because
they send messages to Kafka in an asynchronous
way and assert on the results immediately without
waiting for the send operation to complete.

This commit updates the tests to wait for
send results before asserting on them
(similar to a140a9f5).

(cherry picked from commit 5826111b85)
This commit is contained in:
Mahmoud Ben Hassine
2020-11-17 16:18:05 +01:00
parent 4f2d283d51
commit f9adf9466c

View File

@@ -1,5 +1,5 @@
/*
* Copyright 2019 the original author or authors.
* Copyright 2019-2020 the original author or authors.
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
@@ -26,6 +26,7 @@ import java.util.concurrent.ExecutionException;
import org.apache.kafka.clients.admin.NewTopic;
import org.apache.kafka.clients.consumer.ConsumerConfig;
import org.apache.kafka.clients.consumer.OffsetAndMetadata;
import org.apache.kafka.common.TopicPartition;
import org.apache.kafka.common.serialization.StringDeserializer;
import org.junit.Before;
@@ -37,6 +38,7 @@ import org.springframework.batch.item.ExecutionContext;
import org.springframework.kafka.core.DefaultKafkaProducerFactory;
import org.springframework.kafka.core.KafkaTemplate;
import org.springframework.kafka.core.ProducerFactory;
import org.springframework.kafka.support.SendResult;
import org.springframework.kafka.test.rule.EmbeddedKafkaRule;
import org.springframework.kafka.test.utils.KafkaTestUtils;
import org.springframework.util.concurrent.ListenableFuture;
@@ -183,12 +185,16 @@ public class KafkaItemReaderTests {
}
@Test
public void testReadFromSinglePartition() {
public void testReadFromSinglePartition() throws ExecutionException, InterruptedException {
this.template.setDefaultTopic("topic1");
this.template.sendDefault("val0");
this.template.sendDefault("val1");
this.template.sendDefault("val2");
this.template.sendDefault("val3");
List<ListenableFuture<SendResult<String, String>>> futures = new ArrayList<>();
futures.add(this.template.sendDefault("val0"));
futures.add(this.template.sendDefault("val1"));
futures.add(this.template.sendDefault("val2"));
futures.add(this.template.sendDefault("val3"));
for (ListenableFuture<SendResult<String, String>> future : futures) {
future.get();
}
this.reader = new KafkaItemReader<>(this.consumerProperties, "topic1", 0);
this.reader.setPollTimeout(Duration.ofSeconds(1));
@@ -213,12 +219,16 @@ public class KafkaItemReaderTests {
}
@Test
public void testReadFromMultiplePartitions() {
public void testReadFromMultiplePartitions() throws ExecutionException, InterruptedException {
this.template.setDefaultTopic("topic2");
this.template.sendDefault("val0");
this.template.sendDefault("val1");
this.template.sendDefault("val2");
this.template.sendDefault("val3");
List<ListenableFuture<SendResult<String, String>>> futures = new ArrayList<>();
futures.add(this.template.sendDefault("val0"));
futures.add(this.template.sendDefault("val1"));
futures.add(this.template.sendDefault("val2"));
futures.add(this.template.sendDefault("val3"));
for (ListenableFuture<SendResult<String, String>> future : futures) {
future.get();
}
this.reader = new KafkaItemReader<>(this.consumerProperties, "topic2", 0, 1);
this.reader.setPollTimeout(Duration.ofSeconds(1));
@@ -237,14 +247,17 @@ public class KafkaItemReaderTests {
}
@Test
public void testReadFromSinglePartitionAfterRestart() {
public void testReadFromSinglePartitionAfterRestart() throws ExecutionException, InterruptedException {
this.template.setDefaultTopic("topic3");
this.template.sendDefault("val0");
this.template.sendDefault("val1");
this.template.sendDefault("val2");
this.template.sendDefault("val3");
this.template.sendDefault("val4");
List<ListenableFuture<SendResult<String, String>>> futures = new ArrayList<>();
futures.add(this.template.sendDefault("val0"));
futures.add(this.template.sendDefault("val1"));
futures.add(this.template.sendDefault("val2"));
futures.add(this.template.sendDefault("val3"));
futures.add(this.template.sendDefault("val4"));
for (ListenableFuture<SendResult<String, String>> future : futures) {
future.get();
}
ExecutionContext executionContext = new ExecutionContext();
Map<TopicPartition, Long> offsets = new HashMap<>();
offsets.put(new TopicPartition("topic3", 0), 1L);