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
This commit is contained in:
committed by
Artem Bilan
parent
5f4089db0a
commit
325ee9a4f3
@@ -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<String, BeanMetadataElement> consumerConfigurationsMap = new ManagedMap<String, BeanMetadataElement>();
|
||||
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<String, Integer> topicStreamsMap = new HashMap<String, Integer>();
|
||||
final Map<String, String> topicStreamsMap = new HashMap<String, String>();
|
||||
|
||||
final List<Element> 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);
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -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<K,V> implements BeanFactoryAware, DisposableBean {
|
||||
private Map<String, ConsumerConfiguration<K,V>> consumerConfigurations;
|
||||
public class KafkaConsumerContext<K, V> implements DisposableBean {
|
||||
private Map<String, ConsumerConfiguration<K, V>> consumerConfigurations;
|
||||
|
||||
private String consumerTimeout = KafkaConsumerDefaults.CONSUMER_TIMEOUT;
|
||||
|
||||
private ZookeeperConnect zookeeperConnect;
|
||||
|
||||
public Collection<ConsumerConfiguration<K,V>> getConsumerConfigurations() {
|
||||
return consumerConfigurations.values();
|
||||
}
|
||||
|
||||
@Override
|
||||
@SuppressWarnings("unchecked")
|
||||
public void setBeanFactory(final BeanFactory beanFactory) throws BeansException {
|
||||
consumerConfigurations = (Map<String, ConsumerConfiguration<K,V>>)
|
||||
(Object) ((ListableBeanFactory) beanFactory).getBeansOfType(ConsumerConfiguration.class);
|
||||
}
|
||||
|
||||
public Message<Map<String, Map<Integer, List<Object>>>> receive() {
|
||||
final Map<String, Map<Integer, List<Object>>> consumedData = new HashMap<String, Map<Integer, List<Object>>>();
|
||||
|
||||
for (final ConsumerConfiguration<K,V> consumerConfiguration : getConsumerConfigurations()) {
|
||||
final Map<String, Map<Integer, List<Object>>> 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<K,V> implements BeanFactoryAware, DisposableB
|
||||
}
|
||||
|
||||
public ZookeeperConnect getZookeeperConnect() {
|
||||
return zookeeperConnect;
|
||||
return this.zookeeperConnect;
|
||||
}
|
||||
|
||||
public void setZookeeperConnect(final ZookeeperConnect zookeeperConnect) {
|
||||
this.zookeeperConnect = zookeeperConnect;
|
||||
}
|
||||
|
||||
public void setConsumerConfigurations(Map<String, ConsumerConfiguration<K, V>> consumerConfigurations) {
|
||||
this.consumerConfigurations = consumerConfigurations;
|
||||
}
|
||||
|
||||
public Map<String, ConsumerConfiguration<K, V>> getConsumerConfigurations() {
|
||||
return this.consumerConfigurations;
|
||||
}
|
||||
|
||||
public ConsumerConfiguration<K, V> getConsumerConfiguration(String groupId) {
|
||||
return this.consumerConfigurations.get(groupId);
|
||||
}
|
||||
|
||||
public Message<Map<String, Map<Integer, List<Object>>>> receive() {
|
||||
final Map<String, Map<Integer, List<Object>>> consumedData = new HashMap<String, Map<Integer, List<Object>>>();
|
||||
|
||||
for (final ConsumerConfiguration<K, V> consumerConfiguration : getConsumerConfigurations().values()) {
|
||||
final Map<String, Map<Integer, List<Object>>> 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<K,V> config: consumerConfigurations.values()) {
|
||||
for (ConsumerConfiguration<K, V> config : this.consumerConfigurations.values()) {
|
||||
config.getConsumerConnector().shutdown();
|
||||
}
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -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<K, V> {
|
||||
@Test
|
||||
@SuppressWarnings("unchecked")
|
||||
public void testConsumerContextConfiguration() {
|
||||
final KafkaConsumerContext<K, V> consumerContext =
|
||||
appContext.getBean("consumerContext", KafkaConsumerContext.class);
|
||||
assertNotNull(consumerContext);
|
||||
|
||||
ConsumerMetadata<K, V> cm = appContext.getBean("consumerMetadata_default1", ConsumerMetadata.class);
|
||||
final KafkaConsumerContext<K, V> consumerContext = appContext.getBean("consumerContext",
|
||||
KafkaConsumerContext.class);
|
||||
Assert.assertNotNull(consumerContext);
|
||||
ConsumerConfiguration<K, V> cc = consumerContext.getConsumerConfiguration("default1");
|
||||
ConsumerMetadata<K, V> cm = cc.getConsumerMetadata();
|
||||
assertNotNull(cm);
|
||||
TopicFilterConfiguration topicFilterConfiguration = cm.getTopicFilterConfiguration();
|
||||
assertEquals("foo : 10", topicFilterConfiguration.toString());
|
||||
|
||||
@@ -0,0 +1,69 @@
|
||||
<?xml version="1.0" encoding="UTF-8"?>
|
||||
<beans xmlns="http://www.springframework.org/schema/beans"
|
||||
xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
|
||||
xmlns:context="http://www.springframework.org/schema/context"
|
||||
xmlns:int-kafka="http://www.springframework.org/schema/integration/kafka"
|
||||
xsi:schemaLocation="http://www.springframework.org/schema/integration/kafka
|
||||
http://www.springframework.org/schema/integration/kafka/spring-integration-kafka.xsd
|
||||
http://www.springframework.org/schema/context http://www.springframework.org/schema/context/spring-context.xsd
|
||||
http://www.springframework.org/schema/beans http://www.springframework.org/schema/beans/spring-beans.xsd">
|
||||
|
||||
<context:property-placeholder/>
|
||||
|
||||
<int-kafka:zookeeper-connect id="zookeeperConnect" zk-connect="localhost:2181" zk-connection-timeout="6000"
|
||||
zk-session-timeout="6000"
|
||||
zk-sync-time="2000"/>
|
||||
|
||||
<bean id="valueDecoder"
|
||||
class="org.springframework.integration.kafka.serializer.avro.AvroReflectDatumBackedKafkaDecoder">
|
||||
<constructor-arg type="java.lang.Class" value="java.lang.String"/>
|
||||
</bean>
|
||||
|
||||
<bean id="consumerProperties" class="org.springframework.beans.factory.config.PropertiesFactoryBean">
|
||||
<property name="properties">
|
||||
<props>
|
||||
<prop key="auto.offset.reset">largest</prop>
|
||||
<prop key="socket.receive.buffer.bytes">10485760</prop>
|
||||
<!-- 10M -->
|
||||
<prop key="fetch.message.max.bytes">5242880</prop>
|
||||
<prop key="auto.commit.interval.ms">1000</prop>
|
||||
</props>
|
||||
</property>
|
||||
</bean>
|
||||
|
||||
<int-kafka:consumer-context id="consumerContext1"
|
||||
consumer-timeout="4000"
|
||||
zookeeper-connect="zookeeperConnect"
|
||||
consumer-properties="consumerProperties">
|
||||
<int-kafka:consumer-configurations>
|
||||
<int-kafka:consumer-configuration group-id="${groupId:default1}"
|
||||
value-decoder="valueDecoder"
|
||||
key-decoder="valueDecoder"
|
||||
max-messages="5000">
|
||||
<int-kafka:topic id="test1" streams="${streams:3}"/>
|
||||
<int-kafka:topic id="test2" streams="4"/>
|
||||
</int-kafka:consumer-configuration>
|
||||
<int-kafka:consumer-configuration group-id="${groupId:default2}"
|
||||
value-decoder="valueDecoder"
|
||||
key-decoder="valueDecoder"
|
||||
max-messages="5000">
|
||||
<int-kafka:topic id="test3" streams="${streams:1}"/>
|
||||
</int-kafka:consumer-configuration>
|
||||
</int-kafka:consumer-configurations>
|
||||
</int-kafka:consumer-context>
|
||||
<int-kafka:consumer-context id="consumerContext2"
|
||||
consumer-timeout="4000"
|
||||
zookeeper-connect="zookeeperConnect"
|
||||
consumer-properties="consumerProperties">
|
||||
<int-kafka:consumer-configurations>
|
||||
<int-kafka:consumer-configuration group-id="${groupId:default1}"
|
||||
value-decoder="valueDecoder"
|
||||
key-decoder="valueDecoder"
|
||||
max-messages="5000">
|
||||
<int-kafka:topic id="test4" streams="${streams:3}"/>
|
||||
<int-kafka:topic id="test5" streams="4"/>
|
||||
</int-kafka:consumer-configuration>
|
||||
</int-kafka:consumer-configurations>
|
||||
</int-kafka:consumer-context>
|
||||
|
||||
</beans>
|
||||
@@ -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<K,V> {
|
||||
|
||||
@Autowired
|
||||
private ApplicationContext appContext;
|
||||
|
||||
@SuppressWarnings("unchecked")
|
||||
@Test
|
||||
public void testMultiConsumerContexts() {
|
||||
final KafkaConsumerContext<K,V> consumerContext1 = appContext.getBean("consumerContext1", KafkaConsumerContext.class);
|
||||
Assert.assertNotNull(consumerContext1);
|
||||
final KafkaConsumerContext<K,V> consumerContext2 = appContext.getBean("consumerContext2", KafkaConsumerContext.class);
|
||||
Assert.assertNotNull(consumerContext2);
|
||||
}
|
||||
|
||||
@SuppressWarnings("unchecked")
|
||||
@Test
|
||||
public void testConsumerContextConfigurations() {
|
||||
final KafkaConsumerContext<K,V> consumerContext = appContext.getBean("consumerContext1", KafkaConsumerContext.class);
|
||||
Assert.assertNotNull(consumerContext);
|
||||
final ConsumerConfiguration<K,V> cc = consumerContext.getConsumerConfiguration("default1");
|
||||
final ConsumerMetadata<K,V> cm = cc.getConsumerMetadata();
|
||||
Assert.assertTrue(cm.getTopicStreamMap().get("test1") == 3);
|
||||
Assert.assertTrue(cm.getTopicStreamMap().get("test2") == 4);
|
||||
Assert.assertNotNull(cm);
|
||||
final ConsumerConfiguration<K,V> cc2 = consumerContext.getConsumerConfiguration("default2");
|
||||
final ConsumerMetadata<K,V> cm2 = cc2.getConsumerMetadata();
|
||||
Assert.assertTrue(cm2.getTopicStreamMap().get("test3") == 1);
|
||||
Assert.assertNotNull(cm2);
|
||||
final KafkaConsumerContext<K,V> consumerContext2 = appContext.getBean("consumerContext2", KafkaConsumerContext.class);
|
||||
Assert.assertNotNull(consumerContext2);
|
||||
final ConsumerConfiguration<K,V> otherCC = consumerContext2.getConsumerConfiguration("default1");
|
||||
final ConsumerMetadata<K,V> otherCM = otherCC.getConsumerMetadata();
|
||||
Assert.assertTrue(otherCM.getTopicStreamMap().get("test4") == 3);
|
||||
}
|
||||
}
|
||||
@@ -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<K,V> {
|
||||
public class KafkaConsumerContextTest<K, V> {
|
||||
|
||||
@Test
|
||||
@SuppressWarnings("unchecked")
|
||||
public void testMergeResultsFromMultipleConsumerConfiguration() {
|
||||
final KafkaConsumerContext<K,V> kafkaConsumerContext = new KafkaConsumerContext<K,V>();
|
||||
final KafkaConsumerContext<K, V> kafkaConsumerContext = new KafkaConsumerContext<K, V>();
|
||||
final ListableBeanFactory beanFactory = Mockito.mock(ListableBeanFactory.class);
|
||||
final ConsumerConfiguration<K,V> consumerConfiguration1 = Mockito.mock(ConsumerConfiguration.class);
|
||||
final ConsumerConfiguration<K,V> consumerConfiguration2 = Mockito.mock(ConsumerConfiguration.class);
|
||||
final ConsumerConfiguration<K, V> consumerConfiguration1 = Mockito.mock(ConsumerConfiguration.class);
|
||||
final ConsumerConfiguration<K, V> consumerConfiguration2 = Mockito.mock(ConsumerConfiguration.class);
|
||||
|
||||
final Map<String, ConsumerConfiguration<K,V>> map = new HashMap<String, ConsumerConfiguration<K,V>>();
|
||||
final Map<String, ConsumerConfiguration<K, V>> map = new HashMap<String, ConsumerConfiguration<K, V>>();
|
||||
map.put("config1", consumerConfiguration1);
|
||||
map.put("config2", consumerConfiguration2);
|
||||
|
||||
Mockito.when((Map<String, ConsumerConfiguration<K,V>>) (Object) beanFactory.getBeansOfType(ConsumerConfiguration.class)).thenReturn(
|
||||
map);
|
||||
kafkaConsumerContext.setBeanFactory(beanFactory);
|
||||
kafkaConsumerContext.setConsumerConfigurations(map);
|
||||
|
||||
final Map<String, Map<Integer, List<Object>>> result1 = new HashMap<String, Map<Integer, List<Object>>>();
|
||||
final List<Object> l1 = new ArrayList<Object>();
|
||||
@@ -80,4 +80,5 @@ public class KafkaConsumerContextTest<K,V> {
|
||||
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");
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user