From 325ee9a4f3c7bd57ac60a8ef6fea75aeb5fe2a87 Mon Sep 17 00:00:00 2001 From: Ilayaperumal Gopinathan Date: Wed, 20 Aug 2014 21:10:00 +0300 Subject: [PATCH] Kafka: PP for consumer context bean JIRA: https://jira.spring.io/browse/INTEXT-112 https://jira.spring.io/browse/INTEXT-111 - Avoid using hardcoded bean names for consumer context - Remove Integer value topic streams count, instead use string value - This will allow placeholder values to be used for `streams` attribute - Set `group-id` as a bean name for consumer configuration bean only if the group id is specified explicitly. This will allow placeholder values not being used as a bean name. - Add test cases to verify multiple consumer configurations within the same consumer context. Set consumerConfigurations as propertyValue Better handling of consumerConfigurations - Add getConsumerConfiguration(String groupId) in KafkaConsumerContext - Fix tests - Add additional test case to check multi consumer contexts Fix getConsumerConfiguration(groupId) Polishing --- .../xml/KafkaConsumerContextParser.java | 88 +++++++++++-------- .../kafka/support/KafkaConsumerContext.java | 70 +++++++-------- .../xml/KafkaConsumerContextParserTests.java | 17 ++-- ...ultiConsumerContextParserTests-context.xml | 69 +++++++++++++++ .../KafkaMultiConsumerContextParserTests.java | 68 ++++++++++++++ .../support/KafkaConsumerContextTest.java | 31 +++---- 6 files changed, 249 insertions(+), 94 deletions(-) create mode 100644 spring-integration-kafka/src/test/java/org/springframework/integration/kafka/config/xml/KafkaMultiConsumerContextParserTests-context.xml create mode 100644 spring-integration-kafka/src/test/java/org/springframework/integration/kafka/config/xml/KafkaMultiConsumerContextParserTests.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 dab70c9..7604d0d 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 @@ -22,10 +22,11 @@ import java.util.Map; import org.w3c.dom.Element; +import org.springframework.beans.BeanMetadataElement; import org.springframework.beans.factory.config.BeanDefinition; -import org.springframework.beans.factory.config.BeanDefinitionHolder; import org.springframework.beans.factory.support.AbstractBeanDefinition; import org.springframework.beans.factory.support.BeanDefinitionBuilder; +import org.springframework.beans.factory.support.ManagedMap; import org.springframework.beans.factory.xml.AbstractSingleBeanDefinitionParser; import org.springframework.beans.factory.xml.ParserContext; import org.springframework.integration.config.xml.IntegrationNamespaceUtils; @@ -43,6 +44,7 @@ import org.springframework.util.xml.DomUtils; * @author Soby Chacko * @author Rajasekar Elango * @author Artem Bilan + * @author Ilayaperumal Gopinathan * @since 0.5 */ public class KafkaConsumerContextParser extends AbstractSingleBeanDefinitionParser { @@ -62,20 +64,30 @@ public class KafkaConsumerContextParser extends AbstractSingleBeanDefinitionPars private void parseConsumerConfigurations(final Element consumerConfigurations, final ParserContext parserContext, final BeanDefinitionBuilder builder, final Element parentElem) { + Map consumerConfigurationsMap = new ManagedMap(); for (final Element consumerConfiguration : DomUtils.getChildElementsByTagName(consumerConfigurations, "consumer-configuration")) { - final BeanDefinitionBuilder consumerConfigurationBuilder = BeanDefinitionBuilder.genericBeanDefinition(ConsumerConfiguration.class); - final BeanDefinitionBuilder consumerMetadataBuilder = BeanDefinitionBuilder.genericBeanDefinition(ConsumerMetadata.class); + final BeanDefinitionBuilder consumerConfigurationBuilder = + BeanDefinitionBuilder.genericBeanDefinition(ConsumerConfiguration.class); + final BeanDefinitionBuilder consumerMetadataBuilder = + BeanDefinitionBuilder.genericBeanDefinition(ConsumerMetadata.class); - IntegrationNamespaceUtils.setValueIfAttributeDefined(consumerMetadataBuilder, consumerConfiguration, "group-id"); + IntegrationNamespaceUtils.setValueIfAttributeDefined(consumerMetadataBuilder, consumerConfiguration, + "group-id"); - IntegrationNamespaceUtils.setReferenceIfAttributeDefined(consumerMetadataBuilder, consumerConfiguration, "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, "consumer-timeout"); + IntegrationNamespaceUtils.setReferenceIfAttributeDefined(consumerMetadataBuilder, consumerConfiguration, + "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, + "consumer-timeout"); - final Map topicStreamsMap = new HashMap(); + final Map topicStreamsMap = new HashMap(); final List topicConfigurations = DomUtils.getChildElementsByTagName(consumerConfiguration, "topic"); @@ -83,8 +95,7 @@ public class KafkaConsumerContextParser extends AbstractSingleBeanDefinitionPars for (final Element topicConfiguration : topicConfigurations) { final String topic = topicConfiguration.getAttribute("id"); final String streams = topicConfiguration.getAttribute("streams"); - final Integer streamsInt = Integer.valueOf(streams); - topicStreamsMap.put(topic, streamsInt); + topicStreamsMap.put(topic, streams); } consumerMetadataBuilder.addPropertyValue("topicStreamMap", topicStreamsMap); } @@ -98,20 +109,20 @@ public class KafkaConsumerContextParser extends AbstractSingleBeanDefinitionPars .addConstructorArgValue(topicFilter.getAttribute("streams")) .addConstructorArgValue(topicFilter.getAttribute("exclude")) .getBeanDefinition(); - consumerMetadataBuilder.addPropertyValue("topicFilterConfiguration", topicFilterConfigurationBeanDefinition); + consumerMetadataBuilder.addPropertyValue("topicFilterConfiguration", + topicFilterConfigurationBeanDefinition); } - final BeanDefinition consumerMetadataBeanDef = consumerMetadataBuilder.getBeanDefinition(); - registerBeanDefinition(new BeanDefinitionHolder(consumerMetadataBeanDef, "consumerMetadata_" + consumerConfiguration.getAttribute("group-id")), - parserContext.getRegistry()); + final AbstractBeanDefinition consumerMetadataBeanDefintiion = consumerMetadataBuilder.getBeanDefinition(); final String zookeeperConnectBean = parentElem.getAttribute("zookeeper-connect"); IntegrationNamespaceUtils.setReferenceIfAttributeDefined(builder, parentElem, zookeeperConnectBean); final String consumerPropertiesBean = parentElem.getAttribute("consumer-properties"); - final BeanDefinitionBuilder consumerConfigFactoryBuilder = BeanDefinitionBuilder.genericBeanDefinition(ConsumerConfigFactoryBean.class); - consumerConfigFactoryBuilder.addConstructorArgReference("consumerMetadata_" + consumerConfiguration.getAttribute("group-id")); + final BeanDefinitionBuilder consumerConfigFactoryBuilder = + BeanDefinitionBuilder.genericBeanDefinition(ConsumerConfigFactoryBean.class); + consumerConfigFactoryBuilder.addConstructorArgValue(consumerMetadataBeanDefintiion); if (StringUtils.hasText(zookeeperConnectBean)) { consumerConfigFactoryBuilder.addConstructorArgReference(zookeeperConnectBean); @@ -121,30 +132,31 @@ public class KafkaConsumerContextParser extends AbstractSingleBeanDefinitionPars consumerConfigFactoryBuilder.addConstructorArgReference(consumerPropertiesBean); } - final BeanDefinition consumerConfigFactoryBuilderBeanDefinition = consumerConfigFactoryBuilder.getBeanDefinition(); - registerBeanDefinition(new BeanDefinitionHolder(consumerConfigFactoryBuilderBeanDefinition, "consumerConfigFactory_" + consumerConfiguration.getAttribute("group-id")), parserContext.getRegistry()); + AbstractBeanDefinition consumerConfigFactoryBuilderBeanDefinition = + consumerConfigFactoryBuilder.getBeanDefinition(); - final BeanDefinitionBuilder consumerConnectionProviderBuilder = BeanDefinitionBuilder.genericBeanDefinition(ConsumerConnectionProvider.class); - consumerConnectionProviderBuilder.addConstructorArgReference("consumerConfigFactory_" + consumerConfiguration.getAttribute("group-id")); + BeanDefinitionBuilder consumerConnectionProviderBuilder = + BeanDefinitionBuilder.genericBeanDefinition(ConsumerConnectionProvider.class); + consumerConnectionProviderBuilder.addConstructorArgValue(consumerConfigFactoryBuilderBeanDefinition); - final BeanDefinition consumerConnectionProviderBuilderBeanDefinition = consumerConnectionProviderBuilder.getBeanDefinition(); - registerBeanDefinition(new BeanDefinitionHolder(consumerConnectionProviderBuilderBeanDefinition, "consumerConnectionProvider_" + consumerConfiguration.getAttribute("group-id")), parserContext.getRegistry()); + AbstractBeanDefinition consumerConnectionProviderBuilderBeanDefinition = + consumerConnectionProviderBuilder.getBeanDefinition(); + BeanDefinitionBuilder messageLeftOverBeanDefinitionBuilder = + BeanDefinitionBuilder.genericBeanDefinition(MessageLeftOverTracker.class); + AbstractBeanDefinition messageLeftOverBeanDefinition = + messageLeftOverBeanDefinitionBuilder.getBeanDefinition(); - final BeanDefinitionBuilder messageLeftOverBeanDefinitionBuilder = BeanDefinitionBuilder.genericBeanDefinition(MessageLeftOverTracker.class); - final BeanDefinition messageLeftOverBeanDefinition = messageLeftOverBeanDefinitionBuilder.getBeanDefinition(); - registerBeanDefinition(new BeanDefinitionHolder(messageLeftOverBeanDefinition, "messageLeftOver_" + consumerConfiguration.getAttribute("group-id")), - parserContext.getRegistry()); + consumerConfigurationBuilder.addConstructorArgValue(consumerMetadataBeanDefintiion); + consumerConfigurationBuilder.addConstructorArgValue(consumerConnectionProviderBuilderBeanDefinition); + consumerConfigurationBuilder.addConstructorArgValue(messageLeftOverBeanDefinition); - consumerConfigurationBuilder.addConstructorArgReference("consumerMetadata_" + consumerConfiguration.getAttribute("group-id")); - consumerConfigurationBuilder.addConstructorArgReference("consumerConnectionProvider_" + consumerConfiguration.getAttribute("group-id")); - consumerConfigurationBuilder.addConstructorArgReference("messageLeftOver_" + consumerConfiguration.getAttribute("group-id")); - - final AbstractBeanDefinition consumerConfigurationBeanDefinition = consumerConfigurationBuilder.getBeanDefinition(); - - final String consumerConfigurationBeanName = "consumerConfiguration_" + consumerConfiguration.getAttribute("group-id"); - registerBeanDefinition(new BeanDefinitionHolder(consumerConfigurationBeanDefinition, consumerConfigurationBeanName), - parserContext.getRegistry()); + AbstractBeanDefinition consumerConfigurationBeanDefinition = + consumerConfigurationBuilder.getBeanDefinition(); + consumerConfigurationsMap.put(consumerConfiguration.getAttribute("group-id"), + consumerConfigurationBeanDefinition); } + builder.addPropertyValue("consumerConfigurations", consumerConfigurationsMap); } + } diff --git a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/support/KafkaConsumerContext.java b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/support/KafkaConsumerContext.java index 3c59055..650717c 100644 --- a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/support/KafkaConsumerContext.java +++ b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/support/KafkaConsumerContext.java @@ -1,5 +1,5 @@ /* - * Copyright 2002-2013 the original author or authors. + * Copyright 2002-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,18 +13,14 @@ * See the License for the specific language governing permissions and * limitations under the License. */ + package org.springframework.integration.kafka.support; -import java.util.Collection; import java.util.HashMap; import java.util.List; import java.util.Map; -import org.springframework.beans.BeansException; -import org.springframework.beans.factory.BeanFactory; -import org.springframework.beans.factory.BeanFactoryAware; import org.springframework.beans.factory.DisposableBean; -import org.springframework.beans.factory.ListableBeanFactory; import org.springframework.integration.kafka.core.KafkaConsumerDefaults; import org.springframework.integration.support.MessageBuilder; import org.springframework.messaging.Message; @@ -35,37 +31,15 @@ import org.springframework.util.CollectionUtils; * @author Ilayaperumal Gopinathan * @since 0.5 */ -public class KafkaConsumerContext implements BeanFactoryAware, DisposableBean { - private Map> consumerConfigurations; +public class KafkaConsumerContext implements DisposableBean { + private Map> consumerConfigurations; + private String consumerTimeout = KafkaConsumerDefaults.CONSUMER_TIMEOUT; + private ZookeeperConnect zookeeperConnect; - public Collection> getConsumerConfigurations() { - return consumerConfigurations.values(); - } - - @Override - @SuppressWarnings("unchecked") - public void setBeanFactory(final BeanFactory beanFactory) throws BeansException { - consumerConfigurations = (Map>) - (Object) ((ListableBeanFactory) beanFactory).getBeansOfType(ConsumerConfiguration.class); - } - - public Message>>> receive() { - final Map>> consumedData = new HashMap>>(); - - for (final ConsumerConfiguration consumerConfiguration : getConsumerConfigurations()) { - final Map>> messages = consumerConfiguration.receive(); - - if (!CollectionUtils.isEmpty(messages)){ - consumedData.putAll(messages); - } - } - return consumedData.isEmpty() ? null : MessageBuilder.withPayload(consumedData).build(); - } - public String getConsumerTimeout() { - return consumerTimeout; + return this.consumerTimeout; } public void setConsumerTimeout(final String consumerTimeout) { @@ -73,17 +47,43 @@ public class KafkaConsumerContext implements BeanFactoryAware, DisposableB } public ZookeeperConnect getZookeeperConnect() { - return zookeeperConnect; + return this.zookeeperConnect; } public void setZookeeperConnect(final ZookeeperConnect zookeeperConnect) { this.zookeeperConnect = zookeeperConnect; } + public void setConsumerConfigurations(Map> consumerConfigurations) { + this.consumerConfigurations = consumerConfigurations; + } + + public Map> getConsumerConfigurations() { + return this.consumerConfigurations; + } + + public ConsumerConfiguration getConsumerConfiguration(String groupId) { + return this.consumerConfigurations.get(groupId); + } + + public Message>>> receive() { + final Map>> consumedData = new HashMap>>(); + + for (final ConsumerConfiguration consumerConfiguration : getConsumerConfigurations().values()) { + final Map>> messages = consumerConfiguration.receive(); + + if (!CollectionUtils.isEmpty(messages)) { + consumedData.putAll(messages); + } + } + return consumedData.isEmpty() ? null : MessageBuilder.withPayload(consumedData).build(); + } + @Override public void destroy() throws Exception { - for (ConsumerConfiguration config: consumerConfigurations.values()) { + for (ConsumerConfiguration config : this.consumerConfigurations.values()) { config.getConsumerConnector().shutdown(); } } + } 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 f53d0b5..5d43b73 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 @@ -13,17 +13,22 @@ * See the License for the specific language governing permissions and * limitations under the License. */ + package org.springframework.integration.kafka.config.xml; -import static org.junit.Assert.*; +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.Test; import org.junit.runner.RunWith; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.context.ApplicationContext; +import org.springframework.integration.kafka.support.ConsumerConfiguration; import org.springframework.integration.kafka.support.ConsumerMetadata; import org.springframework.integration.kafka.support.KafkaConsumerContext; import org.springframework.integration.kafka.support.TopicFilterConfiguration; @@ -45,11 +50,11 @@ public class KafkaConsumerContextParserTests { @Test @SuppressWarnings("unchecked") public void testConsumerContextConfiguration() { - final KafkaConsumerContext consumerContext = - appContext.getBean("consumerContext", KafkaConsumerContext.class); - assertNotNull(consumerContext); - - ConsumerMetadata cm = appContext.getBean("consumerMetadata_default1", ConsumerMetadata.class); + final KafkaConsumerContext consumerContext = appContext.getBean("consumerContext", + KafkaConsumerContext.class); + Assert.assertNotNull(consumerContext); + ConsumerConfiguration cc = consumerContext.getConsumerConfiguration("default1"); + ConsumerMetadata cm = cc.getConsumerMetadata(); assertNotNull(cm); TopicFilterConfiguration topicFilterConfiguration = cm.getTopicFilterConfiguration(); assertEquals("foo : 10", topicFilterConfiguration.toString()); diff --git a/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/config/xml/KafkaMultiConsumerContextParserTests-context.xml b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/config/xml/KafkaMultiConsumerContextParserTests-context.xml new file mode 100644 index 0000000..a428f6f --- /dev/null +++ b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/config/xml/KafkaMultiConsumerContextParserTests-context.xml @@ -0,0 +1,69 @@ + + + + + + + + + + + + + + + largest + 10485760 + + 5242880 + 1000 + + + + + + + + + + + + + + + + + + + + + + + + + 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 new file mode 100644 index 0000000..84f9b95 --- /dev/null +++ b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/config/xml/KafkaMultiConsumerContextParserTests.java @@ -0,0 +1,68 @@ +/* + * 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.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.kafka.support.ConsumerConfiguration; +import org.springframework.integration.kafka.support.ConsumerMetadata; +import org.springframework.integration.kafka.support.KafkaConsumerContext; +import org.springframework.test.context.ContextConfiguration; +import org.springframework.test.context.junit4.SpringJUnit4ClassRunner; + +/** + * @author Ilayaperumal Gopinathan + */ +@RunWith(SpringJUnit4ClassRunner.class) +@ContextConfiguration +public class KafkaMultiConsumerContextParserTests { + + @Autowired + private ApplicationContext appContext; + + @SuppressWarnings("unchecked") + @Test + public void testMultiConsumerContexts() { + final KafkaConsumerContext consumerContext1 = appContext.getBean("consumerContext1", KafkaConsumerContext.class); + Assert.assertNotNull(consumerContext1); + final KafkaConsumerContext consumerContext2 = appContext.getBean("consumerContext2", KafkaConsumerContext.class); + Assert.assertNotNull(consumerContext2); + } + + @SuppressWarnings("unchecked") + @Test + public void testConsumerContextConfigurations() { + final KafkaConsumerContext consumerContext = appContext.getBean("consumerContext1", KafkaConsumerContext.class); + Assert.assertNotNull(consumerContext); + final ConsumerConfiguration cc = consumerContext.getConsumerConfiguration("default1"); + final ConsumerMetadata cm = cc.getConsumerMetadata(); + Assert.assertTrue(cm.getTopicStreamMap().get("test1") == 3); + Assert.assertTrue(cm.getTopicStreamMap().get("test2") == 4); + Assert.assertNotNull(cm); + final ConsumerConfiguration cc2 = consumerContext.getConsumerConfiguration("default2"); + final ConsumerMetadata cm2 = cc2.getConsumerMetadata(); + Assert.assertTrue(cm2.getTopicStreamMap().get("test3") == 1); + Assert.assertNotNull(cm2); + final KafkaConsumerContext consumerContext2 = appContext.getBean("consumerContext2", KafkaConsumerContext.class); + Assert.assertNotNull(consumerContext2); + final ConsumerConfiguration otherCC = consumerContext2.getConsumerConfiguration("default1"); + final ConsumerMetadata otherCM = otherCC.getConsumerMetadata(); + Assert.assertTrue(otherCM.getTopicStreamMap().get("test4") == 3); + } +} diff --git a/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/support/KafkaConsumerContextTest.java b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/support/KafkaConsumerContextTest.java index bb28b92..709ab63 100644 --- a/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/support/KafkaConsumerContextTest.java +++ b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/support/KafkaConsumerContextTest.java @@ -1,5 +1,5 @@ /* - * Copyright 2002-2013 the original author or authors. + * Copyright 2002-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,40 +13,40 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package org.springframework.integration.kafka.support; -import org.junit.Assert; -import org.junit.Test; -import org.mockito.Mockito; -import org.springframework.beans.factory.ListableBeanFactory; -import org.springframework.messaging.Message; +package org.springframework.integration.kafka.support; import java.util.ArrayList; import java.util.HashMap; import java.util.List; import java.util.Map; +import org.junit.Assert; +import org.junit.Test; +import org.mockito.Mockito; + +import org.springframework.beans.factory.ListableBeanFactory; +import org.springframework.messaging.Message; + /** * @author Soby Chacko * @since 0.5 */ -public class KafkaConsumerContextTest { +public class KafkaConsumerContextTest { @Test @SuppressWarnings("unchecked") public void testMergeResultsFromMultipleConsumerConfiguration() { - final KafkaConsumerContext kafkaConsumerContext = new KafkaConsumerContext(); + final KafkaConsumerContext kafkaConsumerContext = new KafkaConsumerContext(); final ListableBeanFactory beanFactory = Mockito.mock(ListableBeanFactory.class); - final ConsumerConfiguration consumerConfiguration1 = Mockito.mock(ConsumerConfiguration.class); - final ConsumerConfiguration consumerConfiguration2 = Mockito.mock(ConsumerConfiguration.class); + final ConsumerConfiguration consumerConfiguration1 = Mockito.mock(ConsumerConfiguration.class); + final ConsumerConfiguration consumerConfiguration2 = Mockito.mock(ConsumerConfiguration.class); - final Map> map = new HashMap>(); + final Map> map = new HashMap>(); map.put("config1", consumerConfiguration1); map.put("config2", consumerConfiguration2); - Mockito.when((Map>) (Object) beanFactory.getBeansOfType(ConsumerConfiguration.class)).thenReturn( - map); - kafkaConsumerContext.setBeanFactory(beanFactory); + kafkaConsumerContext.setConsumerConfigurations(map); final Map>> result1 = new HashMap>>(); final List l1 = new ArrayList(); @@ -80,4 +80,5 @@ public class KafkaConsumerContextTest { Assert.assertEquals(messages.getPayload().get("topic2").get(1).get(1), "got message2 - l2"); Assert.assertEquals(messages.getPayload().get("topic2").get(1).get(2), "got message3 - l2"); } + }