From 614308b4826f7c3683dd10e9d01734cb84b29e68 Mon Sep 17 00:00:00 2001 From: Artem Bilan Date: Mon, 17 Nov 2014 13:22:39 +0200 Subject: [PATCH] IE-100: Kafka:o-c-a: topic and message-key attrs JIRA: https://jira.spring.io/browse/INTEXT-100 Add `topic(-expression)` and `message-key(-expression)` attributes to avoid upstream configuration to specify them in the `MessageHeaders` Polishing for XSD Conflicts: spring-integration-kafka/src/main/java/org/springframework/integration/kafka/support/KafkaProducerContext.java spring-integration-kafka/src/main/java/org/springframework/integration/kafka/support/ProducerConfiguration.java spring-integration-kafka/src/main/resources/org/springframework/integration/config/xml/spring-integration-kafka-1.0.xsd spring-integration-kafka/src/test/java/org/springframework/integration/kafka/config/xml/KafkaOutboundAdapterParserTests.java INTEXT-100: Add `Assert.notNull(this.evaluationContext);` INTEXT-100: Fix parser potential NPE INTEXT-100: Add `KafkaHeaders` Merge branch 'INTEXT-100-1' of ..\spring-integration-extensions into INTEXT-100 Conflicts: src/main/java/org/springframework/integration/kafka/outbound/KafkaProducerMessageHandler.java src/main/java/org/springframework/integration/kafka/support/KafkaProducerContext.java src/main/java/org/springframework/integration/kafka/support/ProducerConfiguration.java src/test/java/org/springframework/integration/kafka/config/xml/KafkaOutboundAdapterParserTests.java * Polishing according the rebase to `master` * Polishing for the `KafkaRunning` Rule to use `ZkClient` and its `getAllBrokersInCluster` * Add a note to the `README.md` Addressing PR comments Polishing - Minor doc polish + change tabs to spaces in code - Enhance test to include expressions --- .../KafkaOutboundChannelAdapterParser.java | 24 +- .../outbound/KafkaProducerMessageHandler.java | 46 ++- .../kafka/support/KafkaHeaders.java | 31 ++ .../kafka/support/KafkaProducerContext.java | 33 +- .../kafka/support/ProducerConfiguration.java | 25 +- .../xml/spring-integration-kafka-1.0.xsd | 320 +++++++++--------- ...afkaOutboundAdapterParserTests-context.xml | 71 ++-- .../xml/KafkaOutboundAdapterParserTests.java | 26 +- .../kafka/outbound/OutboundTests.java | 42 ++- .../integration/kafka/rule/KafkaRunning.java | 50 +-- .../support/ProducerConfigurationTests.java | 175 +++++----- 11 files changed, 476 insertions(+), 367 deletions(-) create mode 100644 spring-integration-kafka/src/main/java/org/springframework/integration/kafka/support/KafkaHeaders.java 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 2df108fc93..e348f8ec61 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 @@ -1,5 +1,5 @@ /* - * Copyright 2002-2013 the original author or authors. + * Copyright 2013-2014 the original author or authors. * * Licensed under the Apache License, Version 2.0 (the "License"); * you may not use this file except in compliance with the License. @@ -15,21 +15,26 @@ */ package org.springframework.integration.kafka.config.xml; +import org.w3c.dom.Element; + +import org.springframework.beans.factory.config.BeanDefinition; import org.springframework.beans.factory.support.AbstractBeanDefinition; import org.springframework.beans.factory.support.BeanDefinitionBuilder; import org.springframework.beans.factory.xml.ParserContext; import org.springframework.integration.config.xml.AbstractOutboundChannelAdapterParser; +import org.springframework.integration.config.xml.IntegrationNamespaceUtils; import org.springframework.integration.kafka.outbound.KafkaProducerMessageHandler; import org.springframework.util.StringUtils; -import org.w3c.dom.Element; /** * * @author Soby Chacko + * @author Artem Bilan * @since 0.5 * */ public class KafkaOutboundChannelAdapterParser extends AbstractOutboundChannelAdapterParser { + @Override protected AbstractBeanDefinition parseConsumer(final Element element, final ParserContext parserContext) { final BeanDefinitionBuilder kafkaProducerMessageHandlerBuilder = @@ -41,6 +46,21 @@ public class KafkaOutboundChannelAdapterParser extends AbstractOutboundChannelAd kafkaProducerMessageHandlerBuilder.addConstructorArgReference(kafkaServerBeanName); } + BeanDefinition topicExpressionDef = + IntegrationNamespaceUtils.createExpressionDefinitionFromValueOrExpression("topic", "topic-expression", + parserContext, element, false); + if (topicExpressionDef != null) { + kafkaProducerMessageHandlerBuilder.addPropertyValue("topicExpression", topicExpressionDef); + } + + BeanDefinition messageKeyExpressionDef = + IntegrationNamespaceUtils.createExpressionDefinitionFromValueOrExpression("message-key", + "message-key-expression", parserContext, element, false); + if (messageKeyExpressionDef != null) { + kafkaProducerMessageHandlerBuilder.addPropertyValue("messageKeyExpression", messageKeyExpressionDef); + } + 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 894503b3b9..1cfafbb9f3 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 @@ -13,32 +13,72 @@ * See the License for the specific language governing permissions and * limitations under the License. */ + package org.springframework.integration.kafka.outbound; +import org.springframework.expression.EvaluationContext; +import org.springframework.expression.Expression; +import org.springframework.integration.expression.IntegrationEvaluationContextAware; import org.springframework.integration.handler.AbstractMessageHandler; +import org.springframework.integration.kafka.support.KafkaHeaders; import org.springframework.integration.kafka.support.KafkaProducerContext; import org.springframework.messaging.Message; +import org.springframework.util.Assert; /** * @author Soby Chacko + * @author Artem Bilan * @author Gary Russell * @since 0.5 */ -public class KafkaProducerMessageHandler extends AbstractMessageHandler { +public class KafkaProducerMessageHandler extends AbstractMessageHandler + implements IntegrationEvaluationContextAware { private final KafkaProducerContext kafkaProducerContext; + private EvaluationContext evaluationContext; + + private volatile Expression topicExpression; + + private volatile Expression messageKeyExpression; + public KafkaProducerMessageHandler(final KafkaProducerContext kafkaProducerContext) { this.kafkaProducerContext = kafkaProducerContext; } + @Override + public void setIntegrationEvaluationContext(EvaluationContext evaluationContext) { + this.evaluationContext = evaluationContext; + } + + public void setTopicExpression(Expression topicExpression) { + this.topicExpression = topicExpression; + } + + public void setMessageKeyExpression(Expression messageKeyExpression) { + this.messageKeyExpression = messageKeyExpression; + } + public KafkaProducerContext getKafkaProducerContext() { - return kafkaProducerContext; + return this.kafkaProducerContext; + } + + @Override + protected void onInit() throws Exception { + Assert.notNull(this.evaluationContext); } @Override protected void handleMessageInternal(final Message message) throws Exception { - kafkaProducerContext.send(message); + String topic = this.topicExpression != null + ? this.topicExpression.getValue(this.evaluationContext, message, String.class) + : message.getHeaders().get(KafkaHeaders.TOPIC, String.class); + + Object messageKey = this.messageKeyExpression != null + ? this.messageKeyExpression.getValue(this.evaluationContext, message) + : message.getHeaders().get(KafkaHeaders.MESSAGE_KEY); + + this.kafkaProducerContext.send(topic, messageKey, message); } @Override diff --git a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/support/KafkaHeaders.java b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/support/KafkaHeaders.java new file mode 100644 index 0000000000..148235812b --- /dev/null +++ b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/support/KafkaHeaders.java @@ -0,0 +1,31 @@ +/* + * Copyright 2014 the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.springframework.integration.kafka.support; + +/** + * @author Artem Bilan + * @since 1.0 + */ +public abstract class KafkaHeaders { + + private static final String PREFIX = "kafka_"; + + public static final String TOPIC = PREFIX + "topic"; + + public static final String MESSAGE_KEY = PREFIX + "messageKey"; + +} diff --git a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/support/KafkaProducerContext.java b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/support/KafkaProducerContext.java index d77edb5705..fd66562a71 100644 --- a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/support/KafkaProducerContext.java +++ b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/support/KafkaProducerContext.java @@ -18,7 +18,6 @@ package org.springframework.integration.kafka.support; import java.util.Collection; import java.util.Map; -import java.util.Properties; import java.util.concurrent.atomic.AtomicBoolean; import org.apache.commons.logging.Log; @@ -34,6 +33,7 @@ import org.springframework.messaging.Message; * @author Rajasekar Elango * @author Ilayaperumal Gopinathan * @author Gary Russell + * @author Artem Bilan * @since 0.5 */ public class KafkaProducerContext implements SmartLifecycle, NamedComponent, BeanNameAware { @@ -46,8 +46,6 @@ public class KafkaProducerContext implements SmartLifecycle, NamedComponen private volatile ProducerConfiguration theProducerConfiguration; - private Properties producerProperties; - private String beanName = "not_specified"; private int phase = 0; @@ -83,21 +81,6 @@ public class KafkaProducerContext implements SmartLifecycle, NamedComponen } } - /** - * @param producerProperties - * The producerProperties to set. - */ - public void setProducerProperties(Properties producerProperties) { - this.producerProperties = producerProperties; - } - - /** - * @return the producerProperties. - */ - public Properties getProducerProperties() { - return this.producerProperties; - } - /** * @return the component type. * @since 1.0 @@ -199,21 +182,19 @@ public class KafkaProducerContext implements SmartLifecycle, NamedComponen callback.run(); } - public void send(final Message message) throws Exception { + public void send(String topic, Object messageKey, final Message message) throws Exception { if (!running.get()) { start(); } - if (message.getHeaders().containsKey("topic")) { - ProducerConfiguration producerConfiguration = - getTopicConfiguration(message.getHeaders().get("topic", String.class)); - if (producerConfiguration != null) { - producerConfiguration.send(message); - } + ProducerConfiguration producerConfiguration = getTopicConfiguration(topic); + + if (producerConfiguration != null) { + producerConfiguration.send(topic, messageKey, message); } // if there is a single producer configuration then use that config to send message. else if (this.theProducerConfiguration != null) { - this.theProducerConfiguration.send(message); + this.theProducerConfiguration.send(null, messageKey, message); } else { throw new IllegalStateException("Could not send messages as there are multiple producer configurations " + 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 e7055401f7..b61e3d3362 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 @@ -26,6 +26,7 @@ import org.apache.commons.lang.builder.HashCodeBuilder; import org.springframework.messaging.Message; import org.springframework.messaging.MessageHandlingException; import org.springframework.util.Assert; +import org.springframework.util.StringUtils; import kafka.javaapi.producer.Producer; import kafka.producer.KeyedMessage; @@ -59,19 +60,14 @@ public class ProducerConfiguration { return this.producer; } - public void send(final Message message) throws Exception { + public void send(String topic, Object messageKey, final Message message) throws Exception { final V v = getPayload(message); - String topic = message.getHeaders().containsKey("topic") - ? message.getHeaders().get("topic", String.class) - : this.producerMetadata.getTopic(); + if (!StringUtils.hasText(topic)) { + topic = this.producerMetadata.getTopic(); + } - if (message.getHeaders().containsKey("messageKey")) { - this.producer.send(new KeyedMessage(topic, getKey(message), v)); - } - else { - this.producer.send(new KeyedMessage(topic, v)); - } + this.producer.send(new KeyedMessage(topic, (messageKey != null ? getKey(messageKey) : null), v)); } @SuppressWarnings("unchecked") @@ -86,13 +82,12 @@ public class ProducerConfiguration { } @SuppressWarnings("unchecked") - private K getKey(final Message message) throws Exception { - final Object key = message.getHeaders().get("messageKey"); - + private K getKey(Object messageKey) throws Exception { if (this.producerMetadata.getKeyEncoder() instanceof DefaultEncoder) { - return (K) getByteStream(key); + return (K) getByteStream(messageKey); } - return message.getHeaders().get("messageKey", this.producerMetadata.getKeyClassType()); + + return (K) messageKey; } private static boolean isRawByteArray(final Object obj) { 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 22e20e881f..6495585e09 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 @@ -1,6 +1,6 @@ - + - - - - - @@ -43,48 +38,31 @@ - - - - - + + + - + - - - - - - + - - - - - - + - - - - - @@ -120,79 +98,91 @@ - + - - - - - - - Custom implementation of a Kafka Encoder for encoding message values. - - - - - - - Custom implementation of a Kafka Encoder for encoding message keys. - - - - - - - Class type used for the key - - - - - - - Class type used for the value - - - - - - - - + + Custom implementation of a Kafka Encoder for encoding message values. + + + - + - - Custom Kafka key partitioner. - + + + Custom implementation of a Kafka Encoder for encoding message keys. + + + + + - + + + + + Class type used for the key + + + + + + + + + + + + Class type used for the value + + + + + + + + + + + + + + + + + Custom Kafka key partitioner. + + + + + + + + Indicates if this producer is async or not. + + + - + number of messages to batch at this producer. @@ -206,12 +196,16 @@ - + - - Kafka producer properties to use for all producers - + + + Kafka producer properties to use for all producers + + + + + @@ -252,14 +246,9 @@ - - - - - - + + Regex pattern to match topic + @@ -267,11 +256,6 @@ - - - - - @@ -281,11 +265,6 @@ If exclude is true, it uses blacklist to exclude topics matching given pattern. Default value is false. ]]> - - - - - @@ -294,44 +273,42 @@ - + - - - - - - + + + + + - - + + Custom implementation of a Kafka Decoder for values. + + + - + - - Kafka Server Bean Name - - - - - - - Kafka Server Bean Name - + + + Custom implementation of a Kafka Decoder for keys. + + + + + @@ -342,32 +319,30 @@ - + - - - - - - + Kafka Server Bean Name - + - - Kafka consumer properties to use for all consumers - + + + Kafka consumer properties to use for all consumers + + + + + @@ -402,23 +377,23 @@ - + - - - - - - - Kafka Server Bean Name - + + + Kafka consumer context reference. + + + + + @@ -439,9 +414,48 @@ - - Kafka producer context reference. - + + + Kafka producer context reference. + + + + + + + + + + + + + + + + + + + + + + + + + @@ -451,7 +465,11 @@ 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 d1cdd05443..3906c9df74 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 @@ -1,45 +1,48 @@ - - - + + + - - - + + + - + - - - + + + + + + + + + + - - - - - - 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 0eff80addd..a71e86f8f6 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 @@ -13,9 +13,12 @@ * See the License for the specific language governing permissions and * limitations under the License. */ + package org.springframework.integration.kafka.config.xml; -import org.junit.Assert; +import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertNotNull; + import org.junit.ClassRule; import org.junit.Test; import org.junit.runner.RunWith; @@ -26,11 +29,13 @@ import org.springframework.integration.endpoint.PollingConsumer; import org.springframework.integration.kafka.outbound.KafkaProducerMessageHandler; import org.springframework.integration.kafka.rule.KafkaRunning; import org.springframework.integration.kafka.support.KafkaProducerContext; +import org.springframework.integration.test.util.TestUtils; import org.springframework.test.context.ContextConfiguration; import org.springframework.test.context.junit4.SpringJUnit4ClassRunner; /** * @author Soby Chacko + * @author Artem Bilan * @author Gary Russell * @since 0.5 */ @@ -47,15 +52,16 @@ public class KafkaOutboundAdapterParserTests { @Test @SuppressWarnings("unchecked") public void testOutboundAdapterConfiguration() { - final PollingConsumer pollingConsumer = - appContext.getBean("kafkaOutboundChannelAdapter", PollingConsumer.class); - final KafkaProducerMessageHandler messageHandler = appContext.getBean(KafkaProducerMessageHandler.class); - Assert.assertNotNull(pollingConsumer); - Assert.assertNotNull(messageHandler); - Assert.assertEquals(messageHandler.getOrder(), 3); - final KafkaProducerContext producerContext = messageHandler.getKafkaProducerContext(); - Assert.assertNotNull(producerContext); - Assert.assertEquals(producerContext.getProducerConfigurations().size(), 2); + PollingConsumer pollingConsumer = this.appContext.getBean("kafkaOutboundChannelAdapter", PollingConsumer.class); + KafkaProducerMessageHandler messageHandler = this.appContext.getBean(KafkaProducerMessageHandler.class); + assertNotNull(pollingConsumer); + assertNotNull(messageHandler); + assertEquals(messageHandler.getOrder(), 3); + assertEquals("foo", TestUtils.getPropertyValue(messageHandler, "topicExpression.literalValue")); + assertEquals("'bar'", TestUtils.getPropertyValue(messageHandler, "messageKeyExpression.expression")); + KafkaProducerContext producerContext = messageHandler.getKafkaProducerContext(); + assertNotNull(producerContext); + assertEquals(producerContext.getProducerConfigurations().size(), 2); } } diff --git a/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/outbound/OutboundTests.java b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/outbound/OutboundTests.java index 240f7c0ab6..1471b3bac1 100644 --- a/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/outbound/OutboundTests.java +++ b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/outbound/OutboundTests.java @@ -23,9 +23,18 @@ import java.util.List; import java.util.Map; import java.util.Properties; +import kafka.admin.AdminUtils; +import kafka.consumer.ConsumerConfig; +import kafka.serializer.Decoder; +import kafka.serializer.Encoder; + +import org.junit.AfterClass; import org.junit.ClassRule; import org.junit.Test; +import org.springframework.context.expression.MapAccessor; +import org.springframework.expression.spel.standard.SpelExpressionParser; +import org.springframework.expression.spel.support.StandardEvaluationContext; import org.springframework.integration.kafka.rule.KafkaRunning; import org.springframework.integration.kafka.serializer.common.StringDecoder; import org.springframework.integration.kafka.serializer.common.StringEncoder; @@ -34,6 +43,7 @@ import org.springframework.integration.kafka.support.ConsumerConfiguration; import org.springframework.integration.kafka.support.ConsumerConnectionProvider; import org.springframework.integration.kafka.support.ConsumerMetadata; import org.springframework.integration.kafka.support.KafkaConsumerContext; +import org.springframework.integration.kafka.support.KafkaHeaders; import org.springframework.integration.kafka.support.KafkaProducerContext; import org.springframework.integration.kafka.support.MessageLeftOverTracker; import org.springframework.integration.kafka.support.ProducerConfiguration; @@ -43,10 +53,6 @@ import org.springframework.integration.kafka.support.ZookeeperConnect; import org.springframework.messaging.Message; import org.springframework.messaging.support.MessageBuilder; -import kafka.consumer.ConsumerConfig; -import kafka.serializer.Decoder; -import kafka.serializer.Encoder; - /** * @author Gary Russell * @since 1.0 @@ -59,6 +65,15 @@ public class OutboundTests { @ClassRule public static KafkaRunning kafkaRunning = KafkaRunning.isRunning(); + @AfterClass + public static void tearDown() { + try { + AdminUtils.deleteTopic(kafkaRunning.getZkClient(), TOPIC); + } + catch (Exception e) { + } + } + @Test public void testAsyncProducerFlushed() throws Exception { KafkaConsumerContext consumerContext = createConsumer(); @@ -83,14 +98,29 @@ public class OutboundTests { KafkaProducerMessageHandler handler = new KafkaProducerMessageHandler(kafkaProducerContext); handler.handleMessage(MessageBuilder.withPayload("foo") - .setHeader("messagekey", "3") - .setHeader("topic", TOPIC) + .setHeader(KafkaHeaders.MESSAGE_KEY, "3") + .setHeader(KafkaHeaders.TOPIC, TOPIC) + .build()); + + SpelExpressionParser parser = new SpelExpressionParser(); + handler.setMessageKeyExpression(parser.parseExpression("headers.foo")); + handler.setTopicExpression(parser.parseExpression("headers.bar")); + StandardEvaluationContext evaluationContext = new StandardEvaluationContext(); + evaluationContext.addPropertyAccessor(new MapAccessor()); + handler.setIntegrationEvaluationContext(evaluationContext); + handler.handleMessage(MessageBuilder.withPayload("bar") + .setHeader("foo", "3") + .setHeader("bar", TOPIC) .build()); kafkaProducerContext.stop(); Message>>> received = consumerContext.receive(); assertNotNull(received); + if (((Map) received.getPayload()).size() < 2) { + received = consumerContext.receive(); + assertNotNull(received); + } consumerContext.destroy(); } diff --git a/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/rule/KafkaRunning.java b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/rule/KafkaRunning.java index 9487148c67..fdb4495b2b 100644 --- a/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/rule/KafkaRunning.java +++ b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/rule/KafkaRunning.java @@ -15,11 +15,7 @@ */ package org.springframework.integration.kafka.rule; -import java.io.IOException; -import java.net.Socket; - -import javax.net.SocketFactory; - +import org.I0Itec.zkclient.ZkClient; import org.apache.commons.logging.Log; import org.apache.commons.logging.LogFactory; import org.junit.Assume; @@ -27,6 +23,11 @@ import org.junit.rules.TestWatcher; import org.junit.runner.Description; import org.junit.runners.model.Statement; +import org.springframework.integration.kafka.core.ZookeeperConnectDefaults; + +import kafka.utils.ZKStringSerializer$; +import kafka.utils.ZkUtils; + /** *

* A rule that prevents integration tests from failing if the Kafka server is not running or not @@ -44,9 +45,7 @@ import org.junit.runners.model.Statement; */ public class KafkaRunning extends TestWatcher { - public static final int KAFKA_PORT = 9092; - - public static final int ZOOKEEPER_PORT = 2181; + private static final String ZOOKEEPER_CONNECT_STRING = ZookeeperConnectDefaults.ZK_CONNECT; private static final Log logger = LogFactory.getLog(KafkaRunning.class); @@ -57,37 +56,24 @@ public class KafkaRunning extends TestWatcher { return new KafkaRunning(); } + private ZkClient zkClient; + + public ZkClient getZkClient() { + return zkClient; + } + @Override public Statement apply(Statement base, Description description) { - - Socket kSocket = null; - Socket zSocket = null; try { - kSocket = SocketFactory.getDefault().createSocket("localhost", KAFKA_PORT); - kSocket.getInputStream(); - zSocket = SocketFactory.getDefault().createSocket("localhost", ZOOKEEPER_PORT); - zSocket.getInputStream(); + this.zkClient = new ZkClient(ZOOKEEPER_CONNECT_STRING, 1000, 1000, ZKStringSerializer$.MODULE$); + if (ZkUtils.getAllBrokersInCluster(zkClient).size() == 0) { + throw new IllegalStateException("No running Kafka brokers"); + } } - catch (final Exception e) { + catch (Exception e) { logger.warn("Not executing tests because basic connectivity test failed"); Assume.assumeNoException(e); } - finally { - if (kSocket != null) { - try { - kSocket.close(); - } - catch (IOException e) { - } - } - if (zSocket != null) { - try { - zSocket.close(); - } - catch (IOException e) { - } - } - } return super.apply(base, description); } diff --git a/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/support/ProducerConfigurationTests.java b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/support/ProducerConfigurationTests.java index cbdad2c771..bc736e8ee7 100644 --- a/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/support/ProducerConfigurationTests.java +++ b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/support/ProducerConfigurationTests.java @@ -1,5 +1,5 @@ /* - * Copyright 2002-2013 the original author or authors. + * Copyright 2013-2014 the original author or authors. * * Licensed under the Apache License, Version 2.0 (the "License"); * you may not use this file except in compliance with the License. @@ -13,33 +13,38 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package org.springframework.integration.kafka.support; -import kafka.javaapi.producer.Producer; -import kafka.producer.KeyedMessage; -import kafka.serializer.DefaultEncoder; -import kafka.serializer.StringEncoder; -import org.junit.Assert; -import org.junit.Test; -import org.mockito.ArgumentCaptor; -import org.mockito.Mockito; -import org.springframework.integration.kafka.serializer.avro.AvroReflectDatumBackedKafkaEncoder; -import org.springframework.integration.kafka.test.utils.NonSerializableTestKey; -import org.springframework.integration.kafka.test.utils.NonSerializableTestPayload; -import org.springframework.integration.kafka.test.utils.TestKey; -import org.springframework.integration.kafka.test.utils.TestPayload; -import org.springframework.integration.support.MessageBuilder; -import org.springframework.messaging.Message; +package org.springframework.integration.kafka.support; import java.io.ByteArrayInputStream; import java.io.NotSerializableException; import java.io.ObjectInputStream; +import org.junit.Assert; +import org.junit.Test; +import org.mockito.ArgumentCaptor; +import org.mockito.Mockito; + +import org.springframework.integration.kafka.serializer.avro.AvroReflectDatumBackedKafkaEncoder; +import org.springframework.integration.kafka.test.utils.NonSerializableTestKey; +import org.springframework.integration.kafka.test.utils.NonSerializableTestPayload; +import org.springframework.integration.kafka.test.utils.TestKey; +import org.springframework.integration.kafka.test.utils.TestPayload; +import org.springframework.messaging.Message; +import org.springframework.messaging.support.GenericMessage; + +import kafka.javaapi.producer.Producer; +import kafka.producer.KeyedMessage; +import kafka.serializer.DefaultEncoder; +import kafka.serializer.StringEncoder; + /** * @author Soby Chacko + * @author Artem Bilan * @since 0.5 */ -public class ProducerConfigurationTests { +public class ProducerConfigurationTests { + @Test @SuppressWarnings("unchecked") public void testSendMessageWithNonDefaultKeyAndValueEncoders() throws Exception { @@ -50,20 +55,16 @@ public class ProducerConfigurationTests { producerMetadata.setValueClassType(String.class); final Producer producer = Mockito.mock(Producer.class); - final ProducerConfiguration configuration = new ProducerConfiguration(producerMetadata, producer); + final ProducerConfiguration configuration = + new ProducerConfiguration(producerMetadata, producer); - final Message message = MessageBuilder.withPayload("test message") - .setHeader("messageKey", "key") - .setHeader("topic", "test") - .build(); - - configuration.send(message); + configuration.send("test", "key", new GenericMessage("test message")); Mockito.verify(producer, Mockito.times(1)).send(Mockito.any(KeyedMessage.class)); final ArgumentCaptor> argument = (ArgumentCaptor>) (Object) - ArgumentCaptor.forClass(KeyedMessage.class); + ArgumentCaptor.forClass(KeyedMessage.class); Mockito.verify(producer).send(argument.capture()); final KeyedMessage capturedKeyMessage = argument.getValue(); @@ -78,48 +79,47 @@ public class ProducerConfigurationTests { */ @Test @SuppressWarnings("unchecked") - public void testSendMessageWithDefaultKeyAndValueEncodersAndCustomSerializableKeyAndPayloadObject() throws Exception { + public void testSendMessageWithDefaultKeyAndValueEncodersAndCustomSerializableKeyAndPayloadObject() + throws Exception { final ProducerMetadata producerMetadata = new ProducerMetadata("test"); producerMetadata.setValueEncoder(new DefaultEncoder(null)); producerMetadata.setKeyEncoder(new DefaultEncoder(null)); final Producer producer = Mockito.mock(Producer.class); - final ProducerConfiguration configuration = new ProducerConfiguration(producerMetadata, producer); + final ProducerConfiguration configuration = + new ProducerConfiguration(producerMetadata, producer); - final Message message = MessageBuilder.withPayload(new TestPayload("part1", "part2")) - .setHeader("messageKey", new TestKey("compositePart1", "compositePart2")) - .setHeader("topic", "test") - .build(); + Message message = new GenericMessage(new TestPayload("part1", "part2")); - configuration.send(message); + configuration.send("test", new TestKey("compositePart1", "compositePart2"), message); Mockito.verify(producer, Mockito.times(1)).send(Mockito.any(KeyedMessage.class)); final ArgumentCaptor> argument = (ArgumentCaptor>) (Object) - ArgumentCaptor.forClass(KeyedMessage.class); + ArgumentCaptor.forClass(KeyedMessage.class); Mockito.verify(producer).send(argument.capture()); final KeyedMessage capturedKeyMessage = argument.getValue(); final byte[] keyBytes = capturedKeyMessage.key(); - final ByteArrayInputStream keyInputStream = new ByteArrayInputStream (keyBytes); - final ObjectInputStream keyObjectInputStream = new ObjectInputStream (keyInputStream); + final ByteArrayInputStream keyInputStream = new ByteArrayInputStream(keyBytes); + final ObjectInputStream keyObjectInputStream = new ObjectInputStream(keyInputStream); final Object keyObj = keyObjectInputStream.readObject(); - final TestKey tk = (TestKey)keyObj; + final TestKey tk = (TestKey) keyObj; Assert.assertEquals(tk.getKeyPart1(), "compositePart1"); Assert.assertEquals(tk.getKeyPart2(), "compositePart2"); final byte[] messageBytes = capturedKeyMessage.message(); - final ByteArrayInputStream messageInputStream = new ByteArrayInputStream (messageBytes); - final ObjectInputStream messageObjectInputStream = new ObjectInputStream (messageInputStream); + final ByteArrayInputStream messageInputStream = new ByteArrayInputStream(messageBytes); + final ObjectInputStream messageObjectInputStream = new ObjectInputStream(messageInputStream); final Object messageObj = messageObjectInputStream.readObject(); - final TestPayload tp = (TestPayload)messageObj; + final TestPayload tp = (TestPayload) messageObj; Assert.assertEquals(tp.getPart1(), "part1"); Assert.assertEquals(tp.getPart2(), "part2"); @@ -134,34 +134,32 @@ public class ProducerConfigurationTests { @SuppressWarnings("unchecked") public void testSendMessageWithDefaultKeyEncoderAndNonDefaultValueEncoderAndCorrespondingData() throws Exception { final ProducerMetadata producerMetadata = new ProducerMetadata("test"); - final AvroReflectDatumBackedKafkaEncoder encoder = new AvroReflectDatumBackedKafkaEncoder(TestPayload.class); + final AvroReflectDatumBackedKafkaEncoder encoder = + new AvroReflectDatumBackedKafkaEncoder(TestPayload.class); producerMetadata.setValueEncoder(encoder); producerMetadata.setKeyEncoder(new DefaultEncoder(null)); producerMetadata.setValueClassType(TestPayload.class); final Producer producer = Mockito.mock(Producer.class); - final ProducerConfiguration configuration = new ProducerConfiguration(producerMetadata, producer); - final TestPayload tp = new TestPayload("part1", "part2"); - final Message message = MessageBuilder.withPayload(tp) - .setHeader("messageKey", "key") - .setHeader("topic", "test") - .build(); + final ProducerConfiguration configuration = + new ProducerConfiguration(producerMetadata, producer); - configuration.send(message); + TestPayload tp = new TestPayload("part1", "part2"); + configuration.send("test", "key", new GenericMessage(tp)); Mockito.verify(producer, Mockito.times(1)).send(Mockito.any(KeyedMessage.class)); final ArgumentCaptor> argument = (ArgumentCaptor>) (Object) - ArgumentCaptor.forClass(KeyedMessage.class); + ArgumentCaptor.forClass(KeyedMessage.class); Mockito.verify(producer).send(argument.capture()); final KeyedMessage capturedKeyMessage = argument.getValue(); final byte[] keyBytes = capturedKeyMessage.key(); - final ByteArrayInputStream keyInputStream = new ByteArrayInputStream (keyBytes); - final ObjectInputStream keyObjectInputStream = new ObjectInputStream (keyInputStream); + final ByteArrayInputStream keyInputStream = new ByteArrayInputStream(keyBytes); + final ObjectInputStream keyObjectInputStream = new ObjectInputStream(keyInputStream); final Object keyObj = keyObjectInputStream.readObject(); Assert.assertEquals("key", keyObj); @@ -183,19 +181,18 @@ public class ProducerConfigurationTests { producerMetadata.setKeyClassType(TestKey.class); final Producer producer = Mockito.mock(Producer.class); - final ProducerConfiguration configuration = new ProducerConfiguration(producerMetadata, producer); - final TestKey tk = new TestKey("part1", "part2"); - final Message message = MessageBuilder.withPayload("test message"). - setHeader("messageKey", tk) - .setHeader("topic", "test").build(); + final ProducerConfiguration configuration = + new ProducerConfiguration(producerMetadata, producer); - configuration.send(message); + final TestKey tk = new TestKey("part1", "part2"); + + configuration.send("test", tk, new GenericMessage("test message")); Mockito.verify(producer, Mockito.times(1)).send(Mockito.any(KeyedMessage.class)); final ArgumentCaptor> argument = (ArgumentCaptor>) (Object) - ArgumentCaptor.forClass(KeyedMessage.class); + ArgumentCaptor.forClass(KeyedMessage.class); Mockito.verify(producer).send(argument.capture()); final KeyedMessage capturedKeyMessage = argument.getValue(); @@ -204,8 +201,8 @@ public class ProducerConfigurationTests { final byte[] payloadBytes = capturedKeyMessage.message(); - final ByteArrayInputStream payloadBis = new ByteArrayInputStream (payloadBytes); - final ObjectInputStream payloadOis = new ObjectInputStream (payloadBis); + final ByteArrayInputStream payloadBis = new ByteArrayInputStream(payloadBytes); + final ObjectInputStream payloadOis = new ObjectInputStream(payloadBis); final Object payloadObj = payloadOis.readObject(); Assert.assertEquals("test message", payloadObj); @@ -224,34 +221,31 @@ public class ProducerConfigurationTests { producerMetadata.setKeyEncoder(new DefaultEncoder(null)); final Producer producer = Mockito.mock(Producer.class); - final ProducerConfiguration configuration = new ProducerConfiguration(producerMetadata, producer); + final ProducerConfiguration configuration = + new ProducerConfiguration(producerMetadata, producer); - final Message message = MessageBuilder.withPayload("test message"). - setHeader("messageKey", "key") - .setHeader("topic", "test").build(); - - configuration.send(message); + configuration.send("test", "key", new GenericMessage("test message")); Mockito.verify(producer, Mockito.times(1)).send(Mockito.any(KeyedMessage.class)); final ArgumentCaptor> argument = (ArgumentCaptor>) (Object) - ArgumentCaptor.forClass(KeyedMessage.class); + ArgumentCaptor.forClass(KeyedMessage.class); Mockito.verify(producer).send(argument.capture()); final KeyedMessage capturedKeyMessage = argument.getValue(); final byte[] keyBytes = capturedKeyMessage.key(); - final ByteArrayInputStream keyBis = new ByteArrayInputStream (keyBytes); - final ObjectInputStream keyOis = new ObjectInputStream (keyBis); + final ByteArrayInputStream keyBis = new ByteArrayInputStream(keyBytes); + final ObjectInputStream keyOis = new ObjectInputStream(keyBis); final Object keyObj = keyOis.readObject(); Assert.assertEquals("key", keyObj); final byte[] payloadBytes = capturedKeyMessage.message(); - final ByteArrayInputStream payloadBis = new ByteArrayInputStream (payloadBytes); - final ObjectInputStream payloadOis = new ObjectInputStream (payloadBis); + final ByteArrayInputStream payloadBis = new ByteArrayInputStream(payloadBytes); + final ObjectInputStream payloadOis = new ObjectInputStream(payloadBis); final Object payloadObj = payloadOis.readObject(); Assert.assertEquals("test message", payloadObj); @@ -269,12 +263,13 @@ public class ProducerConfigurationTests { producerMetadata.setKeyEncoder(new DefaultEncoder(null)); final Producer producer = Mockito.mock(Producer.class); - final ProducerConfiguration configuration = new ProducerConfiguration(producerMetadata, producer); + final ProducerConfiguration configuration = + new ProducerConfiguration(producerMetadata, producer); - final Message message = MessageBuilder.withPayload(new NonSerializableTestPayload("part1", "part2")). - setHeader("messageKey", new NonSerializableTestKey("compositePart1", "compositePart2")) - .setHeader("topic", "test").build(); - configuration.send(message); + Message message = + new GenericMessage(new NonSerializableTestPayload("part1", "part2")); + + configuration.send("test", new NonSerializableTestKey("compositePart1", "compositePart2"), message); } /** @@ -282,18 +277,19 @@ public class ProducerConfigurationTests { */ @Test(expected = NotSerializableException.class) @SuppressWarnings("unchecked") - public void testSendMessageWithDefaultKeyAndValueEncodersButNonSerializableKeyAndSerializableValue() throws Exception { + public void testSendMessageWithDefaultKeyAndValueEncodersButNonSerializableKeyAndSerializableValue() + throws Exception { final ProducerMetadata producerMetadata = new ProducerMetadata("test"); producerMetadata.setValueEncoder(new DefaultEncoder(null)); producerMetadata.setKeyEncoder(new DefaultEncoder(null)); final Producer producer = Mockito.mock(Producer.class); - final ProducerConfiguration configuration = new ProducerConfiguration(producerMetadata, producer); + final ProducerConfiguration configuration = + new ProducerConfiguration(producerMetadata, producer); - final Message message = MessageBuilder.withPayload(new TestPayload("part1", "part2")). - setHeader("messageKey", new NonSerializableTestKey("compositePart1", "compositePart2")) - .setHeader("topic", "test").build(); - configuration.send(message); + Message message = new GenericMessage(new TestPayload("part1", "part2")); + + configuration.send("test", new NonSerializableTestKey("compositePart1", "compositePart2"), message); } /** @@ -301,17 +297,20 @@ public class ProducerConfigurationTests { */ @Test(expected = NotSerializableException.class) @SuppressWarnings("unchecked") - public void testSendMessageWithDefaultKeyAndValueEncodersButSerializableKeyAndNonSerializableValue() throws Exception { + public void testSendMessageWithDefaultKeyAndValueEncodersButSerializableKeyAndNonSerializableValue() + throws Exception { final ProducerMetadata producerMetadata = new ProducerMetadata("test"); producerMetadata.setValueEncoder(new DefaultEncoder(null)); producerMetadata.setKeyEncoder(new DefaultEncoder(null)); final Producer producer = Mockito.mock(Producer.class); - final ProducerConfiguration configuration = new ProducerConfiguration(producerMetadata, producer); + final ProducerConfiguration configuration = + new ProducerConfiguration(producerMetadata, producer); - final Message message = MessageBuilder.withPayload(new NonSerializableTestPayload("part1", "part2")). - setHeader("messageKey", new TestKey("compositePart1", "compositePart2")) - .setHeader("topic", "test").build(); - configuration.send(message); + Message message = + new GenericMessage(new NonSerializableTestPayload("part1", "part2")); + + configuration.send("test", new TestKey("compositePart1", "compositePart2"), message); } + }