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
This commit is contained in:
Ilayaperumal Gopinathan
2014-11-12 11:20:04 +02:00
committed by Artem Bilan
parent 94585fe03b
commit 896d2c43f4
8 changed files with 167 additions and 105 deletions

View File

@@ -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<String, BeanMetadataElement> producerConfigurationsMap = new ManagedMap<String, BeanMetadataElement>();
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);
}
}

View File

@@ -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<K,V> implements BeanFactoryAware {
public class KafkaProducerContext<K, V> {
private static final Log LOGGER = LogFactory.getLog(KafkaProducerContext.class);
private Map<String, ProducerConfiguration<K,V>> topicsConfiguration;
private volatile Map<String, ProducerConfiguration<K, V>> producerConfigurations;
private volatile ProducerConfiguration<K, V> theProducerConfiguration;
private Properties producerProperties;
public void send(final Message<?> message) throws Exception {
final ProducerConfiguration<K,V> producerConfiguration =
getTopicConfiguration(message.getHeaders().get("topic", String.class));
if (producerConfiguration != null) {
producerConfiguration.send(message);
if (message.getHeaders().containsKey("topic")) {
ProducerConfiguration<K, V> 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<K, V> getTopicConfiguration(final String topic) {
final Collection<ProducerConfiguration<K,V>> topics = topicsConfiguration.values();
if (this.theProducerConfiguration != null) {
if (topic.matches(this.theProducerConfiguration.getProducerMetadata().getTopic())) {
return this.theProducerConfiguration;
}
}
for (final ProducerConfiguration<K,V> producerConfiguration : topics){
if (topic.matches(producerConfiguration.getProducerMetadata().getTopic())){
Collection<ProducerConfiguration<K, V>> topics = this.producerConfigurations.values();
for (final ProducerConfiguration<K, V> 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<String, ProducerConfiguration<K,V>> getTopicsConfiguration() {
return topicsConfiguration;
public Map<String, ProducerConfiguration<K, V>> getProducerConfigurations() {
return this.producerConfigurations;
}
@Override
@SuppressWarnings("unchecked")
public void setBeanFactory(final BeanFactory beanFactory) throws BeansException {
topicsConfiguration =
(Map<String, ProducerConfiguration<K,V>>) (Object)
((ListableBeanFactory)beanFactory).getBeansOfType(ProducerConfiguration.class);
public void setProducerConfigurations(Map<String, ProducerConfiguration<K, V>> producerConfigurations) {
this.producerConfigurations = producerConfigurations;
if (this.producerConfigurations.size() == 1) {
this.theProducerConfiguration = this.producerConfigurations.values().iterator().next();
}
}
/**
@@ -81,6 +100,7 @@ public class KafkaProducerContext<K,V> implements BeanFactoryAware {
* @return Returns the producerProperties.
*/
public Properties getProducerProperties() {
return producerProperties;
return this.producerProperties;
}
}

View File

@@ -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<K, V> {
private final Producer<K, V> producer;
private final ProducerMetadata<K, V> producerMetadata;
public ProducerConfiguration(final ProducerMetadata<K, V> producerMetadata, final Producer<K, V> producer) {
@@ -33,42 +48,47 @@ public class ProducerConfiguration<K, V> {
}
public ProducerMetadata<K, V> getProducerMetadata() {
return producerMetadata;
return this.producerMetadata;
}
public Producer<K, V> 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<K, V>(topic, getKey(message), v));
this.producer.send(new KeyedMessage<K, V>(topic, getKey(message), v));
}
else {
producer.send(new KeyedMessage<K, V>(topic, v));
this.producer.send(new KeyedMessage<K, V>(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<K, V> {
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<K, V> {
@Override
public String toString() {
StringBuilder builder = new StringBuilder();
builder.append("ProducerConfiguration [producerMetadata=").append(producerMetadata).append("]");
return builder.toString();
return "ProducerConfiguration [producerMetadata=" + this.producerMetadata + "]";
}
}

View File

@@ -135,7 +135,7 @@
<xsd:attribute name="value-encoder" use="optional" type="xsd:string">
<xsd:annotation>
<xsd:documentation>
Custom implemenation of a Kafka Encoder for encoding message values.
Custom implementation of a Kafka Encoder for encoding message values.
</xsd:documentation>
</xsd:annotation>
</xsd:attribute>
@@ -143,7 +143,7 @@
type="xsd:string">
<xsd:annotation>
<xsd:documentation>
Custom implemenation of a Kafka Encoder for encoding message keys.
Custom implementation of a Kafka Encoder for encoding message keys.
</xsd:documentation>
</xsd:annotation>
</xsd:attribute>

View File

@@ -49,7 +49,7 @@ public class KafkaOutboundAdapterParserTests<K, V> {
Assert.assertEquals(messageHandler.getOrder(), 3);
final KafkaProducerContext<K, V> producerContext = messageHandler.getKafkaProducerContext();
Assert.assertNotNull(producerContext);
Assert.assertEquals(producerContext.getTopicsConfiguration().size(), 2);
Assert.assertEquals(producerContext.getProducerConfigurations().size(), 2);
}
}

View File

@@ -1,9 +1,12 @@
<?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: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/beans http://www.springframework.org/schema/beans/spring-beans.xsd">
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">
<bean id="producerProperties" class="org.springframework.beans.factory.config.PropertiesFactoryBean">
<property name="properties">
@@ -15,17 +18,26 @@
</property>
</bean>
<util:properties id="placeholderProperties">
<prop key="brokerList1">localhost:9092</prop>
<prop key="brokerList2">localhost:9091</prop>
<prop key="topic1">test1</prop>
<prop key="topic2">test2</prop>
</util:properties>
<context:property-placeholder properties-ref="placeholderProperties"/>
<int-kafka:producer-context id="producerContext" producer-properties="producerProperties">
<int-kafka:producer-configurations>
<int-kafka:producer-configuration broker-list="localhost:9092"
<int-kafka:producer-configuration broker-list="${brokerList1}"
key-class-type="java.lang.String"
value-class-type="java.lang.String"
key-encoder="valueEncoder"
value-encoder="valueEncoder"
topic="test1"
topic="${topic1}"
compression-codec="default"/>
<int-kafka:producer-configuration broker-list="localhost:9092"
topic="test2"
<int-kafka:producer-configuration broker-list="${brokerList2}"
topic="${topic2}"
compression-codec="default"/>
</int-kafka:producer-configurations>
</int-kafka:producer-context>

View File

@@ -47,10 +47,10 @@ public class KafkaProducerContextParserTests<K,V,T> {
final KafkaProducerContext<K,V> producerContext = appContext.getBean("producerContext", KafkaProducerContext.class);
Assert.assertNotNull(producerContext);
final Map<String, ProducerConfiguration<K,V>> topicConfigurations = producerContext.getTopicsConfiguration();
Assert.assertEquals(topicConfigurations.size(), 2);
final Map<String, ProducerConfiguration<K,V>> producerConfigurations = producerContext.getProducerConfigurations();
Assert.assertEquals(producerConfigurations.size(), 2);
final ProducerConfiguration<K,V> producerConfigurationTest1 = topicConfigurations.get("producerConfiguration_test1");
final ProducerConfiguration<K,V> producerConfigurationTest1 = producerConfigurations.get("test1");
Assert.assertNotNull(producerConfigurationTest1);
final ProducerMetadata<K,V> producerMetadataTest1 = producerConfigurationTest1.getProducerMetadata();
Assert.assertEquals(producerMetadataTest1.getTopic(), "test1");
@@ -62,16 +62,16 @@ public class KafkaProducerContextParserTests<K,V,T> {
Assert.assertEquals(producerMetadataTest1.getValueEncoder(), valueEncoder);
Assert.assertEquals(producerMetadataTest1.getKeyEncoder(), valueEncoder);
final Producer<K,V> producerTest1 = appContext.getBean("prodFactory_test1", Producer.class);
final Producer<K,V> producerTest1 = producerConfigurationTest1.getProducer();
Assert.assertEquals(producerConfigurationTest1, new ProducerConfiguration<K,V>(producerMetadataTest1, producerTest1));
final ProducerConfiguration<K,V> producerConfigurationTest2 = topicConfigurations.get("producerConfiguration_" + "test2");
final ProducerConfiguration<K,V> producerConfigurationTest2 = producerConfigurations.get("test2");
Assert.assertNotNull(producerConfigurationTest2);
final ProducerMetadata<K,V> producerMetadataTest2 = producerConfigurationTest2.getProducerMetadata();
Assert.assertEquals(producerMetadataTest2.getTopic(), "test2");
Assert.assertEquals(producerMetadataTest2.getCompressionCodec(), "0");
final Producer<K,V> producerTest2 = appContext.getBean("prodFactory_test2", Producer.class);
final Producer<K,V> producerTest2 = producerConfigurationTest2.getProducer();
Assert.assertEquals(producerConfigurationTest2, new ProducerConfiguration<K,V>(producerMetadataTest2, producerTest2));
}
}

View File

@@ -37,7 +37,6 @@ public class KafkaProducerContextTests {
public void testTopicRegexForProducerConfiguration(){
final KafkaProducerContext kafkaProducerContext = new KafkaProducerContext();
final ListableBeanFactory beanFactory = Mockito.mock(ListableBeanFactory.class);
final ProducerMetadata<String, String> producerMetadata = Mockito.mock(ProducerMetadata.class);
@@ -48,12 +47,9 @@ public class KafkaProducerContextTests {
final ProducerConfiguration<String, String> producerConfiguration = new ProducerConfiguration<String, String>(producerMetadata, producer);
final Map<String, ProducerConfiguration> topicConfigurations = new HashMap<String, ProducerConfiguration>();
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"));