From bdb46760ad7b0e8fecede2ac54166b5a0e844545 Mon Sep 17 00:00:00 2001 From: Martin Dam Date: Fri, 16 Oct 2015 10:40:11 -0700 Subject: [PATCH] Support for sending messages synchronously to the Kafka broker. Adds two fields to the producer-configuration that allows to set the send mode (async or sync) and an optional timeout in sync mode. Default behavior is async. This allows for an upstream application to verify that the message have been handed off to Kafka following the delivery options specified, and not just residing in memory at the producer. --- .../xml/KafkaProducerContextParser.java | 5 ++++ .../kafka/support/ProducerConfiguration.java | 26 +++++++++++++++++-- .../kafka/support/ProducerMetadata.java | 23 +++++++++++++++- .../main/resources/META-INF/spring.schemas | 4 +-- ...2.xsd => spring-integration-kafka-1.3.xsd} | 14 ++++++++++ 5 files changed, 67 insertions(+), 5 deletions(-) rename spring-integration-kafka/src/main/resources/org/springframework/integration/config/xml/{spring-integration-kafka-1.2.xsd => spring-integration-kafka-1.3.xsd} (98%) 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) + + +