From 89e1d9536345275d78ce4176cc8e25b36a28b035 Mon Sep 17 00:00:00 2001 From: xinhc Date: Thu, 11 Jan 2018 11:00:00 +0800 Subject: [PATCH] Add Kafka sample See gh-11597 --- spring-boot-samples/README.adoc | 3 + spring-boot-samples/pom.xml | 1 + .../spring-boot-sample-kafka/pom.xml | 53 ++++++++++++++++ .../src/main/java/sample/kafka/Consumer.java | 28 +++++++++ .../src/main/java/sample/kafka/Producer.java | 32 ++++++++++ .../sample/kafka/SampleKafkaApplication.java | 27 ++++++++ .../main/java/sample/kafka/SampleMessage.java | 53 ++++++++++++++++ .../src/main/resources/application.properties | 5 ++ .../kafka/SampleKafkaApplicationTests.java | 61 +++++++++++++++++++ 9 files changed, 263 insertions(+) create mode 100644 spring-boot-samples/spring-boot-sample-kafka/pom.xml create mode 100644 spring-boot-samples/spring-boot-sample-kafka/src/main/java/sample/kafka/Consumer.java create mode 100644 spring-boot-samples/spring-boot-sample-kafka/src/main/java/sample/kafka/Producer.java create mode 100644 spring-boot-samples/spring-boot-sample-kafka/src/main/java/sample/kafka/SampleKafkaApplication.java create mode 100644 spring-boot-samples/spring-boot-sample-kafka/src/main/java/sample/kafka/SampleMessage.java create mode 100644 spring-boot-samples/spring-boot-sample-kafka/src/main/resources/application.properties create mode 100644 spring-boot-samples/spring-boot-sample-kafka/src/test/java/sample/kafka/SampleKafkaApplicationTests.java diff --git a/spring-boot-samples/README.adoc b/spring-boot-samples/README.adoc index ab6b03a734..0e48ede4f8 100644 --- a/spring-boot-samples/README.adoc +++ b/spring-boot-samples/README.adoc @@ -230,3 +230,6 @@ The following sample applications are provided: | link:spring-boot-sample-xml[spring-boot-sample-xml] | Example show how Spring Boot can be mixed with traditional XML configuration (we generally recommend using Java `@Configuration` whenever possible + +| link:spring-boot-sample-kafka[spring-boot-sample-kafka] +| consumer and producer using Apache Kafka diff --git a/spring-boot-samples/pom.xml b/spring-boot-samples/pom.xml index 51f5dbed35..83792c9881 100644 --- a/spring-boot-samples/pom.xml +++ b/spring-boot-samples/pom.xml @@ -97,6 +97,7 @@ spring-boot-sample-websocket-undertow spring-boot-sample-webservices spring-boot-sample-xml + spring-boot-sample-kafka diff --git a/spring-boot-samples/spring-boot-sample-kafka/pom.xml b/spring-boot-samples/spring-boot-sample-kafka/pom.xml new file mode 100644 index 0000000000..901631faa6 --- /dev/null +++ b/spring-boot-samples/spring-boot-sample-kafka/pom.xml @@ -0,0 +1,53 @@ + + + 4.0.0 + + + spring-boot-samples + org.springframework.boot + ${revision} + + spring-boot-sample-kafka + Spring Boot Kafka Sample + Spring Boot Kafka Sample + + + ${basedir}/../.. + + + + + + org.springframework.boot + spring-boot-starter + + + org.springframework.kafka + spring-kafka + + + org.springframework.boot + spring-boot-starter-json + + + + org.springframework.boot + spring-boot-starter-test + test + + + org.springframework.kafka + spring-kafka-test + + + + + + + org.springframework.boot + spring-boot-maven-plugin + + + + diff --git a/spring-boot-samples/spring-boot-sample-kafka/src/main/java/sample/kafka/Consumer.java b/spring-boot-samples/spring-boot-sample-kafka/src/main/java/sample/kafka/Consumer.java new file mode 100644 index 0000000000..4e7e2a8b3d --- /dev/null +++ b/spring-boot-samples/spring-boot-sample-kafka/src/main/java/sample/kafka/Consumer.java @@ -0,0 +1,28 @@ +/* + * Copyright 2012-2018 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 sample.kafka; + +import org.springframework.kafka.annotation.KafkaListener; +import org.springframework.stereotype.Component; + +@Component +public class Consumer { + + @KafkaListener(topics = "myTopic") + public void processMessage(SampleMessage message) { + System.out.println("consumer has received message : [" + message + "]"); + } +} \ No newline at end of file diff --git a/spring-boot-samples/spring-boot-sample-kafka/src/main/java/sample/kafka/Producer.java b/spring-boot-samples/spring-boot-sample-kafka/src/main/java/sample/kafka/Producer.java new file mode 100644 index 0000000000..89c50e3533 --- /dev/null +++ b/spring-boot-samples/spring-boot-sample-kafka/src/main/java/sample/kafka/Producer.java @@ -0,0 +1,32 @@ +/* + * Copyright 2012-2018 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 sample.kafka; + +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.kafka.core.KafkaTemplate; +import org.springframework.stereotype.Component; + +@Component +public class Producer { + + @Autowired + private KafkaTemplate kafkaTemplate; + + public void send(SampleMessage message) { + kafkaTemplate.send("myTopic", message); + System.out.println("producer has sent message."); + } +} \ No newline at end of file diff --git a/spring-boot-samples/spring-boot-sample-kafka/src/main/java/sample/kafka/SampleKafkaApplication.java b/spring-boot-samples/spring-boot-sample-kafka/src/main/java/sample/kafka/SampleKafkaApplication.java new file mode 100644 index 0000000000..9769c06ace --- /dev/null +++ b/spring-boot-samples/spring-boot-sample-kafka/src/main/java/sample/kafka/SampleKafkaApplication.java @@ -0,0 +1,27 @@ +/* + * Copyright 2012-2018 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 sample.kafka; + +import org.springframework.boot.SpringApplication; +import org.springframework.boot.autoconfigure.SpringBootApplication; + +@SpringBootApplication +public class SampleKafkaApplication { + + public static void main(String[] args) { + SpringApplication.run(SampleKafkaApplication.class, args); + } +} diff --git a/spring-boot-samples/spring-boot-sample-kafka/src/main/java/sample/kafka/SampleMessage.java b/spring-boot-samples/spring-boot-sample-kafka/src/main/java/sample/kafka/SampleMessage.java new file mode 100644 index 0000000000..c0845fb491 --- /dev/null +++ b/spring-boot-samples/spring-boot-sample-kafka/src/main/java/sample/kafka/SampleMessage.java @@ -0,0 +1,53 @@ +/* + * Copyright 2012-2018 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 sample.kafka; + +public class SampleMessage { + private Integer id; + private String message; + + public Integer getId() { + return id; + } + + public void setId(Integer id) { + this.id = id; + } + + public String getMessage() { + return message; + } + + public void setMessage(String message) { + this.message = message; + } + + public SampleMessage() { + } + + public SampleMessage(Integer id, String message) { + this.id = id; + this.message = message; + } + + @Override + public String toString() { + return "SampleMessage{" + + "id=" + id + + ", message='" + message + '\'' + + '}'; + } +} diff --git a/spring-boot-samples/spring-boot-sample-kafka/src/main/resources/application.properties b/spring-boot-samples/spring-boot-sample-kafka/src/main/resources/application.properties new file mode 100644 index 0000000000..588dc54c56 --- /dev/null +++ b/spring-boot-samples/spring-boot-sample-kafka/src/main/resources/application.properties @@ -0,0 +1,5 @@ +spring.kafka.bootstrap-servers=localhost:9092 +spring.kafka.consumer.group-id=myGroup +spring.kafka.consumer.value-deserializer=org.springframework.kafka.support.serializer.JsonDeserializer +spring.kafka.Producer.value-serializer=org.springframework.kafka.support.serializer.JsonSerializer +spring.kafka.consumer.properties.spring.json.trusted.packages=sample.kafka \ No newline at end of file diff --git a/spring-boot-samples/spring-boot-sample-kafka/src/test/java/sample/kafka/SampleKafkaApplicationTests.java b/spring-boot-samples/spring-boot-sample-kafka/src/test/java/sample/kafka/SampleKafkaApplicationTests.java new file mode 100644 index 0000000000..3f91019946 --- /dev/null +++ b/spring-boot-samples/spring-boot-sample-kafka/src/test/java/sample/kafka/SampleKafkaApplicationTests.java @@ -0,0 +1,61 @@ +/* + * Copyright 2012-2018 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 sample.kafka; + +import org.junit.Rule; +import org.junit.Test; +import org.junit.runner.RunWith; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.boot.test.context.SpringBootTest; +import org.springframework.boot.test.rule.OutputCapture; +import org.springframework.kafka.test.rule.KafkaEmbedded; +import org.springframework.test.context.junit4.SpringRunner; + +import static org.assertj.core.api.Assertions.assertThat; + +/** + * Integration tests for demo application. + * + * @author hcxin + */ +@RunWith(SpringRunner.class) +@SpringBootTest +public class SampleKafkaApplicationTests { + + @Rule + public OutputCapture outputCapture = new OutputCapture(); + + @Autowired + private Producer producer; + + @Test + public void sendSimpleMessage() throws Exception { + initKafkaEmbedded(); + SampleMessage message = new SampleMessage(1, "Test message"); + producer.send(message); + Thread.sleep(1000L); + assertThat(this.outputCapture.toString().contains("Test message")).isTrue(); + } + + public void initKafkaEmbedded() throws Exception { + KafkaEmbedded embeddedKafka = new KafkaEmbedded(1, true); + embeddedKafka.setKafkaPorts(9092); + embeddedKafka.afterPropertiesSet(); + //Need 10s, waiting for the Kafka server start. + Thread.sleep(10000L); + + } +}