From c52be66a1c0185d3c4115f8accc901f45f677b57 Mon Sep 17 00:00:00 2001 From: Rajasekar Elango Date: Tue, 29 Oct 2013 11:41:24 -0700 Subject: [PATCH] INTEXT-92 Kafka: Add support for specifying any producer/consumer property * Update README.md for producer and consumer properties usage --- .../xml/KafkaConsumerContextParser.java | 18 ++++++------- .../xml/KafkaProducerContextParser.java | 11 ++++++-- .../support/ConsumerConfigFactoryBean.java | 26 ++++++++++++++++--- .../kafka/support/ConsumerConfiguration.java | 2 +- .../kafka/support/KafkaProducerContext.java | 17 ++++++++++++ .../kafka/support/ProducerFactoryBean.java | 19 +++++++++++++- .../xml/spring-integration-kafka-1.0.xsd | 16 ++++++++++++ ...afkaConsumerContextParserTests-context.xml | 14 +++++++++- ...afkaProducerContextParserTests-context.xml | 12 ++++++++- 9 files changed, 116 insertions(+), 19 deletions(-) diff --git a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/config/xml/KafkaConsumerContextParser.java b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/config/xml/KafkaConsumerContextParser.java index c21f7f19e3..d7e2463e2e 100644 --- a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/config/xml/KafkaConsumerContextParser.java +++ b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/config/xml/KafkaConsumerContextParser.java @@ -15,9 +15,7 @@ */ package org.springframework.integration.kafka.config.xml; -import java.util.HashMap; -import java.util.List; -import java.util.Map; +import java.util.*; import org.springframework.beans.factory.config.BeanDefinition; import org.springframework.beans.factory.config.BeanDefinitionHolder; @@ -26,13 +24,7 @@ import org.springframework.beans.factory.support.BeanDefinitionBuilder; import org.springframework.beans.factory.xml.AbstractSingleBeanDefinitionParser; import org.springframework.beans.factory.xml.ParserContext; import org.springframework.integration.config.xml.IntegrationNamespaceUtils; -import org.springframework.integration.kafka.support.ConsumerConfigFactoryBean; -import org.springframework.integration.kafka.support.ConsumerConfiguration; -import org.springframework.integration.kafka.support.ConsumerConnectionProvider; -import org.springframework.integration.kafka.support.ConsumerMetadata; -import org.springframework.integration.kafka.support.KafkaConsumerContext; -import org.springframework.integration.kafka.support.MessageLeftOverTracker; -import org.springframework.integration.kafka.support.TopicFilterConfiguration; +import org.springframework.integration.kafka.support.*; import org.springframework.util.StringUtils; import org.springframework.util.xml.DomUtils; import org.w3c.dom.Element; @@ -100,6 +92,8 @@ public class KafkaConsumerContextParser extends AbstractSingleBeanDefinitionPars 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")); @@ -107,6 +101,10 @@ public class KafkaConsumerContextParser extends AbstractSingleBeanDefinitionPars consumerConfigFactoryBuilder.addConstructorArgReference(zookeeperConnectBean); } + if (StringUtils.hasText(consumerPropertiesBean)) { + consumerConfigFactoryBuilder.addConstructorArgReference(consumerPropertiesBean); + } + final BeanDefinition consumerConfigFactoryBuilderBeanDefinition = consumerConfigFactoryBuilder.getBeanDefinition(); registerBeanDefinition(new BeanDefinitionHolder(consumerConfigFactoryBuilderBeanDefinition, "consumerConfigFactory_" + consumerConfiguration.getAttribute("group-id")), parserContext.getRegistry()); 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 ee2a8bb941..b3eb7a64aa 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 @@ -46,10 +46,11 @@ public class KafkaProducerContextParser extends AbstractSimpleBeanDefinitionPars super.doParse(element, parserContext, builder); final Element topics = DomUtils.getChildElementByTagName(element, "producer-configurations"); - parseProducerConfigurations(topics, parserContext); + parseProducerConfigurations(topics, parserContext, builder, element); } - private void parseProducerConfigurations(final Element topics, final ParserContext parserContext) { + 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); @@ -68,6 +69,8 @@ public class KafkaProducerContextParser extends AbstractSimpleBeanDefinitionPars registerBeanDefinition(new BeanDefinitionHolder(producerMetadataBeanDef, "producerMetadata_" + producerConfiguration.getAttribute("topic")), parserContext.getRegistry()); + final String producerPropertiesBean = parentElem.getAttribute("producer-properties"); + final BeanDefinitionBuilder producerFactoryBuilder = BeanDefinitionBuilder.genericBeanDefinition(ProducerFactoryBean.class); producerFactoryBuilder.addConstructorArgReference("producerMetadata_" + producerConfiguration.getAttribute("topic")); @@ -76,6 +79,10 @@ public class KafkaProducerContextParser extends AbstractSimpleBeanDefinitionPars producerFactoryBuilder.addConstructorArgValue(producerConfiguration.getAttribute("broker-list")); } + if (StringUtils.hasText(producerPropertiesBean)) { + producerFactoryBuilder.addConstructorArgReference(producerPropertiesBean); + } + final BeanDefinition producerfactoryBeanDefinition = producerFactoryBuilder.getBeanDefinition(); registerBeanDefinition(new BeanDefinitionHolder(producerfactoryBeanDefinition, "prodFactory_" + producerConfiguration.getAttribute("topic")), parserContext.getRegistry()); diff --git a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/support/ConsumerConfigFactoryBean.java b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/support/ConsumerConfigFactoryBean.java index d8def7d99f..6a39d49306 100644 --- a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/support/ConsumerConfigFactoryBean.java +++ b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/support/ConsumerConfigFactoryBean.java @@ -16,6 +16,9 @@ package org.springframework.integration.kafka.support; import kafka.consumer.ConsumerConfig; + +import org.apache.commons.logging.Log; +import org.apache.commons.logging.LogFactory; import org.springframework.beans.factory.FactoryBean; import java.util.Properties; @@ -26,25 +29,42 @@ import java.util.Properties; */ public class ConsumerConfigFactoryBean implements FactoryBean { + private static final Log LOGGER = LogFactory.getLog(ConsumerConfigFactoryBean.class); private final ConsumerMetadata consumerMetadata; private final ZookeeperConnect zookeeperConnect; + private Properties consumerProperties = new Properties(); public ConsumerConfigFactoryBean(final ConsumerMetadata consumerMetadata, - final ZookeeperConnect zookeeperConnect){ + final ZookeeperConnect zookeeperConnect, final Properties consumerProperties) { this.consumerMetadata = consumerMetadata; this.zookeeperConnect = zookeeperConnect; + if (consumerProperties != null) { + this.consumerProperties = consumerProperties; + } } + public ConsumerConfigFactoryBean(final ConsumerMetadata consumerMetadata, final ZookeeperConnect zookeeperConnect) { + this(consumerMetadata, zookeeperConnect, null); + } + @Override public ConsumerConfig getObject() throws Exception { final Properties properties = new Properties(); + properties.putAll(consumerProperties); properties.put("zookeeper.connect", zookeeperConnect.getZkConnect()); properties.put("zookeeper.session.timeout.ms", zookeeperConnect.getZkSessionTimeout()); properties.put("zookeeper.sync.time.ms", zookeeperConnect.getZkSyncTime()); - properties.put("auto.commit.interval.ms", consumerMetadata.getAutoCommitInterval()); - properties.put("consumer.timeout.ms", consumerMetadata.getConsumerTimeout()); + + // Overriding the default value of -1, which will make the consumer to + // wait indefinitely + if (!properties.containsKey("consumer.timeout.ms")) { + properties.put("consumer.timeout.ms", consumerMetadata.getConsumerTimeout()); + } + properties.put("group.id", consumerMetadata.getGroupId()); + LOGGER.info("Using consumer properties => " + properties); + return new ConsumerConfig(properties); } diff --git a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/support/ConsumerConfiguration.java b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/support/ConsumerConfiguration.java index 8376effbc0..6ad6f2ff3e 100644 --- a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/support/ConsumerConfiguration.java +++ b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/support/ConsumerConfiguration.java @@ -76,7 +76,7 @@ public class ConsumerConfiguration { } } } catch (ConsumerTimeoutException cte) { - LOGGER.info("Consumer timed out"); + LOGGER.debug("Consumer timed out"); } return rawMessages; } 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 0fdb3c2078..bf9b047c2b 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 @@ -17,6 +17,7 @@ package org.springframework.integration.kafka.support; import java.util.Collection; import java.util.Map; +import java.util.Properties; import org.apache.commons.logging.Log; import org.apache.commons.logging.LogFactory; @@ -32,6 +33,7 @@ import org.springframework.integration.Message; public class KafkaProducerContext implements BeanFactoryAware { private static final Log LOGGER = LogFactory.getLog(KafkaProducerContext.class); private Map> topicsConfiguration; + private Properties producerProperties; public void send(final Message message) throws Exception { final ProducerConfiguration producerConfiguration = @@ -65,4 +67,19 @@ public class KafkaProducerContext implements BeanFactoryAware { (Map>) (Object) ((ListableBeanFactory)beanFactory).getBeansOfType(ProducerConfiguration.class); } + + /** + * @param producerProperties + * The producerProperties to set. + */ + public void setProducerProperties(Properties producerProperties) { + this.producerProperties = producerProperties; + } + + /** + * @return Returns the producerProperties. + */ + public Properties getProducerProperties() { + return producerProperties; + } } diff --git a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/support/ProducerFactoryBean.java b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/support/ProducerFactoryBean.java index d49fa2ca73..b622154432 100644 --- a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/support/ProducerFactoryBean.java +++ b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/support/ProducerFactoryBean.java @@ -20,7 +20,11 @@ import kafka.producer.ProducerConfig; import kafka.producer.ProducerPool; import kafka.producer.async.DefaultEventHandler; import kafka.producer.async.EventHandler; + +import org.apache.commons.logging.Log; +import org.apache.commons.logging.LogFactory; import org.springframework.beans.factory.FactoryBean; + import scala.collection.mutable.HashMap; import java.util.Properties; @@ -31,17 +35,29 @@ import java.util.Properties; */ public class ProducerFactoryBean implements FactoryBean> { + private static final Log LOGGER = LogFactory.getLog(ProducerFactoryBean.class); + private final String brokerList; private final ProducerMetadata producerMetadata; + private Properties producerProperties = new Properties(); - public ProducerFactoryBean(final ProducerMetadata producerMetadata, final String brokerList){ + public ProducerFactoryBean(final ProducerMetadata producerMetadata, final String brokerList, + final Properties producerProperties) { this.producerMetadata = producerMetadata; this.brokerList = brokerList; + if (producerProperties != null) { + this.producerProperties = producerProperties; + } + } + + public ProducerFactoryBean(final ProducerMetadata producerMetadata, final String brokerList) { + this(producerMetadata, brokerList, null); } @Override public Producer getObject() throws Exception { final Properties props = new Properties(); + props.putAll(producerProperties); props.put("metadata.broker.list", brokerList); props.put("compression.codec", producerMetadata.getCompressionCodec()); @@ -52,6 +68,7 @@ public class ProducerFactoryBean implements FactoryBean> { } } + LOGGER.info("Using producer properties => " + props); final ProducerConfig config = new ProducerConfig(props); final EventHandler eventHandler = new DefaultEventHandler(config, producerMetadata.getPartitioner() == null ? new DefaultPartitioner() : producerMetadata.getPartitioner(), 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 b3009d85ce..866255f2f4 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 @@ -206,6 +206,14 @@ + + + + Kafka producer properties to use for all producers + + + @@ -350,6 +358,14 @@ + + + + Kafka consumer properties to use for all consumers + + + diff --git a/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/config/xml/KafkaConsumerContextParserTests-context.xml b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/config/xml/KafkaConsumerContextParserTests-context.xml index 576dabd3c0..8be6421fa1 100644 --- a/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/config/xml/KafkaConsumerContextParserTests-context.xml +++ b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/config/xml/KafkaConsumerContextParserTests-context.xml @@ -14,9 +14,21 @@ + + + + largest + 10485760 + 5242880 + 1000 + + + + + zookeeper-connect="zookeeperConnect" + consumer-properties="consumerProperties"> - + + + + 3600000 + 5 + 5242880 + + + + +