From adfd6ee49e794642b9cf6efccc7a6e189acccb4d Mon Sep 17 00:00:00 2001 From: Gary Russell Date: Fri, 16 Sep 2016 09:28:53 -0400 Subject: [PATCH] GH-143: Add sync to XML Namespace Resolves: https://github.com/spring-projects/spring-integration-kafka/issues/143 --- .../xml/KafkaOutboundChannelAdapterParser.java | 2 +- .../kafka/outbound/KafkaProducerMessageHandler.java | 4 ++-- .../kafka/config/spring-integration-kafka-2.1.xsd | 12 ++++++++++++ .../xml/KafkaOutboundAdapterParserTests-context.xml | 1 + .../config/xml/KafkaOutboundAdapterParserTests.java | 3 +++ 5 files changed, 19 insertions(+), 3 deletions(-) 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 29ef3e8740..bff194ea10 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 @@ -66,7 +66,7 @@ public class KafkaOutboundChannelAdapterParser extends AbstractOutboundChannelAd if (partitionIdExpressionDef != null) { kafkaProducerMessageHandlerBuilder.addPropertyValue("partitionIdExpression", partitionIdExpressionDef); } - + IntegrationNamespaceUtils.setValueIfAttributeDefined(kafkaProducerMessageHandlerBuilder, element, "sync"); 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 ba0d09a4c1..0b6ba94e92 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 @@ -84,8 +84,8 @@ public class KafkaProducerMessageHandler extends AbstractMessageHandler { } /** - * The {@code boolean} indicated if {@link KafkaProducerMessageHandler} - * should wait for send operation results or not. Defaults to {@code false}. + * A {@code boolean} indicating if the {@link KafkaProducerMessageHandler} + * should wait for the send operation results or not. Defaults to {@code false}. * In {@code sync} mode a downstream send operation exception will be re-thrown. * @param sync the send mode; async by default. * @since 2.0.1 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 4fafa2f868..51d4b64192 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 @@ -96,6 +96,18 @@ ]]> + + + + + + + + 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 ed7187411f..755b6c732b 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 @@ -15,6 +15,7 @@ auto-startup="false" channel="inputToKafka" order="3" + sync="true" topic="foo" message-key-expression="'bar'" partition-id-expression="'2'"> 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 ff59bf3328..2a21e2ca31 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 @@ -67,16 +67,19 @@ public class KafkaOutboundAdapterParserTests { assertThat(TestUtils.getPropertyValue(messageHandler, "topicExpression.literalValue")).isEqualTo("foo"); assertThat(TestUtils.getPropertyValue(messageHandler, "messageKeyExpression.expression")).isEqualTo("'bar'"); assertThat(TestUtils.getPropertyValue(messageHandler, "partitionIdExpression.expression")).isEqualTo("'2'"); + assertThat(TestUtils.getPropertyValue(messageHandler, "sync", Boolean.class)).isTrue(); 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(); } @Test public void testSyncMode() { + @SuppressWarnings("resource") MockProducer mockProducer = new MockProducer<>(false, new IntegerSerializer(), new StringSerializer()); KafkaTemplate template = new KafkaTemplate<>(() -> mockProducer);