diff --git a/spring-batch-infrastructure/src/main/java/org/springframework/batch/item/kafka/KafkaItemWriter.java b/spring-batch-infrastructure/src/main/java/org/springframework/batch/item/kafka/KafkaItemWriter.java index 3e496a465..0fcdb63a1 100644 --- a/spring-batch-infrastructure/src/main/java/org/springframework/batch/item/kafka/KafkaItemWriter.java +++ b/spring-batch-infrastructure/src/main/java/org/springframework/batch/item/kafka/KafkaItemWriter.java @@ -1,5 +1,5 @@ /* - * Copyright 2019-2022 the original author or authors. + * Copyright 2019-2021 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. @@ -21,10 +21,10 @@ import org.springframework.batch.item.KeyValueItemWriter; import org.springframework.kafka.core.KafkaTemplate; import org.springframework.kafka.support.SendResult; import org.springframework.util.Assert; +import org.springframework.util.concurrent.ListenableFuture; import java.util.ArrayList; import java.util.List; -import java.util.concurrent.CompletableFuture; import java.util.concurrent.TimeUnit; /** @@ -42,24 +42,24 @@ public class KafkaItemWriter extends KeyValueItemWriter { protected KafkaTemplate kafkaTemplate; - private final List>> completableFutures = new ArrayList<>(); + private final List>> listenableFutures = new ArrayList<>(); private long timeout = -1; @Override protected void writeKeyValue(K key, T value) { if (this.delete) { - this.completableFutures.add(this.kafkaTemplate.sendDefault(key, null)); + this.listenableFutures.add(this.kafkaTemplate.sendDefault(key, null)); } else { - this.completableFutures.add(this.kafkaTemplate.sendDefault(key, value)); + this.listenableFutures.add(this.kafkaTemplate.sendDefault(key, value)); } } @Override protected void flush() throws Exception { this.kafkaTemplate.flush(); - for (var future : this.completableFutures) { + for (ListenableFuture> future : this.listenableFutures) { if (this.timeout >= 0) { future.get(this.timeout, TimeUnit.MILLISECONDS); } @@ -67,7 +67,7 @@ public class KafkaItemWriter extends KeyValueItemWriter { future.get(); } } - this.completableFutures.clear(); + this.listenableFutures.clear(); } @Override diff --git a/spring-batch-infrastructure/src/test/java/org/springframework/batch/item/kafka/KafkaItemReaderTests.java b/spring-batch-infrastructure/src/test/java/org/springframework/batch/item/kafka/KafkaItemReaderTests.java index 532b91ac5..7ee5efd8b 100644 --- a/spring-batch-infrastructure/src/test/java/org/springframework/batch/item/kafka/KafkaItemReaderTests.java +++ b/spring-batch-infrastructure/src/test/java/org/springframework/batch/item/kafka/KafkaItemReaderTests.java @@ -22,7 +22,6 @@ import java.util.HashMap; import java.util.List; import java.util.Map; import java.util.Properties; -import java.util.concurrent.CompletableFuture; import java.util.concurrent.ExecutionException; import org.apache.kafka.clients.admin.NewTopic; @@ -40,10 +39,12 @@ import org.springframework.beans.factory.annotation.Autowired; 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.EmbeddedKafkaBroker; import org.springframework.kafka.test.context.EmbeddedKafka; import org.springframework.kafka.test.utils.KafkaTestUtils; import org.springframework.test.context.junit.jupiter.SpringExtension; +import org.springframework.util.concurrent.ListenableFuture; import static org.hamcrest.MatcherAssert.assertThat; import static org.hamcrest.Matchers.containsInAnyOrder; @@ -142,12 +143,12 @@ class KafkaItemReaderTests { @Test void testReadFromSinglePartition() throws ExecutionException, InterruptedException { this.template.setDefaultTopic("topic1"); - var futures = new ArrayList>(); + List>> 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 (var future : futures) { + for (ListenableFuture> future : futures) { future.get(); } @@ -176,12 +177,12 @@ class KafkaItemReaderTests { @Test void testReadFromSinglePartitionFromCustomOffset() throws ExecutionException, InterruptedException { this.template.setDefaultTopic("topic5"); - var futures = new ArrayList>(); + List>> futures = new ArrayList<>(); futures.add(this.template.sendDefault("val0")); // <-- offset 0 futures.add(this.template.sendDefault("val1")); // <-- offset 1 futures.add(this.template.sendDefault("val2")); // <-- offset 2 futures.add(this.template.sendDefault("val3")); // <-- offset 3 - for (var future : futures) { + for (ListenableFuture> future : futures) { future.get(); } @@ -212,10 +213,10 @@ class KafkaItemReaderTests { // first run: read a topic from the beginning this.template.setDefaultTopic("topic6"); - var futures = new ArrayList>(); + List>> futures = new ArrayList<>(); futures.add(this.template.sendDefault("val0")); // <-- offset 0 futures.add(this.template.sendDefault("val1")); // <-- offset 1 - for (var future : futures) { + for (ListenableFuture> future : futures) { future.get(); } this.reader = new KafkaItemReader<>(this.consumerProperties, "topic6", 0); @@ -266,12 +267,12 @@ class KafkaItemReaderTests { @Test void testReadFromMultiplePartitions() throws ExecutionException, InterruptedException { this.template.setDefaultTopic("topic2"); - var futures = new ArrayList>(); + List>> 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 (var future : futures) { + for (ListenableFuture> future : futures) { future.get(); } @@ -294,13 +295,13 @@ class KafkaItemReaderTests { @Test void testReadFromSinglePartitionAfterRestart() throws ExecutionException, InterruptedException { this.template.setDefaultTopic("topic3"); - var futures = new ArrayList>(); + List>> 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 (var future : futures) { + for (ListenableFuture> future : futures) { future.get(); } ExecutionContext executionContext = new ExecutionContext(); @@ -330,7 +331,7 @@ class KafkaItemReaderTests { @Test void testReadFromMultiplePartitionsAfterRestart() throws ExecutionException, InterruptedException { - var futures = new ArrayList>(); + List>> futures = new ArrayList<>(); futures.add(this.template.send("topic4", 0, null, "val0")); futures.add(this.template.send("topic4", 0, null, "val2")); futures.add(this.template.send("topic4", 0, null, "val4")); @@ -340,7 +341,7 @@ class KafkaItemReaderTests { futures.add(this.template.send("topic4", 1, null, "val5")); futures.add(this.template.send("topic4", 1, null, "val7")); - for (var future : futures) { + for (ListenableFuture future : futures) { future.get(); } diff --git a/spring-batch-infrastructure/src/test/java/org/springframework/batch/item/kafka/KafkaItemWriterTests.java b/spring-batch-infrastructure/src/test/java/org/springframework/batch/item/kafka/KafkaItemWriterTests.java index 069376803..cd9d0476b 100644 --- a/spring-batch-infrastructure/src/test/java/org/springframework/batch/item/kafka/KafkaItemWriterTests.java +++ b/spring-batch-infrastructure/src/test/java/org/springframework/batch/item/kafka/KafkaItemWriterTests.java @@ -1,5 +1,5 @@ /* - * Copyright 2019-2022 the original author or authors. + * Copyright 2019-2021 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. @@ -17,7 +17,6 @@ package org.springframework.batch.item.kafka; import java.util.Arrays; import java.util.List; -import java.util.concurrent.CompletableFuture; import java.util.concurrent.TimeUnit; import org.junit.jupiter.api.BeforeEach; @@ -30,6 +29,7 @@ import org.springframework.batch.item.Chunk; import org.springframework.core.convert.converter.Converter; import org.springframework.kafka.core.KafkaTemplate; import org.springframework.kafka.support.SendResult; +import org.springframework.util.concurrent.ListenableFuture; import static org.junit.jupiter.api.Assertions.assertDoesNotThrow; import static org.junit.jupiter.api.Assertions.assertEquals; @@ -46,7 +46,7 @@ class KafkaItemWriterTests { private KafkaTemplate kafkaTemplate; @Mock - private CompletableFuture> future; + private ListenableFuture> future; private KafkaItemKeyMapper itemKeyMapper;