INTEXT-92 Kafka: Add support for specifying any producer/consumer property
* Update README.md for producer and consumer properties usage
This commit is contained in:
committed by
Artem Bilan
parent
d226b7c3c5
commit
c52be66a1c
@@ -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());
|
||||
|
||||
|
||||
@@ -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());
|
||||
|
||||
|
||||
@@ -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<K,V> implements FactoryBean<ConsumerConfig> {
|
||||
|
||||
private static final Log LOGGER = LogFactory.getLog(ConsumerConfigFactoryBean.class);
|
||||
private final ConsumerMetadata<K,V> consumerMetadata;
|
||||
private final ZookeeperConnect zookeeperConnect;
|
||||
private Properties consumerProperties = new Properties();
|
||||
|
||||
public ConsumerConfigFactoryBean(final ConsumerMetadata<K,V> 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);
|
||||
}
|
||||
|
||||
|
||||
@@ -76,7 +76,7 @@ public class ConsumerConfiguration<K, V> {
|
||||
}
|
||||
}
|
||||
} catch (ConsumerTimeoutException cte) {
|
||||
LOGGER.info("Consumer timed out");
|
||||
LOGGER.debug("Consumer timed out");
|
||||
}
|
||||
return rawMessages;
|
||||
}
|
||||
|
||||
@@ -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<K,V> implements BeanFactoryAware {
|
||||
private static final Log LOGGER = LogFactory.getLog(KafkaProducerContext.class);
|
||||
private Map<String, ProducerConfiguration<K,V>> topicsConfiguration;
|
||||
private Properties producerProperties;
|
||||
|
||||
public void send(final Message<?> message) throws Exception {
|
||||
final ProducerConfiguration<K,V> producerConfiguration =
|
||||
@@ -65,4 +67,19 @@ public class KafkaProducerContext<K,V> implements BeanFactoryAware {
|
||||
(Map<String, ProducerConfiguration<K,V>>) (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;
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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<K,V> implements FactoryBean<Producer<K,V>> {
|
||||
|
||||
private static final Log LOGGER = LogFactory.getLog(ProducerFactoryBean.class);
|
||||
|
||||
private final String brokerList;
|
||||
private final ProducerMetadata<K,V> producerMetadata;
|
||||
private Properties producerProperties = new Properties();
|
||||
|
||||
public ProducerFactoryBean(final ProducerMetadata<K,V> producerMetadata, final String brokerList){
|
||||
public ProducerFactoryBean(final ProducerMetadata<K, V> producerMetadata, final String brokerList,
|
||||
final Properties producerProperties) {
|
||||
this.producerMetadata = producerMetadata;
|
||||
this.brokerList = brokerList;
|
||||
if (producerProperties != null) {
|
||||
this.producerProperties = producerProperties;
|
||||
}
|
||||
}
|
||||
|
||||
public ProducerFactoryBean(final ProducerMetadata<K, V> producerMetadata, final String brokerList) {
|
||||
this(producerMetadata, brokerList, null);
|
||||
}
|
||||
|
||||
@Override
|
||||
public Producer<K, V> 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<K,V> implements FactoryBean<Producer<K,V>> {
|
||||
}
|
||||
}
|
||||
|
||||
LOGGER.info("Using producer properties => " + props);
|
||||
final ProducerConfig config = new ProducerConfig(props);
|
||||
final EventHandler<K, V> eventHandler = new DefaultEventHandler<K, V>(config,
|
||||
producerMetadata.getPartitioner() == null ? new DefaultPartitioner<K>() : producerMetadata.getPartitioner(),
|
||||
|
||||
@@ -206,6 +206,14 @@
|
||||
</xsd:element>
|
||||
</xsd:sequence>
|
||||
<xsd:attribute name="id" type="xsd:string" use="required"/>
|
||||
<xsd:attribute name="producer-properties" use="optional"
|
||||
type="xsd:string">
|
||||
<xsd:annotation>
|
||||
<xsd:documentation>
|
||||
Kafka producer properties to use for all producers
|
||||
</xsd:documentation>
|
||||
</xsd:annotation>
|
||||
</xsd:attribute>
|
||||
</xsd:complexType>
|
||||
</xsd:element>
|
||||
|
||||
@@ -350,6 +358,14 @@
|
||||
</xsd:documentation>
|
||||
</xsd:annotation>
|
||||
</xsd:attribute>
|
||||
<xsd:attribute name="consumer-properties" use="optional"
|
||||
type="xsd:string">
|
||||
<xsd:annotation>
|
||||
<xsd:documentation>
|
||||
Kafka consumer properties to use for all consumers
|
||||
</xsd:documentation>
|
||||
</xsd:annotation>
|
||||
</xsd:attribute>
|
||||
</xsd:complexType>
|
||||
</xsd:element>
|
||||
|
||||
|
||||
@@ -14,9 +14,21 @@
|
||||
<constructor-arg type="java.lang.Class" value="java.lang.String"/>
|
||||
</bean>
|
||||
|
||||
<bean id="consumerProperties" class="org.springframework.beans.factory.config.PropertiesFactoryBean">
|
||||
<property name="properties">
|
||||
<props>
|
||||
<prop key="auto.offset.reset">largest</prop>
|
||||
<prop key="socket.receive.buffer.bytes">10485760</prop> <!-- 10M -->
|
||||
<prop key="fetch.message.max.bytes">5242880</prop>
|
||||
<prop key="auto.commit.interval.ms">1000</prop>
|
||||
</props>
|
||||
</property>
|
||||
</bean>
|
||||
|
||||
<int-kafka:consumer-context id="consumerContext"
|
||||
consumer-timeout="4000"
|
||||
zookeeper-connect="zookeeperConnect">
|
||||
zookeeper-connect="zookeeperConnect"
|
||||
consumer-properties="consumerProperties">
|
||||
<int-kafka:consumer-configurations>
|
||||
<int-kafka:consumer-configuration group-id="default1"
|
||||
value-decoder="valueDecoder"
|
||||
|
||||
@@ -5,7 +5,17 @@
|
||||
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">
|
||||
|
||||
<int-kafka:producer-context id="producerContext">
|
||||
<bean id="producerProperties" class="org.springframework.beans.factory.config.PropertiesFactoryBean">
|
||||
<property name="properties">
|
||||
<props>
|
||||
<prop key="topic.metadata.refresh.interval.ms">3600000</prop>
|
||||
<prop key="message.send.max.retries">5</prop>
|
||||
<prop key="send.buffer.bytes">5242880</prop>
|
||||
</props>
|
||||
</property>
|
||||
</bean>
|
||||
|
||||
<int-kafka:producer-context id="producerContext" producer-properties="producerProperties">
|
||||
<int-kafka:producer-configurations>
|
||||
<int-kafka:producer-configuration broker-list="localhost:9092"
|
||||
key-class-type="java.lang.String"
|
||||
|
||||
Reference in New Issue
Block a user