From 0e74cd1cee1a70accf596bd311b30fbf6033c4f5 Mon Sep 17 00:00:00 2001 From: tomvandenberge Date: Mon, 30 Sep 2019 16:23:37 +0200 Subject: [PATCH] Add XML attribute for header-mapper * Added support for header-mapper to outbound-channel-adapter and outbound-gateway XML * Follow up on review * Corrected mistake * Corrected indentation --- .../kafka/config/xml/KafkaParsingUtils.java | 3 +++ .../outbound/KafkaProducerMessageHandler.java | 5 ++++ .../config/spring-integration-kafka-3.2.xsd | 17 ++++++++++++++ ...afkaOutboundAdapterParserTests-context.xml | 3 +++ .../xml/KafkaOutboundAdapterParserTests.java | 4 +++- ...afkaOutboundGatewayParserTests-context.xml | 5 +++- .../xml/KafkaOutboundGatewayParserTests.java | 3 +++ .../KafkaProducerMessageHandlerTests.java | 23 +++++++++++++++++++ 8 files changed, 61 insertions(+), 2 deletions(-) diff --git a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/config/xml/KafkaParsingUtils.java b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/config/xml/KafkaParsingUtils.java index 3f2b6d4246..8c022f00f1 100644 --- a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/config/xml/KafkaParsingUtils.java +++ b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/config/xml/KafkaParsingUtils.java @@ -27,6 +27,7 @@ import org.springframework.integration.config.xml.IntegrationNamespaceUtils; * Utilities to assist with parsing XML. * * @author Gary Russell + * @author Tom van den Berge * @since 3.2 * */ @@ -81,6 +82,8 @@ public final class KafkaParsingUtils { if (timestampExpressionDef != null) { builder.addPropertyValue("timestampExpression", timestampExpressionDef); } + + IntegrationNamespaceUtils.setReferenceIfAttributeDefined(builder, element, "header-mapper"); } } 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 6c4e10b242..38a8609857 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,6 +84,7 @@ import org.springframework.util.concurrent.SettableListenableFuture; * @author Gary Russell * @author Marius Bogoevici * @author Biju Kunjummen + * @author Tom van den Berge * * @since 0.5 */ @@ -186,6 +187,10 @@ public class KafkaProducerMessageHandler extends AbstractReplyProducingMes this.headerMapper = headerMapper; } + public KafkaHeaderMapper getHeaderMapper() { + return this.headerMapper; + } + public KafkaTemplate getKafkaTemplate() { return this.kafkaTemplate; } diff --git a/spring-integration-kafka/src/main/resources/org/springframework/integration/kafka/config/spring-integration-kafka-3.2.xsd b/spring-integration-kafka/src/main/resources/org/springframework/integration/kafka/config/spring-integration-kafka-3.2.xsd index 7b6080ac2b..5eaae94780 100644 --- a/spring-integration-kafka/src/main/resources/org/springframework/integration/kafka/config/spring-integration-kafka-3.2.xsd +++ b/spring-integration-kafka/src/main/resources/org/springframework/integration/kafka/config/spring-integration-kafka-3.2.xsd @@ -31,6 +31,7 @@ + @@ -126,6 +127,7 @@ ]]> + @@ -640,4 +642,19 @@ + + + + + + + + + + + + + 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 8497e22c32..51ecacff1f 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 @@ -23,6 +23,7 @@ error-message-strategy="ems" send-failure-channel="failures" send-success-channel="successes" + header-mapper="customHeaderMapper" > @@ -53,5 +54,7 @@ + + 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 12c4b3edf7..a2fbcbbe63 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 @@ -53,6 +53,7 @@ import org.springframework.test.context.junit4.SpringJUnit4ClassRunner; * @author Artem Bilan * @author Gary Russell * @author Biju Kunjummen + * @author Tom van den Berge * * @since 0.5 */ @@ -84,6 +85,8 @@ public class KafkaOutboundAdapterParserTests { .isSameAs(this.appContext.getBean("failures")); assertThat(TestUtils.getPropertyValue(messageHandler, "sendSuccessChannel")) .isSameAs(this.appContext.getBean("successes")); + assertThat(TestUtils.getPropertyValue(messageHandler, "headerMapper")) + .isSameAs(this.appContext.getBean("customHeaderMapper")); messageHandler = this.appContext.getBean("kafkaOutboundChannelAdapter2.handler", KafkaProducerMessageHandler.class); @@ -94,7 +97,6 @@ public class KafkaOutboundAdapterParserTests { assertThat(TestUtils.getPropertyValue(messageHandler, "sendTimeoutExpression.literalValue")).isEqualTo("500"); } - @Test public void testSyncMode() { MockProducer mockProducer = diff --git a/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/config/xml/KafkaOutboundGatewayParserTests-context.xml b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/config/xml/KafkaOutboundGatewayParserTests-context.xml index 74fe603b4e..31123d4839 100644 --- a/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/config/xml/KafkaOutboundGatewayParserTests-context.xml +++ b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/config/xml/KafkaOutboundGatewayParserTests-context.xml @@ -23,7 +23,9 @@ send-timeout-expression="44" sync="true" timestamp-expression="T(System).currentTimeMillis()" - topic-expression="'topic'"/> + topic-expression="'topic'" + header-mapper="customHeaderMapper" + /> @@ -39,4 +41,5 @@ + diff --git a/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/config/xml/KafkaOutboundGatewayParserTests.java b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/config/xml/KafkaOutboundGatewayParserTests.java index acee42b7f9..9804875ee7 100644 --- a/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/config/xml/KafkaOutboundGatewayParserTests.java +++ b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/config/xml/KafkaOutboundGatewayParserTests.java @@ -31,6 +31,7 @@ import org.springframework.test.context.junit.jupiter.SpringJUnitConfig; /** * @author Gary Russell + * @author Tom van den Berge * @since 3.2 * */ @@ -66,6 +67,8 @@ public class KafkaOutboundGatewayParserTests { .isSameAs(this.context.getBean("failures")); assertThat(TestUtils.getPropertyValue(this.messageHandler, "sendSuccessChannel")) .isSameAs(this.context.getBean("successes")); + assertThat(TestUtils.getPropertyValue(this.messageHandler, "headerMapper")) + .isSameAs(this.context.getBean("customHeaderMapper")); } public static class EMS extends DefaultErrorMessageStrategy { diff --git a/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/outbound/KafkaProducerMessageHandlerTests.java b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/outbound/KafkaProducerMessageHandlerTests.java index 53c25dd31b..c576919c5f 100644 --- a/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/outbound/KafkaProducerMessageHandlerTests.java +++ b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/outbound/KafkaProducerMessageHandlerTests.java @@ -104,6 +104,7 @@ import org.springframework.util.concurrent.SettableListenableFuture; * @author Gary Russell * @author Biju Kunjummen * @author Artem Bilan + * @author Tom van den Berge * * @since 2.0 */ @@ -329,6 +330,28 @@ public class KafkaProducerMessageHandlerTests { producerFactory.destroy(); } + @Test + public void testOutboundWithCustomHeaderMapper() throws Exception { + DefaultKafkaProducerFactory producerFactory = new DefaultKafkaProducerFactory<>( + KafkaTestUtils.producerProps(embeddedKafka)); + KafkaTemplate template = new KafkaTemplate<>(producerFactory); + KafkaProducerMessageHandler handler = new KafkaProducerMessageHandler<>(template); + handler.setBeanFactory(mock(BeanFactory.class)); + handler.setHeaderMapper(new DefaultKafkaHeaderMapper("!*")); + handler.afterPropertiesSet(); + + Message message = MessageBuilder.withPayload("foo") + .setHeader(KafkaHeaders.TOPIC, topic1) + .setHeader("foo-header", "foo-header-value") + .build(); + handler.handleMessage(message); + + ConsumerRecord record = KafkaTestUtils.getSingleRecord(consumer, topic1); + assertThat(record.headers().toArray().length).isEqualTo(0); + + producerFactory.destroy(); + } + @Test public void testOutboundGateway() throws Exception { ConsumerFactory consumerFactory = new DefaultKafkaConsumerFactory<>(