diff --git a/.gitignore b/.gitignore
index afb170e..0890536 100644
--- a/.gitignore
+++ b/.gitignore
@@ -19,4 +19,5 @@ logs
nohup.out
out
target
-spring-integration-aws/src/test/resources/awscredentials.properties
\ No newline at end of file
+spring-integration-aws/src/test/resources/awscredentials.properties
+*.orig
diff --git a/samples/kafka/src/main/resources/org/springframework/integration/samples/kafka/inbound/kafkaInboundAdapterParserTests-context.xml b/samples/kafka/src/main/resources/org/springframework/integration/samples/kafka/inbound/kafkaInboundAdapterParserTests-context.xml
index 5dce6af..ee9cd6a 100644
--- a/samples/kafka/src/main/resources/org/springframework/integration/samples/kafka/inbound/kafkaInboundAdapterParserTests-context.xml
+++ b/samples/kafka/src/main/resources/org/springframework/integration/samples/kafka/inbound/kafkaInboundAdapterParserTests-context.xml
@@ -34,9 +34,20 @@
+
+
+
+ smallest
+ 10485760
+ 5242880
+ 1000
+
+
+
+
+ zookeeper-connect="zookeeperConnect" consumer-properties="consumerProperties">
-
+
\ No newline at end of file
diff --git a/samples/kafka/src/main/resources/org/springframework/integration/samples/kafka/outbound/kafkaOutboundAdapterParserTests-context.xml b/samples/kafka/src/main/resources/org/springframework/integration/samples/kafka/outbound/kafkaOutboundAdapterParserTests-context.xml
index 1fa2155..b1f0d8d 100644
--- a/samples/kafka/src/main/resources/org/springframework/integration/samples/kafka/outbound/kafkaOutboundAdapterParserTests-context.xml
+++ b/samples/kafka/src/main/resources/org/springframework/integration/samples/kafka/outbound/kafkaOutboundAdapterParserTests-context.xml
@@ -31,7 +31,17 @@
-
+
+
+
+ 3600000
+ 5
+ 5242880
+
+
+
+
+
+
+
+ 3600000
+ 5
+ 5242880
+
+
+
+
+
+
+ ...
+ ...
+ ...
+
+
+```
+
Inbound Channel Adapter:
--------------------------------------------
@@ -326,3 +351,32 @@ Map.
If your use case does not require ordering of messages during consumption, then you can easily pass this
payload to a standard SI transformer and just get a full dump of the actual payload sent by Kafka.
+
+#### Tuning Consumer Properties
+Kafka Consumer API provides several [Consumer Configs] (http://kafka.apache.org/documentation.html#consumerconfigs) to fine tune consumers.
+To specify those properties, `consumer-context` element supports optional `consumer-properties` attribute that can reference the spring properties bean.
+This properties will be applied to all Consumer Configurations within the consumer context. For Eg:
+
+```xml
+
+
+
+
+ smallest
+ 10485760
+ 5242880
+ 1000
+
+
+
+
+
+
+ ...
+ ...
+ ...
+ >
+
+```
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 c21f7f1..d7e2463 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 ee2a8bb..b3eb7a6 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 d8def7d..6a39d49 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 8376eff..6ad6f2f 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 0fdb3c2..bf9b047 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 d49fa2c..b622154 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 b3009d8..866255f 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 576dabd..8be6421 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
+
+
+
+
+