Revert "Replace ListenableFuture with CompletableFuture"

This reverts commit a2cf74ad as the change was made after
spring-kafka 3.0.0-M5 was released. This version is required
to release spring-batch 5.0.0-M5

This change will be put back after releasing 5.0.0-M5.
This commit is contained in:
Mahmoud Ben Hassine
2022-08-24 17:35:25 +02:00
parent b524e3287f
commit a165247284
3 changed files with 24 additions and 23 deletions

View File

@@ -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<K, T> extends KeyValueItemWriter<K, T> {
protected KafkaTemplate<K, T> kafkaTemplate;
private final List<CompletableFuture<SendResult<K, T>>> completableFutures = new ArrayList<>();
private final List<ListenableFuture<SendResult<K, T>>> 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<SendResult<K, T>> future : this.listenableFutures) {
if (this.timeout >= 0) {
future.get(this.timeout, TimeUnit.MILLISECONDS);
}
@@ -67,7 +67,7 @@ public class KafkaItemWriter<K, T> extends KeyValueItemWriter<K, T> {
future.get();
}
}
this.completableFutures.clear();
this.listenableFutures.clear();
}
@Override

View File

@@ -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<CompletableFuture<?>>();
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 (var future : futures) {
for (ListenableFuture<SendResult<String, String>> future : futures) {
future.get();
}
@@ -176,12 +177,12 @@ class KafkaItemReaderTests {
@Test
void testReadFromSinglePartitionFromCustomOffset() throws ExecutionException, InterruptedException {
this.template.setDefaultTopic("topic5");
var futures = new ArrayList<CompletableFuture<?>>();
List<ListenableFuture<SendResult<String, String>>> 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<SendResult<String, String>> 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<CompletableFuture<?>>();
List<ListenableFuture<SendResult<String, String>>> 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<SendResult<String, String>> 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<CompletableFuture<?>>();
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 (var future : futures) {
for (ListenableFuture<SendResult<String, String>> future : futures) {
future.get();
}
@@ -294,13 +295,13 @@ class KafkaItemReaderTests {
@Test
void testReadFromSinglePartitionAfterRestart() throws ExecutionException, InterruptedException {
this.template.setDefaultTopic("topic3");
var futures = new ArrayList<CompletableFuture<?>>();
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 (var future : futures) {
for (ListenableFuture<SendResult<String, String>> future : futures) {
future.get();
}
ExecutionContext executionContext = new ExecutionContext();
@@ -330,7 +331,7 @@ class KafkaItemReaderTests {
@Test
void testReadFromMultiplePartitionsAfterRestart() throws ExecutionException, InterruptedException {
var futures = new ArrayList<CompletableFuture<?>>();
List<ListenableFuture<SendResult<String, String>>> 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();
}

View File

@@ -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<String, String> kafkaTemplate;
@Mock
private CompletableFuture<SendResult<String, String>> future;
private ListenableFuture<SendResult<String, String>> future;
private KafkaItemKeyMapper itemKeyMapper;