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 a96f5c6c6..52f67785c 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
@@ -25,6 +25,7 @@ import org.springframework.util.concurrent.ListenableFuture;
import java.util.ArrayList;
import java.util.List;
+import java.util.concurrent.TimeUnit;
/**
*
@@ -33,6 +34,7 @@ import java.util.List;
*
*
* @author Mathieu Ouellet
+ * @author Mahmoud Ben Hassine
* @since 4.2
*
*/
@@ -40,6 +42,7 @@ public class KafkaItemWriter extends KeyValueItemWriter {
protected KafkaTemplate kafkaTemplate;
private final List>> listenableFutures = new ArrayList<>();
+ private long timeout = -1;
@Override
protected void writeKeyValue(K key, T value) {
@@ -55,7 +58,12 @@ public class KafkaItemWriter extends KeyValueItemWriter {
protected void flush() throws Exception{
this.kafkaTemplate.flush();
for(ListenableFuture> future: this.listenableFutures){
- future.get();
+ if (this.timeout >= 0) {
+ future.get(this.timeout, TimeUnit.MILLISECONDS);
+ }
+ else {
+ future.get();
+ }
}
this.listenableFutures.clear();
}
@@ -73,4 +81,15 @@ public class KafkaItemWriter extends KeyValueItemWriter {
public void setKafkaTemplate(KafkaTemplate kafkaTemplate) {
this.kafkaTemplate = kafkaTemplate;
}
+
+ /**
+ * The time limit to wait when flushing items to Kafka.
+ *
+ * @param timeout milliseconds to wait, defaults to -1 (no timeout).
+ * @since 4.3.2
+ */
+ public void setTimeout(long timeout) {
+ this.timeout = timeout;
+ }
+
}
diff --git a/spring-batch-infrastructure/src/main/java/org/springframework/batch/item/kafka/builder/KafkaItemWriterBuilder.java b/spring-batch-infrastructure/src/main/java/org/springframework/batch/item/kafka/builder/KafkaItemWriterBuilder.java
index 30c83e315..09df94027 100644
--- a/spring-batch-infrastructure/src/main/java/org/springframework/batch/item/kafka/builder/KafkaItemWriterBuilder.java
+++ b/spring-batch-infrastructure/src/main/java/org/springframework/batch/item/kafka/builder/KafkaItemWriterBuilder.java
@@ -1,5 +1,5 @@
/*
- * Copyright 2019 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.
@@ -25,6 +25,7 @@ import org.springframework.util.Assert;
* A builder implementation for the {@link KafkaItemWriter}
*
* @author Mathieu Ouellet
+ * @author Mahmoud Ben Hassine
* @since 4.2
*/
public class KafkaItemWriterBuilder {
@@ -35,6 +36,8 @@ public class KafkaItemWriterBuilder {
private boolean delete;
+ private long timeout = -1;
+
/**
* Establish the KafkaTemplate to be used by the KafkaItemWriter.
* @param kafkaTemplate the template to be used
@@ -71,6 +74,19 @@ public class KafkaItemWriterBuilder {
return this;
}
+ /**
+ * The time limit to wait when flushing items to Kafka.
+ *
+ * @param timeout milliseconds to wait, defaults to -1 (no timeout).
+ * @return The current instance of the builder.
+ * @see KafkaItemWriter#setTimeout(long)
+ * @since 4.3.2
+ */
+ public KafkaItemWriterBuilder timeout(long timeout) {
+ this.timeout = timeout;
+ return this;
+ }
+
/**
* Validates and builds a {@link KafkaItemWriter}.
* @return a {@link KafkaItemWriter}
@@ -83,6 +99,7 @@ public class KafkaItemWriterBuilder {
writer.setKafkaTemplate(this.kafkaTemplate);
writer.setItemKeyMapper(this.itemKeyMapper);
writer.setDelete(this.delete);
+ writer.setTimeout(this.timeout);
return writer;
}
}
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 0eac5ac2d..b2fa76648 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
@@ -17,6 +17,7 @@ package org.springframework.batch.item.kafka;
import java.util.Arrays;
import java.util.List;
+import java.util.concurrent.TimeUnit;
import org.junit.Before;
import org.junit.Test;
@@ -57,6 +58,7 @@ public class KafkaItemWriterTests {
this.writer.setKafkaTemplate(this.kafkaTemplate);
this.writer.setItemKeyMapper(this.itemKeyMapper);
this.writer.setDelete(false);
+ this.writer.setTimeout(10L);
this.writer.afterPropertiesSet();
}
@@ -99,7 +101,7 @@ public class KafkaItemWriterTests {
verify(this.kafkaTemplate).sendDefault(items.get(0), items.get(0));
verify(this.kafkaTemplate).sendDefault(items.get(1), items.get(1));
verify(this.kafkaTemplate).flush();
- verify(this.future, times(2)).get();
+ verify(this.future, times(2)).get(10L, TimeUnit.MILLISECONDS);
}
@Test
@@ -112,7 +114,7 @@ public class KafkaItemWriterTests {
verify(this.kafkaTemplate).sendDefault(items.get(0), null);
verify(this.kafkaTemplate).sendDefault(items.get(1), null);
verify(this.kafkaTemplate).flush();
- verify(this.future, times(2)).get();
+ verify(this.future, times(2)).get(10L, TimeUnit.MILLISECONDS);
}
@Test
diff --git a/spring-batch-infrastructure/src/test/java/org/springframework/batch/item/kafka/builder/KafkaItemWriterBuilderTests.java b/spring-batch-infrastructure/src/test/java/org/springframework/batch/item/kafka/builder/KafkaItemWriterBuilderTests.java
index 1ebe70ee9..6aca62998 100644
--- a/spring-batch-infrastructure/src/test/java/org/springframework/batch/item/kafka/builder/KafkaItemWriterBuilderTests.java
+++ b/spring-batch-infrastructure/src/test/java/org/springframework/batch/item/kafka/builder/KafkaItemWriterBuilderTests.java
@@ -1,5 +1,5 @@
/*
- * Copyright 2019 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.
@@ -33,6 +33,7 @@ import static org.junit.Assert.assertTrue;
/**
* @author Mathieu Ouellet
+ * @author Mahmoud Ben Hassine
*/
public class KafkaItemWriterBuilderTests {
@@ -70,13 +71,19 @@ public class KafkaItemWriterBuilderTests {
public void testKafkaItemWriterBuild() {
// given
boolean delete = true;
+ long timeout = 10L;
// when
KafkaItemWriter writer = new KafkaItemWriterBuilder()
- .kafkaTemplate(this.kafkaTemplate).itemKeyMapper(this.itemKeyMapper).delete(delete).build();
+ .kafkaTemplate(this.kafkaTemplate)
+ .itemKeyMapper(this.itemKeyMapper)
+ .delete(delete)
+ .timeout(timeout)
+ .build();
// then
assertTrue((Boolean) ReflectionTestUtils.getField(writer, "delete"));
+ assertEquals(timeout, ReflectionTestUtils.getField(writer, "timeout"));
assertEquals(this.itemKeyMapper, ReflectionTestUtils.getField(writer, "itemKeyMapper"));
assertEquals(this.kafkaTemplate, ReflectionTestUtils.getField(writer, "kafkaTemplate"));
}