From 0e2808c0140aef3b7edf1fd913dd09d629218cfb Mon Sep 17 00:00:00 2001 From: Gary Russell Date: Tue, 2 Dec 2014 21:50:37 +0200 Subject: [PATCH] INTEXT-121 Add SmartLifecycle to ProducerContext JIRA: https://jira.spring.io/browse/INTEXT-121 When using `async` producers, the buffered messages need to be flushed. Previously, when the context was stopped, such messages were lost. - Implement `SmartLifecyle` in the producer context and propagate the `stop()` to the underlying producer(s). - Remove unused attributes from the consumer parser. - Add a test case to show the buffered message is received after closing the producer. - Add KafkaRunning JUnit `@Rule` Add KafkaRunning Rule to Parser Tests Polishing --- .../xml/KafkaConsumerContextParser.java | 5 +- .../xml/KafkaProducerContextParser.java | 20 ++- .../outbound/KafkaProducerMessageHandler.java | 9 +- .../kafka/support/KafkaProducerContext.java | 162 +++++++++++++++--- .../kafka/support/ProducerConfiguration.java | 12 +- .../xml/spring-integration-kafka-1.0.xsd | 1 + .../xml/KafkaConsumerContextParserTests.java | 10 +- .../xml/KafkaInboundAdapterParserTests.java | 9 +- .../KafkaMultiConsumerContextParserTests.java | 7 + .../xml/KafkaOutboundAdapterParserTests.java | 8 +- ...afkaProducerContextParserTests-context.xml | 3 +- .../xml/KafkaProducerContextParserTests.java | 41 +++-- .../xml/ZookeeperConnectParserTests.java | 9 +- .../kafka/outbound/OutboundTests.java | 122 +++++++++++++ .../integration/kafka/rule/KafkaRunning.java | 95 ++++++++++ 15 files changed, 459 insertions(+), 54 deletions(-) create mode 100644 spring-integration-kafka/src/test/java/org/springframework/integration/kafka/outbound/OutboundTests.java create mode 100644 spring-integration-kafka/src/test/java/org/springframework/integration/kafka/rule/KafkaRunning.java diff --git a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/config/xml/KafkaConsumerContextParser.java b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/config/xml/KafkaConsumerContextParser.java index 7604d0d9a8..fa1700f76e 100644 --- a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/config/xml/KafkaConsumerContextParser.java +++ b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/config/xml/KafkaConsumerContextParser.java @@ -45,6 +45,7 @@ import org.springframework.util.xml.DomUtils; * @author Rajasekar Elango * @author Artem Bilan * @author Ilayaperumal Gopinathan + * @author Gary Russell * @since 0.5 */ public class KafkaConsumerContextParser extends AbstractSingleBeanDefinitionParser { @@ -78,10 +79,6 @@ public class KafkaConsumerContextParser extends AbstractSingleBeanDefinitionPars "value-decoder"); IntegrationNamespaceUtils.setReferenceIfAttributeDefined(consumerMetadataBuilder, consumerConfiguration, "key-decoder"); - IntegrationNamespaceUtils.setValueIfAttributeDefined(consumerMetadataBuilder, consumerConfiguration, - "key-class-type"); - IntegrationNamespaceUtils.setValueIfAttributeDefined(consumerMetadataBuilder, consumerConfiguration, - "value-class-type"); IntegrationNamespaceUtils.setValueIfAttributeDefined(consumerConfigurationBuilder, consumerConfiguration, "max-messages"); IntegrationNamespaceUtils.setValueIfAttributeDefined(consumerMetadataBuilder, parentElem, diff --git a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/config/xml/KafkaProducerContextParser.java b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/config/xml/KafkaProducerContextParser.java index cfb4c3c650..9a6c0418d0 100644 --- a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/config/xml/KafkaProducerContextParser.java +++ b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/config/xml/KafkaProducerContextParser.java @@ -37,6 +37,7 @@ import org.springframework.util.xml.DomUtils; /** * @author Soby Chacko * @author Ilayaperumal Gopinathan + * @author Gary Russell * @since 0.5 */ public class KafkaProducerContextParser extends AbstractSimpleBeanDefinitionParser { @@ -50,17 +51,24 @@ public class KafkaProducerContextParser extends AbstractSimpleBeanDefinitionPars protected void doParse(final Element element, final ParserContext parserContext, final BeanDefinitionBuilder builder) { super.doParse(element, parserContext, builder); + IntegrationNamespaceUtils.setValueIfAttributeDefined(builder, element, "phase"); + IntegrationNamespaceUtils.setValueIfAttributeDefined(builder, element, "auto-startup"); + final Element topics = DomUtils.getChildElementByTagName(element, "producer-configurations"); parseProducerConfigurations(topics, parserContext, builder, element); } + + @Override + protected boolean isEligibleAttribute(String attributeName) { + return !"producer-properties".equals(attributeName) && super.isEligibleAttribute(attributeName); + } + private void parseProducerConfigurations(Element topics, ParserContext parserContext, BeanDefinitionBuilder builder, Element parentElem) { Map producerConfigurationsMap = new ManagedMap(); for (Element producerConfiguration : DomUtils.getChildElementsByTagName(topics, "producer-configuration")) { - BeanDefinitionBuilder producerConfigurationBuilder = - BeanDefinitionBuilder.genericBeanDefinition(ProducerConfiguration.class); BeanDefinitionBuilder producerMetadataBuilder = BeanDefinitionBuilder.genericBeanDefinition(ProducerMetadata.class); @@ -100,11 +108,11 @@ public class KafkaProducerContextParser extends AbstractSimpleBeanDefinitionPars AbstractBeanDefinition producerFactoryBeanDefinition = producerFactoryBuilder.getBeanDefinition(); - producerConfigurationBuilder.addConstructorArgValue(producerMetadataBeanDefinition); - producerConfigurationBuilder.addConstructorArgValue(producerFactoryBeanDefinition); - AbstractBeanDefinition producerConfigurationBeanDefinition = - producerConfigurationBuilder.getBeanDefinition(); + BeanDefinitionBuilder.genericBeanDefinition(ProducerConfiguration.class) + .addConstructorArgValue(producerMetadataBeanDefinition) + .addConstructorArgValue(producerFactoryBeanDefinition) + .getBeanDefinition(); producerConfigurationsMap.put(producerConfiguration.getAttribute("topic"), producerConfigurationBeanDefinition); } 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 a9062e7b06..894503b3b9 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 @@ -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. @@ -21,6 +21,7 @@ import org.springframework.messaging.Message; /** * @author Soby Chacko + * @author Gary Russell * @since 0.5 */ public class KafkaProducerMessageHandler extends AbstractMessageHandler { @@ -39,4 +40,10 @@ public class KafkaProducerMessageHandler extends AbstractMessageHandler { protected void handleMessageInternal(final Message message) throws Exception { kafkaProducerContext.send(message); } + + @Override + public String getComponentType() { + return "kafka:outbound-channel-adapter"; + } + } 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 87b5678b6a..d77edb5705 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 @@ -1,5 +1,5 @@ /* - * Copyright 2002-2014 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. @@ -19,21 +19,28 @@ 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; import org.apache.commons.logging.LogFactory; +import org.springframework.beans.factory.BeanNameAware; +import org.springframework.context.SmartLifecycle; +import org.springframework.integration.support.context.NamedComponent; import org.springframework.messaging.Message; /** * @author Soby Chacko * @author Rajasekar Elango * @author Ilayaperumal Gopinathan + * @author Gary Russell * @since 0.5 */ -public class KafkaProducerContext { +public class KafkaProducerContext implements SmartLifecycle, NamedComponent, BeanNameAware { - private static final Log LOGGER = LogFactory.getLog(KafkaProducerContext.class); + private static final Log logger = LogFactory.getLog(KafkaProducerContext.class); + + private final AtomicBoolean running = new AtomicBoolean(); private volatile Map> producerConfigurations; @@ -41,23 +48,11 @@ public class KafkaProducerContext { private Properties producerProperties; - public void send(final Message message) throws Exception { - if (message.getHeaders().containsKey("topic")) { - ProducerConfiguration producerConfiguration = - getTopicConfiguration(message.getHeaders().get("topic", String.class)); - if (producerConfiguration != null) { - producerConfiguration.send(message); - } - } - // if there is a single producer configuration then use that config to send message. - else if (this.theProducerConfiguration != null) { - this.theProducerConfiguration.send(message); - } - else { - throw new IllegalStateException("Could not send messages as there are multiple producer configurations " + - "with no topic information found from the message header."); - } - } + private String beanName = "not_specified"; + + private int phase = 0; + + private boolean autoStartup = true; public ProducerConfiguration getTopicConfiguration(final String topic) { if (this.theProducerConfiguration != null) { @@ -73,7 +68,7 @@ public class KafkaProducerContext { return producerConfiguration; } } - LOGGER.error("No producer-configuration defined for topic " + topic + ". Cannot send message"); + logger.error("No producer-configuration defined for topic " + topic + ". Cannot send message"); return null; } @@ -97,10 +92,133 @@ public class KafkaProducerContext { } /** - * @return Returns the producerProperties. + * @return the producerProperties. */ public Properties getProducerProperties() { return this.producerProperties; } + /** + * @return the component type. + * @since 1.0 + */ + @Override + public String getComponentType() { + return "kafka:producer-context"; + } + + /** + * @param name the bean name. + * @since 1.0 + */ + @Override + public void setBeanName(String name) { + this.beanName = name; + } + + /** + * @param phase the phase to set. + * @see SmartLifecycle + * @since 1.0 + */ + public void setPhase(int phase) { + this.phase = phase; + } + + /** + * @param autoStartup the autoStartup to set. + * @see SmartLifecycle + * @since 1.0 + */ + public void setAutoStartup(boolean autoStartup) { + this.autoStartup = autoStartup; + } + + /** + * @return the component name. + * @since 1.0 + */ + @Override + public String getComponentName() { + return this.beanName; + } + + protected void doStart() { + } + + protected void doStop() { + if (this.producerConfigurations != null) { + for (ProducerConfiguration producerConfiguration : this.producerConfigurations.values()) { + producerConfiguration.stop(); + } + } + } + + @Override + public final void start() { + if (this.running.compareAndSet(false, true)) { + doStart(); + } + else { + if (logger.isDebugEnabled()) { + logger.debug(getComponentType() + ":" + getComponentName() + " is already running"); + } + } + } + + @Override + public final void stop() { + if (this.running.compareAndSet(true, false)) { + doStop(); + } + else { + if (logger.isDebugEnabled()) { + logger.debug(getComponentType() + ":" + getComponentName() + " is not running"); + } + } + } + + @Override + public boolean isRunning() { + return this.running.get(); + } + + @Override + public int getPhase() { + return this.phase; + } + + @Override + public boolean isAutoStartup() { + return this.autoStartup; + } + + @Override + public void stop(Runnable callback) { + stop(); + callback.run(); + } + + public void send(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); + } + } + // if there is a single producer configuration then use that config to send message. + else if (this.theProducerConfiguration != null) { + this.theProducerConfiguration.send(message); + } + else { + throw new IllegalStateException("Could not send messages as there are multiple producer configurations " + + "with no topic information found from the message header."); + } + } + } 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 cab8925b60..e7055401f7 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 @@ -1,5 +1,5 @@ /* - * Copyright 2002-2014 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. @@ -25,6 +25,7 @@ import org.apache.commons.lang.builder.HashCodeBuilder; import org.springframework.messaging.Message; import org.springframework.messaging.MessageHandlingException; +import org.springframework.util.Assert; import kafka.javaapi.producer.Producer; import kafka.producer.KeyedMessage; @@ -34,6 +35,7 @@ import kafka.serializer.DefaultEncoder; * @author Soby Chacko * @author Rajasekar Elango * @author Ilayaperumal Gopinathan + * @author Gary Russell * @since 0.5 */ public class ProducerConfiguration { @@ -42,7 +44,9 @@ public class ProducerConfiguration { private final ProducerMetadata producerMetadata; - public ProducerConfiguration(final ProducerMetadata producerMetadata, final Producer producer) { + public ProducerConfiguration(ProducerMetadata producerMetadata, Producer producer) { + Assert.notNull(producerMetadata); + Assert.notNull(producer); this.producerMetadata = producerMetadata; this.producer = producer; } @@ -120,4 +124,8 @@ public class ProducerConfiguration { return "ProducerConfiguration [producerMetadata=" + this.producerMetadata + "]"; } + public void stop() { + this.producer.close(); + } + } 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 66c65d1196..22e20e881f 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 @@ -214,6 +214,7 @@ + diff --git a/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/config/xml/KafkaConsumerContextParserTests.java b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/config/xml/KafkaConsumerContextParserTests.java index 5d43b73583..84e1320320 100644 --- a/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/config/xml/KafkaConsumerContextParserTests.java +++ b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/config/xml/KafkaConsumerContextParserTests.java @@ -1,5 +1,5 @@ /* - * Copyright 2002-2014 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. @@ -19,15 +19,17 @@ package org.springframework.integration.kafka.config.xml; import static org.junit.Assert.assertEquals; import static org.junit.Assert.assertNotNull; import static org.junit.Assert.assertThat; - import kafka.consumer.Blacklist; + import org.hamcrest.Matchers; import org.junit.Assert; +import org.junit.ClassRule; 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.kafka.rule.KafkaRunning; import org.springframework.integration.kafka.support.ConsumerConfiguration; import org.springframework.integration.kafka.support.ConsumerMetadata; import org.springframework.integration.kafka.support.KafkaConsumerContext; @@ -38,12 +40,16 @@ import org.springframework.test.context.junit4.SpringJUnit4ClassRunner; /** * @author Soby Chacko * @author Artem Bilan + * @author Gary Russell * @since 0.5 */ @RunWith(SpringJUnit4ClassRunner.class) @ContextConfiguration public class KafkaConsumerContextParserTests { + @ClassRule + public static KafkaRunning kafkaRunning = KafkaRunning.isRunning(); + @Autowired private ApplicationContext appContext; diff --git a/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/config/xml/KafkaInboundAdapterParserTests.java b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/config/xml/KafkaInboundAdapterParserTests.java index f5e9486e3a..28fb204eb3 100644 --- a/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/config/xml/KafkaInboundAdapterParserTests.java +++ b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/config/xml/KafkaInboundAdapterParserTests.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. @@ -16,22 +16,29 @@ package org.springframework.integration.kafka.config.xml; import org.junit.Assert; +import org.junit.ClassRule; 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.SourcePollingChannelAdapter; +import org.springframework.integration.kafka.rule.KafkaRunning; import org.springframework.test.context.ContextConfiguration; import org.springframework.test.context.junit4.SpringJUnit4ClassRunner; /** * @author Soby Chacko + * @author Gary Russell * @since 0.5 */ @RunWith(SpringJUnit4ClassRunner.class) @ContextConfiguration public class KafkaInboundAdapterParserTests { + @ClassRule + public static KafkaRunning kafkaRunning = KafkaRunning.isRunning(); + @Autowired private ApplicationContext appContext; diff --git a/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/config/xml/KafkaMultiConsumerContextParserTests.java b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/config/xml/KafkaMultiConsumerContextParserTests.java index 84f9b958b8..e0ee27f97c 100644 --- a/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/config/xml/KafkaMultiConsumerContextParserTests.java +++ b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/config/xml/KafkaMultiConsumerContextParserTests.java @@ -16,10 +16,13 @@ package org.springframework.integration.kafka.config.xml; import org.junit.Assert; +import org.junit.ClassRule; 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.kafka.rule.KafkaRunning; import org.springframework.integration.kafka.support.ConsumerConfiguration; import org.springframework.integration.kafka.support.ConsumerMetadata; import org.springframework.integration.kafka.support.KafkaConsumerContext; @@ -28,11 +31,15 @@ import org.springframework.test.context.junit4.SpringJUnit4ClassRunner; /** * @author Ilayaperumal Gopinathan + * @author Gary Russell */ @RunWith(SpringJUnit4ClassRunner.class) @ContextConfiguration public class KafkaMultiConsumerContextParserTests { + @ClassRule + public static KafkaRunning kafkaRunning = KafkaRunning.isRunning(); + @Autowired private ApplicationContext appContext; 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 19487ee656..0eff80addd 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 @@ -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. @@ -16,6 +16,7 @@ package org.springframework.integration.kafka.config.xml; import org.junit.Assert; +import org.junit.ClassRule; import org.junit.Test; import org.junit.runner.RunWith; @@ -23,18 +24,23 @@ import org.springframework.beans.factory.annotation.Autowired; import org.springframework.context.ApplicationContext; 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.test.context.ContextConfiguration; import org.springframework.test.context.junit4.SpringJUnit4ClassRunner; /** * @author Soby Chacko + * @author Gary Russell * @since 0.5 */ @RunWith(SpringJUnit4ClassRunner.class) @ContextConfiguration public class KafkaOutboundAdapterParserTests { + @ClassRule + public static KafkaRunning kafkaRunning = KafkaRunning.isRunning(); + @Autowired private ApplicationContext appContext; diff --git a/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/config/xml/KafkaProducerContextParserTests-context.xml b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/config/xml/KafkaProducerContextParserTests-context.xml index 1c0945a41b..9e133649c4 100644 --- a/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/config/xml/KafkaProducerContextParserTests-context.xml +++ b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/config/xml/KafkaProducerContextParserTests-context.xml @@ -27,7 +27,8 @@ - + { + @ClassRule + public static KafkaRunning kafkaRunning = KafkaRunning.isRunning(); + @Autowired private ApplicationContext appContext; @@ -73,5 +85,8 @@ public class KafkaProducerContextParserTests { final Producer producerTest2 = producerConfigurationTest2.getProducer(); Assert.assertEquals(producerConfigurationTest2, new ProducerConfiguration(producerMetadataTest2, producerTest2)); + + assertFalse(TestUtils.getPropertyValue(producerContext, "autoStartup", Boolean.class)); + assertEquals(123, TestUtils.getPropertyValue(producerContext, "phase")); } } diff --git a/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/config/xml/ZookeeperConnectParserTests.java b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/config/xml/ZookeeperConnectParserTests.java index f520fbb9a2..37adf3bf04 100644 --- a/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/config/xml/ZookeeperConnectParserTests.java +++ b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/config/xml/ZookeeperConnectParserTests.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. @@ -16,23 +16,30 @@ package org.springframework.integration.kafka.config.xml; import org.junit.Assert; +import org.junit.ClassRule; 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.kafka.core.ZookeeperConnectDefaults; +import org.springframework.integration.kafka.rule.KafkaRunning; import org.springframework.integration.kafka.support.ZookeeperConnect; import org.springframework.test.context.ContextConfiguration; import org.springframework.test.context.junit4.SpringJUnit4ClassRunner; /** * @author Soby Chacko + * @author Gary Russell * @since 0.5 */ @RunWith(SpringJUnit4ClassRunner.class) @ContextConfiguration public class ZookeeperConnectParserTests { + @ClassRule + public static KafkaRunning kafkaRunning = KafkaRunning.isRunning(); + @Autowired private ApplicationContext appContext; 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 new file mode 100644 index 0000000000..240f7c0ab6 --- /dev/null +++ b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/outbound/OutboundTests.java @@ -0,0 +1,122 @@ +/* + * 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.outbound; + +import static org.junit.Assert.assertNotNull; + +import java.util.Collections; +import java.util.List; +import java.util.Map; +import java.util.Properties; + +import org.junit.ClassRule; +import org.junit.Test; + +import org.springframework.integration.kafka.rule.KafkaRunning; +import org.springframework.integration.kafka.serializer.common.StringDecoder; +import org.springframework.integration.kafka.serializer.common.StringEncoder; +import org.springframework.integration.kafka.support.ConsumerConfigFactoryBean; +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.KafkaProducerContext; +import org.springframework.integration.kafka.support.MessageLeftOverTracker; +import org.springframework.integration.kafka.support.ProducerConfiguration; +import org.springframework.integration.kafka.support.ProducerFactoryBean; +import org.springframework.integration.kafka.support.ProducerMetadata; +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 + * + */ +public class OutboundTests { + + private static final String TOPIC = "springintegrationtest"; + + @ClassRule + public static KafkaRunning kafkaRunning = KafkaRunning.isRunning(); + + @Test + public void testAsyncProducerFlushed() throws Exception { + KafkaConsumerContext consumerContext = createConsumer(); + // pre-consume to start the receiver because the high-level API doesn't support --from-beginning + consumerContext.receive(); + + KafkaProducerContext kafkaProducerContext = new KafkaProducerContext(); + ProducerMetadata producerMetadata = new ProducerMetadata(TOPIC); + producerMetadata.setValueClassType(String.class); + producerMetadata.setKeyClassType(String.class); + Encoder encoder = new StringEncoder(); + producerMetadata.setValueEncoder(encoder); + producerMetadata.setKeyEncoder(encoder); + producerMetadata.setAsync(true); + Properties props = new Properties(); + props.put("queue.buffering.max.ms", "15000"); + ProducerFactoryBean producer = + new ProducerFactoryBean(producerMetadata, "localhost:9092", props); + ProducerConfiguration config = + new ProducerConfiguration(producerMetadata, producer.getObject()); + kafkaProducerContext.setProducerConfigurations(Collections.singletonMap(TOPIC, config)); + KafkaProducerMessageHandler handler = new KafkaProducerMessageHandler(kafkaProducerContext); + + handler.handleMessage(MessageBuilder.withPayload("foo") + .setHeader("messagekey", "3") + .setHeader("topic", TOPIC) + .build()); + + kafkaProducerContext.stop(); + + Message>>> received = consumerContext.receive(); + assertNotNull(received); + consumerContext.destroy(); + } + + private KafkaConsumerContext createConsumer() throws Exception { + KafkaConsumerContext consumerContext = new KafkaConsumerContext(); + ZookeeperConnect zookeeperConnect = new ZookeeperConnect(); + zookeeperConnect.setZkConnect("localhost:2181"); + consumerContext.setZookeeperConnect(zookeeperConnect); + ConsumerMetadata consumerMetadata = new ConsumerMetadata(); + consumerMetadata.setGroupId("foo"); + Decoder decoder = new StringDecoder(); + consumerMetadata.setValueDecoder(decoder); + consumerMetadata.setKeyDecoder(decoder); + consumerMetadata.setTopicStreamMap(Collections.singletonMap(TOPIC, 1)); + Properties consumerProps = new Properties(); + consumerProps.put("consumer.timeout.ms", "500"); + ConsumerConfigFactoryBean consumerConfigFactoryBean = new ConsumerConfigFactoryBean( + consumerMetadata, zookeeperConnect, consumerProps); + ConsumerConfig consumerConfig = consumerConfigFactoryBean.getObject(); + ConsumerConnectionProvider consumerConnectionProvider = new ConsumerConnectionProvider(consumerConfig); + MessageLeftOverTracker messageLeftOverTracker = new MessageLeftOverTracker(); + ConsumerConfiguration cConfig = new ConsumerConfiguration(consumerMetadata, + consumerConnectionProvider, messageLeftOverTracker); + cConfig.setMaxMessages(1); + consumerContext.setConsumerConfigurations(Collections.singletonMap("foo", cConfig)); + return consumerContext; + } + +} 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 new file mode 100644 index 0000000000..9487148c67 --- /dev/null +++ b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/rule/KafkaRunning.java @@ -0,0 +1,95 @@ +/* + * 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.rule; + +import java.io.IOException; +import java.net.Socket; + +import javax.net.SocketFactory; + +import org.apache.commons.logging.Log; +import org.apache.commons.logging.LogFactory; +import org.junit.Assume; +import org.junit.rules.TestWatcher; +import org.junit.runner.Description; +import org.junit.runners.model.Statement; + +/** + *

+ * A rule that prevents integration tests from failing if the Kafka server is not running or not + * accessible. If the Kafka server is not running in the background all the tests here will simply be skipped because + * of a violated assumption (showing as successful). + *

+ * The rule can be declared as static so that it only has to check once for all tests in the enclosing test case, but + * there isn't a lot of overhead in making it non-static. + * + * @author Dave Syer + * @author Artem Bilan + * @author Gary Russell + * + * @since 1.0 + */ +public class KafkaRunning extends TestWatcher { + + public static final int KAFKA_PORT = 9092; + + public static final int ZOOKEEPER_PORT = 2181; + + private static final Log logger = LogFactory.getLog(KafkaRunning.class); + + /** + * @return a new rule that assumes an existing running broker + */ + public static KafkaRunning isRunning() { + return new KafkaRunning(); + } + + @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(); + } + catch (final 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); + } + +}