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 b3eb7a64aa..cfb4c3c650 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 @@ -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,12 +13,17 @@ * See the License for the specific language governing permissions and * limitations under the License. */ + package org.springframework.integration.kafka.config.xml; -import org.springframework.beans.factory.config.BeanDefinition; -import org.springframework.beans.factory.config.BeanDefinitionHolder; +import java.util.Map; + +import org.w3c.dom.Element; + +import org.springframework.beans.BeanMetadataElement; 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.AbstractSimpleBeanDefinitionParser; import org.springframework.beans.factory.xml.ParserContext; import org.springframework.integration.config.xml.IntegrationNamespaceUtils; @@ -28,10 +33,10 @@ import org.springframework.integration.kafka.support.ProducerFactoryBean; import org.springframework.integration.kafka.support.ProducerMetadata; import org.springframework.util.StringUtils; import org.springframework.util.xml.DomUtils; -import org.w3c.dom.Element; /** * @author Soby Chacko + * @author Ilayaperumal Gopinathan * @since 0.5 */ public class KafkaProducerContextParser extends AbstractSimpleBeanDefinitionParser { @@ -49,30 +54,40 @@ public class KafkaProducerContextParser extends AbstractSimpleBeanDefinitionPars parseProducerConfigurations(topics, parserContext, builder, element); } - private void parseProducerConfigurations(final Element topics, final ParserContext parserContext, - final BeanDefinitionBuilder builder, final Element parentElem) { - for (final Element producerConfiguration : DomUtils.getChildElementsByTagName(topics, "producer-configuration")){ - final BeanDefinitionBuilder producerConfigurationBuilder = BeanDefinitionBuilder.genericBeanDefinition(ProducerConfiguration.class); + private void parseProducerConfigurations(Element topics, ParserContext parserContext, + BeanDefinitionBuilder builder, Element parentElem) { + Map producerConfigurationsMap = new ManagedMap(); - final BeanDefinitionBuilder producerMetadataBuilder = BeanDefinitionBuilder.genericBeanDefinition(ProducerMetadata.class); + for (Element producerConfiguration : DomUtils.getChildElementsByTagName(topics, "producer-configuration")) { + BeanDefinitionBuilder producerConfigurationBuilder = + BeanDefinitionBuilder.genericBeanDefinition(ProducerConfiguration.class); + + BeanDefinitionBuilder producerMetadataBuilder = + BeanDefinitionBuilder.genericBeanDefinition(ProducerMetadata.class); producerMetadataBuilder.addConstructorArgValue(producerConfiguration.getAttribute("topic")); - IntegrationNamespaceUtils.setReferenceIfAttributeDefined(producerMetadataBuilder, producerConfiguration, "value-encoder"); - IntegrationNamespaceUtils.setReferenceIfAttributeDefined(producerMetadataBuilder, producerConfiguration, "key-encoder"); - IntegrationNamespaceUtils.setValueIfAttributeDefined(producerMetadataBuilder, producerConfiguration, "key-class-type"); - IntegrationNamespaceUtils.setValueIfAttributeDefined(producerMetadataBuilder, producerConfiguration, "value-class-type"); - IntegrationNamespaceUtils.setReferenceIfAttributeDefined(producerMetadataBuilder, producerConfiguration, "partitioner"); - IntegrationNamespaceUtils.setValueIfAttributeDefined(producerMetadataBuilder, producerConfiguration, "compression-codec"); - IntegrationNamespaceUtils.setValueIfAttributeDefined(producerMetadataBuilder, producerConfiguration, "async"); - IntegrationNamespaceUtils.setValueIfAttributeDefined(producerMetadataBuilder, producerConfiguration, "batch-num-messages"); + IntegrationNamespaceUtils.setReferenceIfAttributeDefined(producerMetadataBuilder, producerConfiguration, + "value-encoder"); + IntegrationNamespaceUtils.setReferenceIfAttributeDefined(producerMetadataBuilder, producerConfiguration, + "key-encoder"); + IntegrationNamespaceUtils.setValueIfAttributeDefined(producerMetadataBuilder, producerConfiguration, + "key-class-type"); + IntegrationNamespaceUtils.setValueIfAttributeDefined(producerMetadataBuilder, producerConfiguration, + "value-class-type"); + IntegrationNamespaceUtils.setReferenceIfAttributeDefined(producerMetadataBuilder, producerConfiguration, + "partitioner"); + IntegrationNamespaceUtils.setValueIfAttributeDefined(producerMetadataBuilder, producerConfiguration, + "compression-codec"); + IntegrationNamespaceUtils.setValueIfAttributeDefined(producerMetadataBuilder, producerConfiguration, + "async"); + IntegrationNamespaceUtils.setValueIfAttributeDefined(producerMetadataBuilder, producerConfiguration, + "batch-num-messages"); + AbstractBeanDefinition producerMetadataBeanDefinition = producerMetadataBuilder.getBeanDefinition(); - final BeanDefinition producerMetadataBeanDef = producerMetadataBuilder.getBeanDefinition(); - registerBeanDefinition(new BeanDefinitionHolder(producerMetadataBeanDef, "producerMetadata_" + producerConfiguration.getAttribute("topic")), - parserContext.getRegistry()); + String producerPropertiesBean = parentElem.getAttribute("producer-properties"); - final String producerPropertiesBean = parentElem.getAttribute("producer-properties"); - - final BeanDefinitionBuilder producerFactoryBuilder = BeanDefinitionBuilder.genericBeanDefinition(ProducerFactoryBean.class); - producerFactoryBuilder.addConstructorArgReference("producerMetadata_" + producerConfiguration.getAttribute("topic")); + BeanDefinitionBuilder producerFactoryBuilder = + BeanDefinitionBuilder.genericBeanDefinition(ProducerFactoryBean.class); + producerFactoryBuilder.addConstructorArgValue(producerMetadataBeanDefinition); final String brokerList = producerConfiguration.getAttribute("broker-list"); if (StringUtils.hasText(brokerList)) { @@ -83,16 +98,18 @@ public class KafkaProducerContextParser extends AbstractSimpleBeanDefinitionPars producerFactoryBuilder.addConstructorArgReference(producerPropertiesBean); } - final BeanDefinition producerfactoryBeanDefinition = producerFactoryBuilder.getBeanDefinition(); - registerBeanDefinition(new BeanDefinitionHolder(producerfactoryBeanDefinition, "prodFactory_" + producerConfiguration.getAttribute("topic")), parserContext.getRegistry()); + AbstractBeanDefinition producerFactoryBeanDefinition = producerFactoryBuilder.getBeanDefinition(); - producerConfigurationBuilder.addConstructorArgReference("producerMetadata_" + producerConfiguration.getAttribute("topic")); - producerConfigurationBuilder.addConstructorArgReference("prodFactory_" + producerConfiguration.getAttribute("topic")); + producerConfigurationBuilder.addConstructorArgValue(producerMetadataBeanDefinition); + producerConfigurationBuilder.addConstructorArgValue(producerFactoryBeanDefinition); - final AbstractBeanDefinition producerConfigurationBeanDefinition = producerConfigurationBuilder.getBeanDefinition(); - final String producerConfigurationBeanName = "producerConfiguration_" + producerConfiguration.getAttribute("topic"); - registerBeanDefinition(new BeanDefinitionHolder(producerConfigurationBeanDefinition, producerConfigurationBeanName), - parserContext.getRegistry()); + AbstractBeanDefinition producerConfigurationBeanDefinition = + producerConfigurationBuilder.getBeanDefinition(); + producerConfigurationsMap.put(producerConfiguration.getAttribute("topic"), + producerConfigurationBeanDefinition); } + + builder.addPropertyValue("producerConfigurations", producerConfigurationsMap); } + } 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 8151e7d22b..87b5678b6a 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-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,60 +13,79 @@ * 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.Map; import java.util.Properties; -import org.apache.commons.lang.StringUtils; import org.apache.commons.logging.Log; import org.apache.commons.logging.LogFactory; -import org.springframework.beans.BeansException; -import org.springframework.beans.factory.*; + import org.springframework.messaging.Message; /** * @author Soby Chacko * @author Rajasekar Elango + * @author Ilayaperumal Gopinathan * @since 0.5 */ -public class KafkaProducerContext implements BeanFactoryAware { +public class KafkaProducerContext { + private static final Log LOGGER = LogFactory.getLog(KafkaProducerContext.class); - private Map> topicsConfiguration; + + private volatile Map> producerConfigurations; + + private volatile ProducerConfiguration theProducerConfiguration; + private Properties producerProperties; public void send(final Message message) throws Exception { - final ProducerConfiguration producerConfiguration = - getTopicConfiguration(message.getHeaders().get("topic", String.class)); - - if (producerConfiguration != null) { - producerConfiguration.send(message); + 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."); } } public ProducerConfiguration getTopicConfiguration(final String topic) { - final Collection> topics = topicsConfiguration.values(); + if (this.theProducerConfiguration != null) { + if (topic.matches(this.theProducerConfiguration.getProducerMetadata().getTopic())) { + return this.theProducerConfiguration; + } + } - for (final ProducerConfiguration producerConfiguration : topics){ - if (topic.matches(producerConfiguration.getProducerMetadata().getTopic())){ + Collection> topics = this.producerConfigurations.values(); + + for (final ProducerConfiguration producerConfiguration : topics) { + if (topic.matches(producerConfiguration.getProducerMetadata().getTopic())) { 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; } - public Map> getTopicsConfiguration() { - return topicsConfiguration; + public Map> getProducerConfigurations() { + return this.producerConfigurations; } - @Override - @SuppressWarnings("unchecked") - public void setBeanFactory(final BeanFactory beanFactory) throws BeansException { - topicsConfiguration = - (Map>) (Object) - ((ListableBeanFactory)beanFactory).getBeansOfType(ProducerConfiguration.class); + public void setProducerConfigurations(Map> producerConfigurations) { + this.producerConfigurations = producerConfigurations; + if (this.producerConfigurations.size() == 1) { + this.theProducerConfiguration = this.producerConfigurations.values().iterator().next(); + } } /** @@ -81,6 +100,7 @@ public class KafkaProducerContext implements BeanFactoryAware { * @return Returns the producerProperties. */ public Properties getProducerProperties() { - return producerProperties; + return this.producerProperties; } + } 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 93a4d7cc43..cab8925b60 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,30 +1,45 @@ /* - * Copyright 2002-2013 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. + * 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. + * 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; -import java.io.*; +import java.io.ByteArrayOutputStream; +import java.io.IOException; +import java.io.ObjectOutputStream; + +import org.apache.commons.lang.builder.EqualsBuilder; +import org.apache.commons.lang.builder.HashCodeBuilder; + +import org.springframework.messaging.Message; +import org.springframework.messaging.MessageHandlingException; import kafka.javaapi.producer.Producer; import kafka.producer.KeyedMessage; import kafka.serializer.DefaultEncoder; -import org.apache.commons.lang.builder.EqualsBuilder; -import org.apache.commons.lang.builder.HashCodeBuilder; -import org.springframework.messaging.Message; - /** * @author Soby Chacko * @author Rajasekar Elango + * @author Ilayaperumal Gopinathan * @since 0.5 */ public class ProducerConfiguration { + private final Producer producer; + private final ProducerMetadata producerMetadata; public ProducerConfiguration(final ProducerMetadata producerMetadata, final Producer producer) { @@ -33,42 +48,47 @@ public class ProducerConfiguration { } public ProducerMetadata getProducerMetadata() { - return producerMetadata; + return this.producerMetadata; + } + + public Producer getProducer() { + return this.producer; } public void send(final Message message) throws Exception { final V v = getPayload(message); - String topic = message.getHeaders().get("topic", String.class); + String topic = message.getHeaders().containsKey("topic") + ? message.getHeaders().get("topic", String.class) + : this.producerMetadata.getTopic(); + if (message.getHeaders().containsKey("messageKey")) { - producer.send(new KeyedMessage(topic, getKey(message), v)); + this.producer.send(new KeyedMessage(topic, getKey(message), v)); } else { - producer.send(new KeyedMessage(topic, v)); + this.producer.send(new KeyedMessage(topic, v)); } } @SuppressWarnings("unchecked") private V getPayload(final Message message) throws Exception { - if (producerMetadata.getValueEncoder().getClass().isAssignableFrom(DefaultEncoder.class)) { + if (this.producerMetadata.getValueEncoder() instanceof DefaultEncoder) { return (V) getByteStream(message.getPayload()); } - else if (message.getPayload().getClass().isAssignableFrom(producerMetadata.getValueClassType())) { + else if (producerMetadata.getValueClassType().isAssignableFrom(message.getPayload().getClass())) { return producerMetadata.getValueClassType().cast(message.getPayload()); } - - throw new Exception("Message payload type is not matching with what is configured"); + throw new MessageHandlingException(message, "Message payload type is not matching with what is configured"); } @SuppressWarnings("unchecked") private K getKey(final Message message) throws Exception { final Object key = message.getHeaders().get("messageKey"); - if (producerMetadata.getKeyEncoder().getClass().isAssignableFrom(DefaultEncoder.class)) { + if (this.producerMetadata.getKeyEncoder() instanceof DefaultEncoder) { return (K) getByteStream(key); } - - return message.getHeaders().get("messageKey", producerMetadata.getKeyClassType()); + return message.getHeaders().get("messageKey", this.producerMetadata.getKeyClassType()); } private static boolean isRawByteArray(final Object obj) { @@ -79,11 +99,9 @@ public class ProducerConfiguration { if (isRawByteArray(obj)) { return (byte[]) obj; } - final ByteArrayOutputStream out = new ByteArrayOutputStream(); final ObjectOutputStream os = new ObjectOutputStream(out); os.writeObject(obj); - return out.toByteArray(); } @@ -99,8 +117,7 @@ public class ProducerConfiguration { @Override public String toString() { - StringBuilder builder = new StringBuilder(); - builder.append("ProducerConfiguration [producerMetadata=").append(producerMetadata).append("]"); - return builder.toString(); + return "ProducerConfiguration [producerMetadata=" + this.producerMetadata + "]"; } + } 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 1e8efe8143..66c65d1196 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 @@ -135,7 +135,7 @@ - Custom implemenation of a Kafka Encoder for encoding message values. + Custom implementation of a Kafka Encoder for encoding message values. @@ -143,7 +143,7 @@ type="xsd:string"> - Custom implemenation of a Kafka Encoder for encoding message keys. + Custom implementation of a Kafka Encoder for encoding message keys. 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 6e2f56440f..19487ee656 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 @@ -49,7 +49,7 @@ public class KafkaOutboundAdapterParserTests { Assert.assertEquals(messageHandler.getOrder(), 3); final KafkaProducerContext producerContext = messageHandler.getKafkaProducerContext(); Assert.assertNotNull(producerContext); - Assert.assertEquals(producerContext.getTopicsConfiguration().size(), 2); + Assert.assertEquals(producerContext.getProducerConfigurations().size(), 2); } } 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 ee83d5c206..1c0945a41b 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 @@ -1,9 +1,12 @@ + xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance" + xmlns:int-kafka="http://www.springframework.org/schema/integration/kafka" + xmlns:context="http://www.springframework.org/schema/context" + xmlns:util="http://www.springframework.org/schema/util" + 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/beans http://www.springframework.org/schema/beans/spring-beans.xsd http://www.springframework.org/schema/context http://www.springframework.org/schema/context/spring-context.xsd http://www.springframework.org/schema/util http://www.springframework.org/schema/util/spring-util.xsd"> + @@ -15,17 +18,26 @@ + + localhost:9092 + localhost:9091 + test1 + test2 + + + + - - diff --git a/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/config/xml/KafkaProducerContextParserTests.java b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/config/xml/KafkaProducerContextParserTests.java index fa9241bbac..54792be04f 100644 --- a/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/config/xml/KafkaProducerContextParserTests.java +++ b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/config/xml/KafkaProducerContextParserTests.java @@ -47,10 +47,10 @@ public class KafkaProducerContextParserTests { final KafkaProducerContext producerContext = appContext.getBean("producerContext", KafkaProducerContext.class); Assert.assertNotNull(producerContext); - final Map> topicConfigurations = producerContext.getTopicsConfiguration(); - Assert.assertEquals(topicConfigurations.size(), 2); + final Map> producerConfigurations = producerContext.getProducerConfigurations(); + Assert.assertEquals(producerConfigurations.size(), 2); - final ProducerConfiguration producerConfigurationTest1 = topicConfigurations.get("producerConfiguration_test1"); + final ProducerConfiguration producerConfigurationTest1 = producerConfigurations.get("test1"); Assert.assertNotNull(producerConfigurationTest1); final ProducerMetadata producerMetadataTest1 = producerConfigurationTest1.getProducerMetadata(); Assert.assertEquals(producerMetadataTest1.getTopic(), "test1"); @@ -62,16 +62,16 @@ public class KafkaProducerContextParserTests { Assert.assertEquals(producerMetadataTest1.getValueEncoder(), valueEncoder); Assert.assertEquals(producerMetadataTest1.getKeyEncoder(), valueEncoder); - final Producer producerTest1 = appContext.getBean("prodFactory_test1", Producer.class); + final Producer producerTest1 = producerConfigurationTest1.getProducer(); Assert.assertEquals(producerConfigurationTest1, new ProducerConfiguration(producerMetadataTest1, producerTest1)); - final ProducerConfiguration producerConfigurationTest2 = topicConfigurations.get("producerConfiguration_" + "test2"); + final ProducerConfiguration producerConfigurationTest2 = producerConfigurations.get("test2"); Assert.assertNotNull(producerConfigurationTest2); final ProducerMetadata producerMetadataTest2 = producerConfigurationTest2.getProducerMetadata(); Assert.assertEquals(producerMetadataTest2.getTopic(), "test2"); Assert.assertEquals(producerMetadataTest2.getCompressionCodec(), "0"); - final Producer producerTest2 = appContext.getBean("prodFactory_test2", Producer.class); + final Producer producerTest2 = producerConfigurationTest2.getProducer(); Assert.assertEquals(producerConfigurationTest2, new ProducerConfiguration(producerMetadataTest2, producerTest2)); } } diff --git a/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/support/KafkaProducerContextTests.java b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/support/KafkaProducerContextTests.java index c1bc1ed32a..5e8e39316d 100644 --- a/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/support/KafkaProducerContextTests.java +++ b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/support/KafkaProducerContextTests.java @@ -37,7 +37,6 @@ public class KafkaProducerContextTests { public void testTopicRegexForProducerConfiguration(){ final KafkaProducerContext kafkaProducerContext = new KafkaProducerContext(); - final ListableBeanFactory beanFactory = Mockito.mock(ListableBeanFactory.class); final ProducerMetadata producerMetadata = Mockito.mock(ProducerMetadata.class); @@ -48,12 +47,9 @@ public class KafkaProducerContextTests { final ProducerConfiguration producerConfiguration = new ProducerConfiguration(producerMetadata, producer); - final Map topicConfigurations = new HashMap(); topicConfigurations.put(testRegex, producerConfiguration); - - Mockito.when(beanFactory.getBeansOfType(ProducerConfiguration.class)).thenReturn(topicConfigurations); - kafkaProducerContext.setBeanFactory(beanFactory); + kafkaProducerContext.setProducerConfigurations(topicConfigurations); Assert.assertNotNull(kafkaProducerContext.getTopicConfiguration("test1")); Assert.assertNotNull(kafkaProducerContext.getTopicConfiguration("test2"));