From bacf8d2cebeaa40fcbefbfc53decf14e25ed4a92 Mon Sep 17 00:00:00 2001 From: Artem Bilan Date: Tue, 12 Aug 2014 22:57:50 +0300 Subject: [PATCH] Kafka-GH-92: Allow PP for the `` Fixes https://github.com/spring-projects/spring-integration-extensions/issues/92 --- .../xml/KafkaConsumerContextParser.java | 34 +- .../support/TopicFilterConfiguration.java | 35 +- .../xml/spring-integration-kafka-1.0.xsd | 767 +++++++++--------- ...afkaConsumerContextParserTests-context.xml | 24 +- .../xml/KafkaConsumerContextParserTests.java | 25 +- .../xml/kafkaInboundAdapterCommon-context.xml | 7 +- 6 files changed, 473 insertions(+), 419 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 d7e2463e2e..dab70c9962 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 @@ -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,9 +13,14 @@ * See the License for the specific language governing permissions and * limitations under the License. */ + package org.springframework.integration.kafka.config.xml; -import java.util.*; +import java.util.HashMap; +import java.util.List; +import java.util.Map; + +import org.w3c.dom.Element; import org.springframework.beans.factory.config.BeanDefinition; import org.springframework.beans.factory.config.BeanDefinitionHolder; @@ -24,14 +29,20 @@ 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.*; +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.util.StringUtils; import org.springframework.util.xml.DomUtils; -import org.w3c.dom.Element; /** * @author Soby Chacko * @author Rajasekar Elango + * @author Artem Bilan * @since 0.5 */ public class KafkaConsumerContextParser extends AbstractSingleBeanDefinitionParser { @@ -50,7 +61,7 @@ public class KafkaConsumerContextParser extends AbstractSingleBeanDefinitionPars } private void parseConsumerConfigurations(final Element consumerConfigurations, final ParserContext parserContext, - final BeanDefinitionBuilder builder, final Element parentElem) { + final BeanDefinitionBuilder builder, final Element parentElem) { for (final Element consumerConfiguration : DomUtils.getChildElementsByTagName(consumerConfigurations, "consumer-configuration")) { final BeanDefinitionBuilder consumerConfigurationBuilder = BeanDefinitionBuilder.genericBeanDefinition(ConsumerConfiguration.class); final BeanDefinitionBuilder consumerMetadataBuilder = BeanDefinitionBuilder.genericBeanDefinition(ConsumerMetadata.class); @@ -68,7 +79,7 @@ public class KafkaConsumerContextParser extends AbstractSingleBeanDefinitionPars final List topicConfigurations = DomUtils.getChildElementsByTagName(consumerConfiguration, "topic"); - if (topicConfigurations != null){ + if (topicConfigurations != null) { for (final Element topicConfiguration : topicConfigurations) { final String topic = topicConfiguration.getAttribute("id"); final String streams = topicConfiguration.getAttribute("streams"); @@ -80,9 +91,14 @@ public class KafkaConsumerContextParser extends AbstractSingleBeanDefinitionPars final Element topicFilter = DomUtils.getChildElementByTagName(consumerConfiguration, "topic-filter"); - if (topicFilter != null){ - final TopicFilterConfiguration topicFilterConfiguration = new TopicFilterConfiguration(topicFilter.getAttribute("pattern"),Integer.valueOf(topicFilter.getAttribute("streams")), Boolean.valueOf(topicFilter.getAttribute("exclude"))); - consumerMetadataBuilder.addPropertyValue("topicFilterConfiguration", topicFilterConfiguration); + if (topicFilter != null) { + BeanDefinition topicFilterConfigurationBeanDefinition = + BeanDefinitionBuilder.genericBeanDefinition(TopicFilterConfiguration.class) + .addConstructorArgValue(topicFilter.getAttribute("pattern")) + .addConstructorArgValue(topicFilter.getAttribute("streams")) + .addConstructorArgValue(topicFilter.getAttribute("exclude")) + .getBeanDefinition(); + consumerMetadataBuilder.addPropertyValue("topicFilterConfiguration", topicFilterConfigurationBeanDefinition); } final BeanDefinition consumerMetadataBeanDef = consumerMetadataBuilder.getBeanDefinition(); diff --git a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/support/TopicFilterConfiguration.java b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/support/TopicFilterConfiguration.java index 7170c6b01d..576fa0b721 100644 --- a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/support/TopicFilterConfiguration.java +++ b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/support/TopicFilterConfiguration.java @@ -1,11 +1,19 @@ /* - * 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 kafka.consumer.Blacklist; @@ -14,33 +22,36 @@ import kafka.consumer.Whitelist; /** * @author Rajasekar Elango + * @author Artem Bilan * @since 0.5 */ public class TopicFilterConfiguration { private final int numberOfStreams; - private TopicFilter topicFilter; + + private final TopicFilter topicFilter; public TopicFilterConfiguration(final String pattern, final int numberOfStreams, final boolean exclude) { this.numberOfStreams = numberOfStreams; if (exclude) { - topicFilter = new Blacklist(pattern); + this.topicFilter = new Blacklist(pattern); } else { - topicFilter = new Whitelist(pattern); + this.topicFilter = new Whitelist(pattern); } } public TopicFilter getTopicFilter() { - return topicFilter; + return this.topicFilter; } public int getNumberOfStreams() { - return numberOfStreams; + return this.numberOfStreams; } @Override public String toString() { - return new StringBuilder(topicFilter.toString()).append(" : ").append(numberOfStreams).toString(); + return this.topicFilter.toString() + " : " + this.numberOfStreams; } + } 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 92a4e1b6c1..1e8efe8143 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 @@ -1,389 +1,392 @@ + xmlns:xsd="http://www.w3.org/2001/XMLSchema" xmlns:beans="http://www.springframework.org/schema/beans" + xmlns:tool="http://www.springframework.org/schema/tool" + xmlns:integration="http://www.springframework.org/schema/integration" + targetNamespace="http://www.springframework.org/schema/integration/kafka" + elementFormDefault="qualified" attributeFormDefault="unqualified"> - - - + + + - - + - + - - - + + - - - - - - + + + + + - - - - - - - - - - + + + + + + + + + - - - - - - - + + + + + + + - - - + + - - - - - - - - - - + + + + + + + + + - - - - - - - - - - + + + + + + + + + - - - - - - - - - + + + + + + + + + - - - + + - - - - - - + + + + + - - - - - - + + + + + - - - - - - The topic configured by this configuration. - - - - - - + + + + + The topic configured by this configuration. + + + + + + - - - - - - - - - - - Custom implemenation of a Kafka Encoder for encoding message values. - - - - - - - Custom implemenation of a Kafka Encoder for encoding message keys. - - - - - - - Class type used for the key - - - - - - - Class type used for the value - - - - - - + + + + + + + + + + Custom implemenation of a Kafka Encoder for encoding message values. + + + + + + + Custom implemenation of a Kafka Encoder for encoding message keys. + + + + + + + Class type used for the key + + + + + + + Class type used for the value + + + + + + - - - - - - - - - - - Custom Kafka key partitioner. - - - - - - - Indicates if this producer is async or not. - - - - - - - number of messages to batch at this producer. - - - - - - - - - - - - - - Kafka producer properties to use for all producers - - - - - + + + + + + + + + + + Custom Kafka key partitioner. + + + + + + + Indicates if this producer is async or not. + + + + + + + number of messages to batch at this producer. + + + + + + + + + + + + + + Kafka producer properties to use for all producers + + + + + - - - + + - - - - - - + + + + + - - - - - - + + + + + - - - - - - - - - - - - - - + + + + + + + + + + + + + - - - - - - - - - - + + + + + + + + + - - - - - - - - - - + + + + + + + + + - - - - - - - - - - - - - + + + + + + + + + + + + + + + - - - - - - - - - - + + + + + + + + + - - - - - - - - - - - Kafka Server Bean Name - - - - - - - Kafka Server Bean Name - - - - - - - + + + + + + + + + + + Kafka Server Bean Name + + + + + + + Kafka Server Bean Name + + + + + + + - - - - - - + + + + + - - - - - - - - - - - Kafka Server Bean Name - - - - - - - Kafka consumer properties to use for all consumers - - - - - + + + + + + + + + + + Kafka Server Bean Name + + + + + + + Kafka consumer properties to use for all consumers + + + + + - - - - The definition for the Spring Integration Kafka - Inbound Channel Adapter. - - - - - - - - - - + + + The definition for the Spring Integration Kafka + Inbound Channel Adapter. + + + + + + + + + + - - - - - - + + + + + - - - - - - - - - - - Kafka Server Bean Name - - - - - + + + + + + + + + + + Kafka Server Bean Name + + + + + - - - - Defines kafka outbound channel adapter that writes the contents of the - Message to kafka broker. - - + + + + Defines kafka outbound channel adapter that writes the contents of the + Message to kafka broker. + + - - - - + + + + - - - - Kafka producer context reference. - - - - - - - Specifies the order for invocation when this endpoint is connected as a - subscriber to a SubscribableChannel. - - - - - + + + + Kafka producer context reference. + + + + + + + Specifies the order for invocation when this endpoint is connected as a + subscriber to a SubscribableChannel. + + + + + 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 8be6421fa1..d787a2dc5a 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 @@ -1,9 +1,11 @@ + 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"> + + foo + 10 + true + + + + - - + + diff --git a/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/config/xml/KafkaConsumerContextParserTests.java b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/config/xml/KafkaConsumerContextParserTests.java index d91210c2a1..f53d0b5559 100644 --- a/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/config/xml/KafkaConsumerContextParserTests.java +++ b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/config/xml/KafkaConsumerContextParserTests.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. @@ -15,23 +15,29 @@ */ package org.springframework.integration.kafka.config.xml; -import org.junit.Assert; +import static org.junit.Assert.*; + +import kafka.consumer.Blacklist; +import org.hamcrest.Matchers; import org.junit.Test; import org.junit.runner.RunWith; + import org.springframework.beans.factory.annotation.Autowired; import org.springframework.context.ApplicationContext; import org.springframework.integration.kafka.support.ConsumerMetadata; import org.springframework.integration.kafka.support.KafkaConsumerContext; +import org.springframework.integration.kafka.support.TopicFilterConfiguration; import org.springframework.test.context.ContextConfiguration; import org.springframework.test.context.junit4.SpringJUnit4ClassRunner; /** * @author Soby Chacko + * @author Artem Bilan * @since 0.5 */ @RunWith(SpringJUnit4ClassRunner.class) @ContextConfiguration -public class KafkaConsumerContextParserTests { +public class KafkaConsumerContextParserTests { @Autowired private ApplicationContext appContext; @@ -39,10 +45,15 @@ public class KafkaConsumerContextParserTests { @Test @SuppressWarnings("unchecked") public void testConsumerContextConfiguration() { - final KafkaConsumerContext consumerContext = appContext.getBean("consumerContext", KafkaConsumerContext.class); - Assert.assertNotNull(consumerContext); + final KafkaConsumerContext consumerContext = + appContext.getBean("consumerContext", KafkaConsumerContext.class); + assertNotNull(consumerContext); - final ConsumerMetadata cm = appContext.getBean("consumerMetadata_default1", ConsumerMetadata.class); - Assert.assertNotNull(cm); + ConsumerMetadata cm = appContext.getBean("consumerMetadata_default1", ConsumerMetadata.class); + assertNotNull(cm); + TopicFilterConfiguration topicFilterConfiguration = cm.getTopicFilterConfiguration(); + assertEquals("foo : 10", topicFilterConfiguration.toString()); + assertThat(topicFilterConfiguration.getTopicFilter(), Matchers.instanceOf(Blacklist.class)); } + } diff --git a/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/config/xml/kafkaInboundAdapterCommon-context.xml b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/config/xml/kafkaInboundAdapterCommon-context.xml index f7fecae4b0..42e907ef89 100644 --- a/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/config/xml/kafkaInboundAdapterCommon-context.xml +++ b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/config/xml/kafkaInboundAdapterCommon-context.xml @@ -9,9 +9,10 @@ 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"> - - - + + + +