From 896d2c43f45b6b63eec426f88ceb66cb68fbb96e Mon Sep 17 00:00:00 2001 From: Ilayaperumal Gopinathan Date: Wed, 12 Nov 2014 11:20:04 +0200 Subject: [PATCH] INTEXT-111: Allow PP for Kafka producer context JIRA: https://jira.spring.io/browse/INTEXT-111 - Remove hard coded names for bean names at the producer context parser - Fix and update tests - Remove the use of `BeanFactory` to get the producerConfigurations in KakfaProducerContext and use setter to set the producerConfigurations - Add logic to send the message from producer configuration when there is no header specified for the topic and single producer configuration is used. Add tests to verify placeholder values Fix code review comment and formatting Polishing code style. Add `KafkaProducerContext#theProducerConfiguration` for single `producerConfigurations` entry to avoid iterators --- .../xml/KafkaProducerContextParser.java | 81 +++++++++++-------- .../kafka/support/KafkaProducerContext.java | 68 ++++++++++------ .../kafka/support/ProducerConfiguration.java | 71 +++++++++------- .../xml/spring-integration-kafka-1.0.xsd | 4 +- .../xml/KafkaOutboundAdapterParserTests.java | 2 +- ...afkaProducerContextParserTests-context.xml | 28 +++++-- .../xml/KafkaProducerContextParserTests.java | 12 +-- .../support/KafkaProducerContextTests.java | 6 +- 8 files changed, 167 insertions(+), 105 deletions(-) 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"));