diff --git a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/config/xml/KafkaProducerContextParser.java b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/config/xml/KafkaProducerContextParser.java index 6edf8258ce..a4435e8aa9 100644 --- a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/config/xml/KafkaProducerContextParser.java +++ b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/config/xml/KafkaProducerContextParser.java @@ -127,6 +127,11 @@ public class KafkaProducerContextParser extends AbstractSimpleBeanDefinitionPars "compression-type"); IntegrationNamespaceUtils.setValueIfAttributeDefined(producerMetadataBuilder, producerConfiguration, "batch-bytes"); + IntegrationNamespaceUtils.setValueIfAttributeDefined(producerMetadataBuilder, producerConfiguration, + "sync"); + IntegrationNamespaceUtils.setValueIfAttributeDefined(producerMetadataBuilder, producerConfiguration, + "send-timeout"); + AbstractBeanDefinition producerMetadataBeanDefinition = producerMetadataBuilder.getBeanDefinition(); String producerPropertiesBean = parentElem.getAttribute("producer-properties"); diff --git a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/support/ProducerConfiguration.java b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/support/ProducerConfiguration.java index a6a9e53e78..5d47afea58 100644 --- a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/support/ProducerConfiguration.java +++ b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/support/ProducerConfiguration.java @@ -16,12 +16,15 @@ package org.springframework.integration.kafka.support; +import java.util.concurrent.ExecutionException; import java.util.concurrent.Future; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.TimeoutException; import org.apache.kafka.clients.producer.Producer; import org.apache.kafka.clients.producer.ProducerRecord; import org.apache.kafka.clients.producer.RecordMetadata; - +import org.apache.kafka.common.KafkaException; import org.springframework.core.convert.ConversionService; import org.springframework.core.convert.support.GenericConversionService; import org.springframework.core.serializer.support.SerializingConverter; @@ -35,6 +38,7 @@ import org.springframework.util.StringUtils; * @author Gary Russell * @author Marius Bogoevici * @author Artem Bilan + * @author Martin Dam * @since 0.5 */ public class ProducerConfiguration { @@ -76,7 +80,25 @@ public class ProducerConfiguration { public Future send(String topic, Integer partition, K messageKey, V messagePayload) { String targetTopic = StringUtils.hasText(topic) ? topic : this.producerMetadata.getTopic(); - return this.producer.send(new ProducerRecord<>(targetTopic, partition, messageKey, messagePayload)); + Future future = this.producer.send(new ProducerRecord<>(targetTopic, partition, messageKey, messagePayload)); + + if (!producerMetadata.isSync()) { + return future; + } + else { + try { + if (producerMetadata.getSendTimeout() <= 0) { + future.get(); + } + else { + future.get(producerMetadata.getSendTimeout(), TimeUnit.MILLISECONDS); + } + } + catch (InterruptedException | ExecutionException | TimeoutException e) { + throw new KafkaException(e); + } + return future; + } } public Future convertAndSend(String topic, Integer partition, Object messageKey, diff --git a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/support/ProducerMetadata.java b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/support/ProducerMetadata.java index 1d3e494f81..9dbcc131ff 100644 --- a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/support/ProducerMetadata.java +++ b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/support/ProducerMetadata.java @@ -49,6 +49,10 @@ public class ProducerMetadata { private int batchBytes = 16384; + private int sendTimeout = 0; + + private boolean sync = false; + public ProducerMetadata(final String topic, Class keyClassType, Class valueClassType, Serializer keySerializer, Serializer valueSerializer) { Assert.notNull(topic, "Topic cannot be null"); @@ -109,6 +113,22 @@ public class ProducerMetadata { return valueClassType; } + public int getSendTimeout() { + return sendTimeout; + } + + public void setSendTimeout(int sendTimeout) { + this.sendTimeout = sendTimeout; + } + + public boolean isSync() { + return sync; + } + + public void setSync(boolean sync) { + this.sync = sync; + } + @Override public String toString() { return "ProducerMetadata{" + @@ -120,6 +140,8 @@ public class ProducerMetadata { ", partitioner=" + this.partitioner + ", compressionType=" + this.compressionType + ", batchBytes=" + this.batchBytes + + ", sync=" + this.sync + + ", sendTimeout=" + this.sendTimeout + '}'; } @@ -128,5 +150,4 @@ public class ProducerMetadata { gzip, snappy } - } diff --git a/spring-integration-kafka/src/main/resources/META-INF/spring.schemas b/spring-integration-kafka/src/main/resources/META-INF/spring.schemas index 73df17583d..78a4988e6a 100644 --- a/spring-integration-kafka/src/main/resources/META-INF/spring.schemas +++ b/spring-integration-kafka/src/main/resources/META-INF/spring.schemas @@ -1,2 +1,2 @@ -http\://www.springframework.org/schema/integration/kafka/spring-integration-kafka-1.2.xsd=org/springframework/integration/config/xml/spring-integration-kafka-1.2.xsd -http\://www.springframework.org/schema/integration/kafka/spring-integration-kafka.xsd=org/springframework/integration/config/xml/spring-integration-kafka-1.2.xsd +http\://www.springframework.org/schema/integration/kafka/spring-integration-kafka-1.3.xsd=org/springframework/integration/config/xml/spring-integration-kafka-1.3.xsd +http\://www.springframework.org/schema/integration/kafka/spring-integration-kafka.xsd=org/springframework/integration/config/xml/spring-integration-kafka-1.3.xsd diff --git a/spring-integration-kafka/src/main/resources/org/springframework/integration/config/xml/spring-integration-kafka-1.2.xsd b/spring-integration-kafka/src/main/resources/org/springframework/integration/config/xml/spring-integration-kafka-1.3.xsd similarity index 98% rename from spring-integration-kafka/src/main/resources/org/springframework/integration/config/xml/spring-integration-kafka-1.2.xsd rename to spring-integration-kafka/src/main/resources/org/springframework/integration/config/xml/spring-integration-kafka-1.3.xsd index fd8fba5f69..c63a78b77c 100644 --- a/spring-integration-kafka/src/main/resources/org/springframework/integration/config/xml/spring-integration-kafka-1.2.xsd +++ b/spring-integration-kafka/src/main/resources/org/springframework/integration/config/xml/spring-integration-kafka-1.3.xsd @@ -215,6 +215,20 @@ + + + + If true, sends messages synchronously. Asynchronously (default) otherwise + + + + + + + If sync=true, specify timeout in milliseconds. 0 or below indicate forever (default) + + +