From a21a2444ef952a31698d59a53d8aab236b6afef7 Mon Sep 17 00:00:00 2001 From: Artem Bilan Date: Wed, 21 Dec 2016 17:15:09 -0500 Subject: [PATCH] Change Kafka Samples to use Boot auto-config Remove `@IntegrationComponentScan` from Boot samples Since the latest Spring Boot utilize already `@IntegrationComponentScan` for the `@SpringBootApplication` class package, there is no reason to worry about lost `@MessagingGateway` Reflect Kafka samples changes according latest Spring Boot fixes --- .../DynamicTcpClientApplication.java | 2 - .../samples/filesplit/Application.java | 18 +----- .../samples/kafka/Application.java | 58 ++++++------------- .../kafka/src/main/resources/application.yml | 18 +++++- .../integration/samples/mqtt/Application.java | 2 - .../samples/dsl/cafe/lambda/Application.java | 2 - .../samples/dsl/kafka/Application.java | 47 ++------------- .../src/main/resources/application.yml | 16 ++++- 8 files changed, 57 insertions(+), 106 deletions(-) diff --git a/advanced/dynamic-tcp-client/src/main/java/org/springframework/integration/samples/dynamictcp/DynamicTcpClientApplication.java b/advanced/dynamic-tcp-client/src/main/java/org/springframework/integration/samples/dynamictcp/DynamicTcpClientApplication.java index d32ddd90..640fa079 100644 --- a/advanced/dynamic-tcp-client/src/main/java/org/springframework/integration/samples/dynamictcp/DynamicTcpClientApplication.java +++ b/advanced/dynamic-tcp-client/src/main/java/org/springframework/integration/samples/dynamictcp/DynamicTcpClientApplication.java @@ -10,7 +10,6 @@ import org.springframework.boot.SpringApplication; import org.springframework.boot.autoconfigure.SpringBootApplication; import org.springframework.context.ConfigurableApplicationContext; import org.springframework.context.annotation.Bean; -import org.springframework.integration.annotation.IntegrationComponentScan; import org.springframework.integration.annotation.MessagingGateway; import org.springframework.integration.channel.QueueChannel; import org.springframework.integration.config.EnableMessageHistory; @@ -28,7 +27,6 @@ import org.springframework.messaging.handler.annotation.Header; import org.springframework.util.Assert; @SpringBootApplication -@IntegrationComponentScan @EnableMessageHistory public class DynamicTcpClientApplication { diff --git a/applications/file-split-ftp/src/main/java/org/springframework/integration/samples/filesplit/Application.java b/applications/file-split-ftp/src/main/java/org/springframework/integration/samples/filesplit/Application.java index aec6c840..458cf7bb 100644 --- a/applications/file-split-ftp/src/main/java/org/springframework/integration/samples/filesplit/Application.java +++ b/applications/file-split-ftp/src/main/java/org/springframework/integration/samples/filesplit/Application.java @@ -22,6 +22,7 @@ import java.io.StringWriter; import org.aopalliance.intercept.MethodInterceptor; import org.apache.commons.net.ftp.FTPFile; + import org.springframework.beans.factory.annotation.Autowired; import org.springframework.boot.SpringApplication; import org.springframework.boot.autoconfigure.SpringBootApplication; @@ -44,12 +45,9 @@ import org.springframework.integration.mail.dsl.Mail; import org.springframework.messaging.Message; import org.springframework.messaging.MessageHeaders; import org.springframework.messaging.MessagingException; -import org.springframework.web.servlet.config.annotation.CorsRegistry; -import org.springframework.web.servlet.config.annotation.WebMvcConfigurer; -import org.springframework.web.servlet.config.annotation.WebMvcConfigurerAdapter; @SpringBootApplication -@EnableIntegrationGraphController +@EnableIntegrationGraphController(allowedOrigins = "http://localhost:8082") public class Application { private static final String EMAIL_SUCCESS_SUFFIX = "emailSuccessSuffix"; @@ -221,18 +219,6 @@ public class Application { }; } - // Integration Graph CORS - @Bean - public WebMvcConfigurer corsConfigurer() { - return new WebMvcConfigurerAdapter() { - - @Override - public void addCorsMappings(CorsRegistry registry) { - registry.addMapping("/integration").allowedOrigins("http://localhost:8082"); - } - }; - } - private String getStackTraceAsString(Throwable cause) { StringWriter stringWriter = new StringWriter(); PrintWriter printWriter = new PrintWriter(stringWriter, true); diff --git a/basic/kafka/src/main/java/org/springframework/integration/samples/kafka/Application.java b/basic/kafka/src/main/java/org/springframework/integration/samples/kafka/Application.java index 8af9f386..c7a19ad5 100644 --- a/basic/kafka/src/main/java/org/springframework/integration/samples/kafka/Application.java +++ b/basic/kafka/src/main/java/org/springframework/integration/samples/kafka/Application.java @@ -16,18 +16,17 @@ package org.springframework.integration.samples.kafka; -import java.util.HashMap; import java.util.Map; import java.util.Properties; import org.I0Itec.zkclient.ZkClient; import org.apache.kafka.clients.consumer.ConsumerConfig; import org.apache.kafka.clients.producer.ProducerConfig; -import org.apache.kafka.common.serialization.StringDeserializer; -import org.apache.kafka.common.serialization.StringSerializer; +import org.apache.kafka.common.errors.TopicExistsException; import org.springframework.beans.factory.annotation.Value; import org.springframework.boot.autoconfigure.SpringBootApplication; +import org.springframework.boot.autoconfigure.kafka.KafkaProperties; import org.springframework.boot.builder.SpringApplicationBuilder; import org.springframework.context.ConfigurableApplicationContext; import org.springframework.context.SmartLifecycle; @@ -53,7 +52,6 @@ import org.springframework.messaging.PollableChannel; import org.springframework.messaging.support.GenericMessage; import kafka.admin.AdminUtils; -import org.apache.kafka.common.errors.TopicExistsException; import kafka.utils.ZKStringSerializer$; import kafka.utils.ZkUtils; @@ -70,9 +68,6 @@ public class Application { @Value("${kafka.messageKey}") private String messageKey; - @Value("${kafka.broker.address}") - private String brokerAddress; - @Value("${kafka.zookeeper.connect}") private String zookeeperConnect; @@ -96,53 +91,38 @@ public class Application { System.exit(0); } + @Bean + public ProducerFactory kafkaProducerFactory(KafkaProperties properties) { + Map producerProperties = properties.buildProducerProperties(); + producerProperties.put(ProducerConfig.LINGER_MS_CONFIG, 1); + return new DefaultKafkaProducerFactory<>(producerProperties); + } + @ServiceActivator(inputChannel = "toKafka") @Bean - public MessageHandler handler() throws Exception { + public MessageHandler handler(KafkaTemplate kafkaTemplate) { KafkaProducerMessageHandler handler = - new KafkaProducerMessageHandler<>(kafkaTemplate()); + new KafkaProducerMessageHandler<>(kafkaTemplate); handler.setTopicExpression(new LiteralExpression(this.topic)); handler.setMessageKeyExpression(new LiteralExpression(this.messageKey)); return handler; } @Bean - public KafkaTemplate kafkaTemplate() { - return new KafkaTemplate<>(producerFactory()); + public ConsumerFactory kafkaConsumerFactory(KafkaProperties properties) { + Map consumerProperties = properties + .buildConsumerProperties(); + consumerProperties.put(ConsumerConfig.SESSION_TIMEOUT_MS_CONFIG, 15000); + return new DefaultKafkaConsumerFactory<>(consumerProperties); } @Bean - public ProducerFactory producerFactory() { - Map props = new HashMap<>(); - props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, this.brokerAddress); - props.put(ProducerConfig.RETRIES_CONFIG, 0); - props.put(ProducerConfig.BATCH_SIZE_CONFIG, 16384); - props.put(ProducerConfig.LINGER_MS_CONFIG, 1); - props.put(ProducerConfig.BUFFER_MEMORY_CONFIG, 33554432); - props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class); - props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class); - return new DefaultKafkaProducerFactory<>(props); - } - - @Bean - public KafkaMessageListenerContainer container() throws Exception { - return new KafkaMessageListenerContainer<>(consumerFactory(), + public KafkaMessageListenerContainer container( + ConsumerFactory kafkaConsumerFactory) { + return new KafkaMessageListenerContainer<>(kafkaConsumerFactory, new ContainerProperties(new TopicPartitionInitialOffset(this.topic, 0))); } - @Bean - public ConsumerFactory consumerFactory() { - Map props = new HashMap<>(); - props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, this.brokerAddress); - props.put(ConsumerConfig.GROUP_ID_CONFIG, "siTestGroup"); - props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, true); - props.put(ConsumerConfig.AUTO_COMMIT_INTERVAL_MS_CONFIG, 100); - props.put(ConsumerConfig.SESSION_TIMEOUT_MS_CONFIG, 15000); - props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class); - props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class); - return new DefaultKafkaConsumerFactory<>(props); - } - @Bean public KafkaMessageDrivenChannelAdapter adapter(KafkaMessageListenerContainer container) { diff --git a/basic/kafka/src/main/resources/application.yml b/basic/kafka/src/main/resources/application.yml index 02f3a96c..6bf3b6fb 100644 --- a/basic/kafka/src/main/resources/application.yml +++ b/basic/kafka/src/main/resources/application.yml @@ -1,7 +1,21 @@ kafka: - broker: - address: localhost:9092 zookeeper: connect: localhost:2181 topic: si.topic messageKey: si.key +spring: + kafka: + consumer: + group-id: siTestGroup + enable-auto-commit: true + auto-commit-interval: 100 + value-deserializer: org.apache.kafka.common.serialization.StringDeserializer + key-deserializer: org.apache.kafka.common.serialization.StringDeserializer + producer: + batch-size: 16384 + buffer-memory: 33554432 + retries: 0 + key-serializer: org.apache.kafka.common.serialization.StringSerializer + value-serializer: org.apache.kafka.common.serialization.StringSerializer + + diff --git a/basic/mqtt/src/main/java/org/springframework/integration/samples/mqtt/Application.java b/basic/mqtt/src/main/java/org/springframework/integration/samples/mqtt/Application.java index c2dd1bca..de3e01ae 100644 --- a/basic/mqtt/src/main/java/org/springframework/integration/samples/mqtt/Application.java +++ b/basic/mqtt/src/main/java/org/springframework/integration/samples/mqtt/Application.java @@ -20,7 +20,6 @@ import org.apache.log4j.Logger; import org.springframework.boot.SpringApplication; import org.springframework.boot.autoconfigure.SpringBootApplication; import org.springframework.context.annotation.Bean; -import org.springframework.integration.annotation.IntegrationComponentScan; import org.springframework.integration.dsl.IntegrationFlow; import org.springframework.integration.dsl.IntegrationFlows; import org.springframework.integration.dsl.Pollers; @@ -41,7 +40,6 @@ import org.springframework.messaging.MessageHandler; * */ @SpringBootApplication -@IntegrationComponentScan public class Application { private static final Logger LOGGER = Logger.getLogger(Application.class); diff --git a/dsl/cafe-dsl/src/main/java/org/springframework/integration/samples/dsl/cafe/lambda/Application.java b/dsl/cafe-dsl/src/main/java/org/springframework/integration/samples/dsl/cafe/lambda/Application.java index 4af38c92..688625ba 100644 --- a/dsl/cafe-dsl/src/main/java/org/springframework/integration/samples/dsl/cafe/lambda/Application.java +++ b/dsl/cafe-dsl/src/main/java/org/springframework/integration/samples/dsl/cafe/lambda/Application.java @@ -26,7 +26,6 @@ import org.springframework.boot.autoconfigure.SpringBootApplication; import org.springframework.context.ConfigurableApplicationContext; import org.springframework.context.annotation.Bean; import org.springframework.integration.annotation.Gateway; -import org.springframework.integration.annotation.IntegrationComponentScan; import org.springframework.integration.annotation.MessagingGateway; import org.springframework.integration.dsl.IntegrationFlow; import org.springframework.integration.dsl.Pollers; @@ -45,7 +44,6 @@ import com.google.common.util.concurrent.Uninterruptibles; * @since 3.0 */ @SpringBootApplication -@IntegrationComponentScan public class Application { public static void main(String[] args) throws Exception { diff --git a/dsl/kafka-dsl/src/main/java/org/springframework/integration/samples/dsl/kafka/Application.java b/dsl/kafka-dsl/src/main/java/org/springframework/integration/samples/dsl/kafka/Application.java index 96a91b76..e16b1648 100644 --- a/dsl/kafka-dsl/src/main/java/org/springframework/integration/samples/dsl/kafka/Application.java +++ b/dsl/kafka-dsl/src/main/java/org/springframework/integration/samples/dsl/kafka/Application.java @@ -16,17 +16,12 @@ package org.springframework.integration.samples.dsl.kafka; -import java.util.HashMap; -import java.util.Map; import java.util.Properties; import javax.annotation.PostConstruct; import org.I0Itec.zkclient.ZkClient; -import org.apache.kafka.clients.consumer.ConsumerConfig; -import org.apache.kafka.clients.producer.ProducerConfig; -import org.apache.kafka.common.serialization.StringDeserializer; -import org.apache.kafka.common.serialization.StringSerializer; +import org.apache.kafka.common.errors.TopicExistsException; import org.springframework.beans.factory.annotation.Value; import org.springframework.boot.autoconfigure.SpringBootApplication; @@ -34,19 +29,15 @@ import org.springframework.boot.builder.SpringApplicationBuilder; import org.springframework.context.ConfigurableApplicationContext; import org.springframework.context.annotation.Bean; import org.springframework.integration.annotation.Gateway; -import org.springframework.integration.annotation.IntegrationComponentScan; import org.springframework.integration.annotation.MessagingGateway; import org.springframework.integration.dsl.IntegrationFlow; import org.springframework.integration.dsl.IntegrationFlows; import org.springframework.integration.kafka.dsl.Kafka; import org.springframework.kafka.core.ConsumerFactory; -import org.springframework.kafka.core.DefaultKafkaConsumerFactory; -import org.springframework.kafka.core.DefaultKafkaProducerFactory; -import org.springframework.kafka.core.ProducerFactory; +import org.springframework.kafka.core.KafkaTemplate; import org.springframework.messaging.Message; import kafka.admin.AdminUtils; -import org.apache.kafka.common.errors.TopicExistsException; import kafka.utils.ZKStringSerializer$; import kafka.utils.ZkUtils; @@ -56,7 +47,6 @@ import kafka.utils.ZkUtils; * @since 4.3 */ @SpringBootApplication -@IntegrationComponentScan public class Application { public static void main(String[] args) throws Exception { @@ -88,9 +78,6 @@ public class Application { @Value("${kafka.messageKey}") private String messageKey; - @Value("${kafka.broker.address}") - private String brokerAddress; - @Value("${kafka.zookeeper.connect}") private String zookeeperConnect; @@ -117,40 +104,18 @@ public class Application { } - @Bean - public ProducerFactory producerFactory() { - Map props = new HashMap<>(); - props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, this.brokerAddress); - props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class); - props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class); - return new DefaultKafkaProducerFactory<>(props); - } - - @Bean - public IntegrationFlow toKafka() { + public IntegrationFlow toKafka(KafkaTemplate kafkaTemplate) { return f -> f - .handle(Kafka.outboundChannelAdapter(producerFactory()) + .handle(Kafka.outboundChannelAdapter(kafkaTemplate) .topic(this.topic) .messageKey(this.messageKey)); } @Bean - public ConsumerFactory consumerFactory() { - Map props = new HashMap<>(); - props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, this.brokerAddress); - props.put(ConsumerConfig.GROUP_ID_CONFIG, "siTestGroup"); - props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest"); - props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, true); - props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class); - props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class); - return new DefaultKafkaConsumerFactory<>(props); - } - - @Bean - public IntegrationFlow fromKafka() { + public IntegrationFlow fromKafka(ConsumerFactory consumerFactory) { return IntegrationFlows - .from(Kafka.messageDrivenChannelAdapter(consumerFactory(), this.topic)) + .from(Kafka.messageDrivenChannelAdapter(consumerFactory, this.topic)) .channel(c -> c.queue("fromKafka")) .get(); } diff --git a/dsl/kafka-dsl/src/main/resources/application.yml b/dsl/kafka-dsl/src/main/resources/application.yml index 02f3a96c..d0efe847 100644 --- a/dsl/kafka-dsl/src/main/resources/application.yml +++ b/dsl/kafka-dsl/src/main/resources/application.yml @@ -1,7 +1,19 @@ kafka: - broker: - address: localhost:9092 zookeeper: connect: localhost:2181 topic: si.topic messageKey: si.key +spring: + kafka: + consumer: + group-id: siTestGroup + enable-auto-commit: true + auto-commit-interval: 100 + value-deserializer: org.apache.kafka.common.serialization.StringDeserializer + key-deserializer: org.apache.kafka.common.serialization.StringDeserializer + producer: + batch-size: 16384 + buffer-memory: 33554432 + retries: 0 + key-serializer: org.apache.kafka.common.serialization.StringSerializer + value-serializer: org.apache.kafka.common.serialization.StringSerializer