Kafka Sample Improvements

- get rid of topic creator and scala deps
- use boot properties
- demonstrate using the DSL to dynamically add adapters
This commit is contained in:
Gary Russell
2017-02-01 16:04:32 -05:00
committed by Artem Bilan
parent 65368f7959
commit 038ee4576a
10 changed files with 282 additions and 130 deletions

View File

@@ -1,5 +1,5 @@
/*
* Copyright 2015 the original author or authors.
* Copyright 2015-2017 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.
@@ -16,24 +16,27 @@
package org.springframework.integration.samples.kafka;
import java.util.Collections;
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.errors.TopicExistsException;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.boot.WebApplicationType;
import org.springframework.boot.autoconfigure.SpringBootApplication;
import org.springframework.boot.autoconfigure.kafka.KafkaProperties;
import org.springframework.boot.builder.SpringApplicationBuilder;
import org.springframework.boot.context.properties.EnableConfigurationProperties;
import org.springframework.context.ConfigurableApplicationContext;
import org.springframework.context.SmartLifecycle;
import org.springframework.context.annotation.Bean;
import org.springframework.expression.common.LiteralExpression;
import org.springframework.integration.annotation.ServiceActivator;
import org.springframework.integration.channel.QueueChannel;
import org.springframework.integration.dsl.IntegrationFlow;
import org.springframework.integration.dsl.IntegrationFlows;
import org.springframework.integration.dsl.context.IntegrationFlowContext;
import org.springframework.integration.kafka.dsl.Kafka;
import org.springframework.integration.kafka.inbound.KafkaMessageDrivenChannelAdapter;
import org.springframework.integration.kafka.outbound.KafkaProducerMessageHandler;
import org.springframework.kafka.core.ConsumerFactory;
@@ -43,6 +46,7 @@ import org.springframework.kafka.core.KafkaTemplate;
import org.springframework.kafka.core.ProducerFactory;
import org.springframework.kafka.listener.KafkaMessageListenerContainer;
import org.springframework.kafka.listener.config.ContainerProperties;
import org.springframework.kafka.support.KafkaHeaders;
import org.springframework.kafka.support.KafkaNull;
import org.springframework.kafka.support.TopicPartitionInitialOffset;
import org.springframework.messaging.Message;
@@ -51,44 +55,54 @@ import org.springframework.messaging.MessageHandler;
import org.springframework.messaging.PollableChannel;
import org.springframework.messaging.support.GenericMessage;
import kafka.admin.AdminUtils;
import kafka.utils.ZKStringSerializer$;
import kafka.utils.ZkUtils;
/**
* @author Gary Russell
* @since 4.2
*/
@SpringBootApplication
@EnableConfigurationProperties(KafkaAppProperties.class)
public class Application {
@Value("${kafka.topic}")
private String topic;
@Value("${kafka.messageKey}")
private String messageKey;
@Value("${kafka.zookeeper.connect}")
private String zookeeperConnect;
@Autowired
private KafkaAppProperties properties;
public static void main(String[] args) throws Exception {
ConfigurableApplicationContext context
= new SpringApplicationBuilder(Application.class)
.web(false)
.run(args);
.web(WebApplicationType.NONE)
.run(args);
context.getBean(Application.class).runDemo(context);
context.close();
}
private void runDemo(ConfigurableApplicationContext context) {
MessageChannel toKafka = context.getBean("toKafka", MessageChannel.class);
System.out.println("Sending 10 messages...");
Map<String, Object> headers = Collections.singletonMap(KafkaHeaders.TOPIC, this.properties.getTopic());
for (int i = 0; i < 10; i++) {
toKafka.send(new GenericMessage<>("foo" + i));
toKafka.send(new GenericMessage<>("foo" + i, headers));
}
toKafka.send(new GenericMessage<>(KafkaNull.INSTANCE));
PollableChannel fromKafka = context.getBean("received", PollableChannel.class);
System.out.println("Sending a null message...");
toKafka.send(new GenericMessage<>(KafkaNull.INSTANCE, headers));
PollableChannel fromKafka = context.getBean("fromKafka", PollableChannel.class);
Message<?> received = fromKafka.receive(10000);
int count = 0;
while (received != null) {
System.out.println(received);
received = fromKafka.receive(10000);
received = fromKafka.receive(++count < 11 ? 10000 : 1000);
}
System.out.println("Adding an adapter for a second topic and sending 10 messages...");
addAnotherListenerForTopics(this.properties.getNewTopic());
headers = Collections.singletonMap(KafkaHeaders.TOPIC, this.properties.getNewTopic());
for (int i = 0; i < 10; i++) {
toKafka.send(new GenericMessage<>("foo" + i, headers));
}
received = fromKafka.receive(10000);
count = 0;
while (received != null) {
System.out.println(received);
received = fromKafka.receive(++count < 10 ? 10000 : 1000);
}
context.close();
System.exit(0);
}
@Bean
@@ -103,8 +117,7 @@ public class Application {
public MessageHandler handler(KafkaTemplate<String, String> kafkaTemplate) {
KafkaProducerMessageHandler<String, String> handler =
new KafkaProducerMessageHandler<>(kafkaTemplate);
handler.setTopicExpression(new LiteralExpression(this.topic));
handler.setMessageKeyExpression(new LiteralExpression(this.messageKey));
handler.setMessageKeyExpression(new LiteralExpression(this.properties.getMessageKey()));
return handler;
}
@@ -120,7 +133,7 @@ public class Application {
public KafkaMessageListenerContainer<String, String> container(
ConsumerFactory<String, String> kafkaConsumerFactory) {
return new KafkaMessageListenerContainer<>(kafkaConsumerFactory,
new ContainerProperties(new TopicPartitionInitialOffset(this.topic, 0)));
new ContainerProperties(new TopicPartitionInitialOffset(this.properties.getTopic(), 0)));
}
@Bean
@@ -128,70 +141,33 @@ public class Application {
adapter(KafkaMessageListenerContainer<String, String> container) {
KafkaMessageDrivenChannelAdapter<String, String> kafkaMessageDrivenChannelAdapter =
new KafkaMessageDrivenChannelAdapter<>(container);
kafkaMessageDrivenChannelAdapter.setOutputChannel(received());
kafkaMessageDrivenChannelAdapter.setOutputChannel(fromKafka());
return kafkaMessageDrivenChannelAdapter;
}
@Bean
public PollableChannel received() {
public PollableChannel fromKafka() {
return new QueueChannel();
}
@Bean
public TopicCreator topicCreator() {
return new TopicCreator(this.topic, this.zookeeperConnect);
}
@Autowired
private IntegrationFlowContext flowContext;
public static class TopicCreator implements SmartLifecycle {
private final String topic;
private final String zkConnect;
private volatile boolean running;
public TopicCreator(String topic, String zkConnect) {
this.topic = topic;
this.zkConnect = zkConnect;
}
@Override
public void start() {
ZkUtils zkUtils = new ZkUtils(new ZkClient(this.zkConnect, 6000, 6000,
ZKStringSerializer$.MODULE$), null, false);
try {
AdminUtils.createTopic(zkUtils, topic, 1, 1, new Properties(), null);
}
catch (TopicExistsException e) {
// no-op
}
this.running = true;
}
@Override
public void stop() {
}
@Override
public boolean isRunning() {
return this.running;
}
@Override
public int getPhase() {
return Integer.MIN_VALUE;
}
@Override
public boolean isAutoStartup() {
return true;
}
@Override
public void stop(Runnable callback) {
callback.run();
}
@Autowired
private KafkaProperties kafkaProperties;
public void addAnotherListenerForTopics(String... topics) {
Map<String, Object> consumerProperties = kafkaProperties.buildConsumerProperties();
// change the group id so we don't revoke the other partitions.
consumerProperties.put(ConsumerConfig.GROUP_ID_CONFIG,
consumerProperties.get(ConsumerConfig.GROUP_ID_CONFIG) + "x");
IntegrationFlow flow =
IntegrationFlows
.from(Kafka.messageDrivenChannelAdapter(
new DefaultKafkaConsumerFactory<String, String>(consumerProperties), topics))
.channel("fromKafka")
.get();
this.flowContext.registration(flow).register();
}
}

View File

@@ -0,0 +1,61 @@
/*
* Copyright 2017 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.samples.kafka;
import org.springframework.boot.context.properties.ConfigurationProperties;
/**
* Properties for the kafka sample app.
*
* @author Gary Russell
* @since 5.0
*
*/
@ConfigurationProperties("kafka")
public class KafkaAppProperties {
private String topic;
private String newTopic;
private String messageKey;
public String getTopic() {
return this.topic;
}
public void setTopic(String topic) {
this.topic = topic;
}
public String getNewTopic() {
return this.newTopic;
}
public void setNewTopic(String newTopic) {
this.newTopic = newTopic;
}
public String getMessageKey() {
return this.messageKey;
}
public void setMessageKey(String messageKey) {
this.messageKey = messageKey;
}
}

View File

@@ -1,14 +1,13 @@
kafka:
zookeeper:
connect: localhost:2181
topic: si.topic
newTopic: si.new.topic
messageKey: si.key
spring:
kafka:
consumer:
group-id: siTestGroup
enable-auto-commit: true
auto-commit-interval: 100
auto-offset-reset: earliest
enable-auto-commit: false
value-deserializer: org.apache.kafka.common.serialization.StringDeserializer
key-deserializer: org.apache.kafka.common.serialization.StringDeserializer
producer: