diff --git a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/config/xml/KafkaOutboundChannelAdapterParser.java b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/config/xml/KafkaOutboundChannelAdapterParser.java index bff194ea10..c166e012da 100644 --- a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/config/xml/KafkaOutboundChannelAdapterParser.java +++ b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/config/xml/KafkaOutboundChannelAdapterParser.java @@ -67,6 +67,14 @@ public class KafkaOutboundChannelAdapterParser extends AbstractOutboundChannelAd kafkaProducerMessageHandlerBuilder.addPropertyValue("partitionIdExpression", partitionIdExpressionDef); } IntegrationNamespaceUtils.setValueIfAttributeDefined(kafkaProducerMessageHandlerBuilder, element, "sync"); + + BeanDefinition sendTimeoutExpressionDef = + IntegrationNamespaceUtils.createExpressionDefinitionFromValueOrExpression("send-timeout", + "send-timeout-expression", parserContext, element, false); + if (sendTimeoutExpressionDef != null) { + kafkaProducerMessageHandlerBuilder.addPropertyValue("sendTimeoutExpression", sendTimeoutExpressionDef); + } + return kafkaProducerMessageHandlerBuilder.getBeanDefinition(); } diff --git a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/outbound/KafkaProducerMessageHandler.java b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/outbound/KafkaProducerMessageHandler.java index 0b6ba94e92..58e0f15fff 100644 --- a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/outbound/KafkaProducerMessageHandler.java +++ b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/outbound/KafkaProducerMessageHandler.java @@ -23,6 +23,7 @@ import org.springframework.expression.EvaluationContext; import org.springframework.expression.Expression; import org.springframework.integration.MessageTimeoutException; import org.springframework.integration.expression.ExpressionUtils; +import org.springframework.integration.expression.ValueExpression; import org.springframework.integration.handler.AbstractMessageHandler; import org.springframework.kafka.core.KafkaTemplate; import org.springframework.kafka.support.KafkaHeaders; @@ -60,7 +61,7 @@ public class KafkaProducerMessageHandler extends AbstractMessageHandler { private boolean sync; - private long sendTimeout = DEFAULT_SEND_TIMEOUT; + private Expression sendTimeoutExpression = new ValueExpression<>(DEFAULT_SEND_TIMEOUT); public KafkaProducerMessageHandler(final KafkaTemplate kafkaTemplate) { Assert.notNull(kafkaTemplate, "kafkaTemplate cannot be null"); @@ -95,12 +96,27 @@ public class KafkaProducerMessageHandler extends AbstractMessageHandler { } /** - * Specify a timeout in milliseconds how long {@link KafkaProducerMessageHandler} - * should wait wait for send operation results. Defaults to 10 seconds. + * Specify a timeout in milliseconds for how long this + * {@link KafkaProducerMessageHandler} should wait wait for send operation + * results. Defaults to 10 seconds. The timeout is applied only in {@link #sync} mode. * @param sendTimeout the timeout to wait for result fo send operation. + * @since 2.0.1 */ public void setSendTimeout(long sendTimeout) { - this.sendTimeout = sendTimeout; + setSendTimeoutExpression(new ValueExpression<>(sendTimeout)); + } + + /** + * Specify a SpEL expression to evaluate a timeout in milliseconds for how long this + * {@link KafkaProducerMessageHandler} should wait wait for send operation + * results. Defaults to 10 seconds. The timeout is applied only in {@link #sync} mode. + * @param sendTimeoutExpression the {@link Expression} for timeout to wait for result + * fo send operation. + * @since 2.1.1 + */ + public void setSendTimeoutExpression(Expression sendTimeoutExpression) { + Assert.notNull(sendTimeoutExpression, "'sendTimeoutExpression' must not be null"); + this.sendTimeoutExpression = sendTimeoutExpression; } @Override @@ -149,12 +165,13 @@ public class KafkaProducerMessageHandler extends AbstractMessageHandler { } } if (this.sync) { - if (this.sendTimeout < 0) { + Long sendTimeout = this.sendTimeoutExpression.getValue(this.evaluationContext, message, Long.class); + if (sendTimeout == null || sendTimeout < 0) { future.get(); } else { try { - future.get(this.sendTimeout, TimeUnit.MILLISECONDS); + future.get(sendTimeout, TimeUnit.MILLISECONDS); } catch (TimeoutException te) { throw new MessageTimeoutException(message, "Timeout waiting for response from KafkaProducer", te); diff --git a/spring-integration-kafka/src/main/resources/org/springframework/integration/kafka/config/spring-integration-kafka-2.1.xsd b/spring-integration-kafka/src/main/resources/org/springframework/integration/kafka/config/spring-integration-kafka-2.1.xsd index ff2f1ebfd0..a04f42f661 100644 --- a/spring-integration-kafka/src/main/resources/org/springframework/integration/kafka/config/spring-integration-kafka-2.1.xsd +++ b/spring-integration-kafka/src/main/resources/org/springframework/integration/kafka/config/spring-integration-kafka-2.1.xsd @@ -99,15 +99,37 @@ - + + + + + Specifies a timeout in milliseconds for how long the 'KafkaProducerMessageHandler' + should wait wait for send operation results. Defaults to 10 seconds. + The timeout is applied only in 'sync' mode. + + + + + + + + + + Specifies an expression that is evaluated to determine a timeout in milliseconds + for how long the 'KafkaProducerMessageHandler' + should wait wait for send operation results. Defaults to 10 seconds. + The timeout is applied only in 'sync' mode. + + + diff --git a/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/config/xml/KafkaOutboundAdapterParserTests-context.xml b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/config/xml/KafkaOutboundAdapterParserTests-context.xml index 755b6c732b..16a6606b76 100644 --- a/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/config/xml/KafkaOutboundAdapterParserTests-context.xml +++ b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/config/xml/KafkaOutboundAdapterParserTests-context.xml @@ -3,10 +3,9 @@ xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance" xmlns:int="http://www.springframework.org/schema/integration" xmlns:int-kafka="http://www.springframework.org/schema/integration/kafka" - xmlns:context="http://www.springframework.org/schema/context" xsi:schemaLocation="http://www.springframework.org/schema/integration/kafka http://www.springframework.org/schema/integration/kafka/spring-integration-kafka.xsd http://www.springframework.org/schema/integration http://www.springframework.org/schema/integration/spring-integration.xsd - http://www.springframework.org/schema/beans http://www.springframework.org/schema/beans/spring-beans.xsd http://www.springframework.org/schema/context http://www.springframework.org/schema/context/spring-context.xsd"> + http://www.springframework.org/schema/beans http://www.springframework.org/schema/beans/spring-beans.xsd"> @@ -18,7 +17,8 @@ sync="true" topic="foo" message-key-expression="'bar'" - partition-id-expression="'2'"> + partition-id-expression="'2'" + send-timeout-expression="1000"> @@ -40,6 +40,7 @@ kafka-template="template" auto-startup="false" channel="inputToKafka" - partition-id="#{0}" /> + partition-id="#{0}" + send-timeout="500" /> diff --git a/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/config/xml/KafkaOutboundAdapterParserTests.java b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/config/xml/KafkaOutboundAdapterParserTests.java index 2a21e2ca31..3c16c141f3 100644 --- a/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/config/xml/KafkaOutboundAdapterParserTests.java +++ b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/config/xml/KafkaOutboundAdapterParserTests.java @@ -68,12 +68,15 @@ public class KafkaOutboundAdapterParserTests { assertThat(TestUtils.getPropertyValue(messageHandler, "messageKeyExpression.expression")).isEqualTo("'bar'"); assertThat(TestUtils.getPropertyValue(messageHandler, "partitionIdExpression.expression")).isEqualTo("'2'"); assertThat(TestUtils.getPropertyValue(messageHandler, "sync", Boolean.class)).isTrue(); + assertThat(TestUtils.getPropertyValue(messageHandler, "sendTimeoutExpression.expression")).isEqualTo("1000"); messageHandler = this.appContext.getBean("kafkaOutboundChannelAdapter2.handler", KafkaProducerMessageHandler.class); assertThat(messageHandler).isNotNull(); assertThat(TestUtils.getPropertyValue(messageHandler, "partitionIdExpression.literalValue")).isEqualTo("0"); assertThat(TestUtils.getPropertyValue(messageHandler, "sync", Boolean.class)).isFalse(); + + assertThat(TestUtils.getPropertyValue(messageHandler, "sendTimeoutExpression.literalValue")).isEqualTo("500"); }