From 46d7567a77d4dd39845471876ee050577c326215 Mon Sep 17 00:00:00 2001 From: Soby Chacko Date: Tue, 22 Jul 2014 23:44:36 -0400 Subject: [PATCH] INTEXT-104: Kafka: `order` for o-channel-adapter JIRA: https://jira.spring.io/browse/INTEXT-104 `order` attribute now honored if kafka outbound adapter is connected to a subscribable channel unit tests and samples are updated Polishing: use `` for adapter tags to cover `SmartLifecycle` options. --- .../kafka/outbound/UserTransformer.java | 27 +++++++ ...afkaOutboundAdapterParserTests-context.xml | 29 ++++--- .../xml/spring-integration-kafka-1.0.xsd | 76 +++---------------- ...afkaOutboundAdapterParserTests-context.xml | 1 + .../xml/KafkaOutboundAdapterParserTests.java | 14 ++-- 5 files changed, 65 insertions(+), 82 deletions(-) create mode 100644 samples/kafka/src/main/java/org/springframework/integration/samples/kafka/outbound/UserTransformer.java diff --git a/samples/kafka/src/main/java/org/springframework/integration/samples/kafka/outbound/UserTransformer.java b/samples/kafka/src/main/java/org/springframework/integration/samples/kafka/outbound/UserTransformer.java new file mode 100644 index 0000000..093bbcb --- /dev/null +++ b/samples/kafka/src/main/java/org/springframework/integration/samples/kafka/outbound/UserTransformer.java @@ -0,0 +1,27 @@ +package org.springframework.integration.samples.kafka.outbound; + +import org.apache.commons.logging.Log; +import org.apache.commons.logging.LogFactory; +import org.springframework.integration.samples.kafka.user.User; +import org.springframework.integration.support.MessageBuilder; +import org.springframework.integration.transformer.Transformer; +import org.springframework.messaging.Message; + +/** + * @author Soby Chacko + */ +public class UserTransformer implements Transformer { + + Log logger = LogFactory.getLog(getClass()); + + @Override + public Message transform(Message message) { + if(message.getPayload().getClass().isAssignableFrom(User.class)) { + User user = (User) message.getPayload(); + user.setFirstName(user.getFirstName().toString()+user.getFirstName()); + logger.info("user confirmed " + user.getFirstName()); + return MessageBuilder.withPayload(user).copyHeaders(message.getHeaders()).build(); + } + return message; + } +} diff --git a/samples/kafka/src/main/resources/org/springframework/integration/samples/kafka/outbound/kafkaOutboundAdapterParserTests-context.xml b/samples/kafka/src/main/resources/org/springframework/integration/samples/kafka/outbound/kafkaOutboundAdapterParserTests-context.xml index b1f0d8d..01c4689 100644 --- a/samples/kafka/src/main/resources/org/springframework/integration/samples/kafka/outbound/kafkaOutboundAdapterParserTests-context.xml +++ b/samples/kafka/src/main/resources/org/springframework/integration/samples/kafka/outbound/kafkaOutboundAdapterParserTests-context.xml @@ -1,24 +1,31 @@ + http://www.springframework.org/schema/task http://www.springframework.org/schema/task/spring-task.xsd + http://www.springframework.org/schema/integration/stream + http://www.springframework.org/schema/integration/stream/spring-integration-stream.xsd"> - - - + - + channel="inputToKafka" + order="1"> + + + + + + diff --git a/spring-integration-kafka/src/main/resources/org/springframework/integration/config/xml/spring-integration-kafka-1.0.xsd b/spring-integration-kafka/src/main/resources/org/springframework/integration/config/xml/spring-integration-kafka-1.0.xsd index 866255f..92a4e1b 100644 --- a/spring-integration-kafka/src/main/resources/org/springframework/integration/config/xml/spring-integration-kafka-1.0.xsd +++ b/spring-integration-kafka/src/main/resources/org/springframework/integration/config/xml/spring-integration-kafka-1.0.xsd @@ -9,7 +9,7 @@ + schemaLocation="http://www.springframework.org/schema/integration/spring-integration-4.0.xsd"/> - - - - - - - - - - + - + Kafka Server Bean Name @@ -430,30 +420,6 @@ - - - - - Identifies the underlying Spring bean definition, which is an - instance of either 'EventDrivenConsumer' or 'PollingConsumer', - depending on whether the component's input channel is a - 'SubscribableChannel' or 'PollableChannel'. - - - - - - - Flag to indicate that the component should start automatically - on startup (default true). - - - - - - - - @@ -466,42 +432,20 @@ - - - - Identifies the underlying Spring bean definition, which is an - instance of either 'EventDrivenConsumer' or 'PollingConsumer', - depending on whether the component's input channel is a - 'SubscribableChannel' or 'PollableChannel'. - - - - - - - Flag to indicate that the component should start automatically - on startup (default true). - - - - - - - + + Kafka producer context reference. - + - - - - - + + Specifies the order for invocation when this endpoint is connected as a + subscriber to a SubscribableChannel. + 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 89df6ff..d1cdd05 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 @@ -17,6 +17,7 @@ kafka-producer-context-ref="kafkaProducerContext" auto-startup="false" channel="inputToKafka" + order="3" > 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 d2b0909..6e2f564 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 @@ -18,6 +18,7 @@ package org.springframework.integration.kafka.config.xml; import org.junit.Assert; import org.junit.Test; import org.junit.runner.RunWith; + import org.springframework.beans.factory.annotation.Autowired; import org.springframework.context.ApplicationContext; import org.springframework.integration.endpoint.PollingConsumer; @@ -32,20 +33,23 @@ import org.springframework.test.context.junit4.SpringJUnit4ClassRunner; */ @RunWith(SpringJUnit4ClassRunner.class) @ContextConfiguration -public class KafkaOutboundAdapterParserTests { +public class KafkaOutboundAdapterParserTests { @Autowired private ApplicationContext appContext; @Test @SuppressWarnings("unchecked") - public void testOutboundAdapterConfiguration(){ - final PollingConsumer pollingConsumer = appContext.getBean("kafkaOutboundChannelAdapter", PollingConsumer.class); - final KafkaProducerMessageHandler messageHandler = appContext.getBean(KafkaProducerMessageHandler.class); + public void testOutboundAdapterConfiguration() { + final PollingConsumer pollingConsumer = + appContext.getBean("kafkaOutboundChannelAdapter", PollingConsumer.class); + final KafkaProducerMessageHandler messageHandler = appContext.getBean(KafkaProducerMessageHandler.class); Assert.assertNotNull(pollingConsumer); Assert.assertNotNull(messageHandler); - final KafkaProducerContext producerContext = messageHandler.getKafkaProducerContext(); + Assert.assertEquals(messageHandler.getOrder(), 3); + final KafkaProducerContext producerContext = messageHandler.getKafkaProducerContext(); Assert.assertNotNull(producerContext); Assert.assertEquals(producerContext.getTopicsConfiguration().size(), 2); } + }