From 2703ab2c8d4dd3bca9c82d63ce8664332894b528 Mon Sep 17 00:00:00 2001 From: Soby Chacko Date: Mon, 28 Oct 2019 19:43:09 -0400 Subject: [PATCH] Update kafka-streams-message-channel sample --- .../kafka-streams-message-channel/pom.xml | 100 +++++++++++++++++- .../KafkaStreamsWordCountApplication.java | 56 ++++------ .../src/main/resources/application.yml | 17 ++- 3 files changed, 121 insertions(+), 52 deletions(-) diff --git a/kafka-streams-samples/kafka-streams-message-channel/pom.xml b/kafka-streams-samples/kafka-streams-message-channel/pom.xml index 1afabb6..6dc5070 100644 --- a/kafka-streams-samples/kafka-streams-message-channel/pom.xml +++ b/kafka-streams-samples/kafka-streams-message-channel/pom.xml @@ -11,12 +11,28 @@ Demo project for Spring Boot - io.spring.cloud.stream.sample - spring-cloud-stream-samples-parent - 0.0.1-SNAPSHOT - ../.. + org.springframework.boot + spring-boot-starter-parent + 2.2.0.RELEASE + + + Hoxton.BUILD-SNAPSHOT + + + + + + org.springframework.cloud + spring-cloud-dependencies + ${spring-cloud.version} + pom + import + + + + org.springframework.cloud @@ -26,14 +42,38 @@ org.springframework.cloud spring-cloud-stream-binder-kafka - org.springframework.boot spring-boot-starter-test test + + org.springframework.kafka + spring-kafka-test + test + + + org.apache.kafka + kafka-streams-test-utils + ${kafka.version} + test + + + + org.springframework.boot + spring-boot-starter-actuator + + + org.springframework.boot + spring-boot-starter + + + org.springframework.boot + spring-boot-starter-web + + @@ -43,4 +83,54 @@ + + + spring-snapshots + Spring Snapshots + https://repo.spring.io/libs-snapshot-local + + true + + + false + + + + spring-milestones + Spring Milestones + https://repo.spring.io/libs-milestone-local + + false + + + + + + spring-snapshots + Spring Snapshots + https://repo.spring.io/libs-snapshot-local + + true + + + false + + + + spring-milestones + Spring Milestones + https://repo.spring.io/libs-milestone-local + + false + + + + spring-releases + Spring Releases + https://repo.spring.io/libs-release-local + + false + + + diff --git a/kafka-streams-samples/kafka-streams-message-channel/src/main/java/kafka/streams/message/channel/KafkaStreamsWordCountApplication.java b/kafka-streams-samples/kafka-streams-message-channel/src/main/java/kafka/streams/message/channel/KafkaStreamsWordCountApplication.java index 53c4b8f..c0db086 100644 --- a/kafka-streams-samples/kafka-streams-message-channel/src/main/java/kafka/streams/message/channel/KafkaStreamsWordCountApplication.java +++ b/kafka-streams-samples/kafka-streams-message-channel/src/main/java/kafka/streams/message/channel/KafkaStreamsWordCountApplication.java @@ -16,23 +16,21 @@ package kafka.streams.message.channel; +import java.time.Duration; +import java.util.Arrays; +import java.util.Date; +import java.util.function.Consumer; +import java.util.function.Function; + import org.apache.kafka.common.serialization.Serdes; import org.apache.kafka.streams.KeyValue; +import org.apache.kafka.streams.kstream.Grouped; import org.apache.kafka.streams.kstream.KStream; import org.apache.kafka.streams.kstream.Materialized; -import org.apache.kafka.streams.kstream.Serialized; import org.apache.kafka.streams.kstream.TimeWindows; import org.springframework.boot.SpringApplication; import org.springframework.boot.autoconfigure.SpringBootApplication; -import org.springframework.cloud.stream.annotation.EnableBinding; -import org.springframework.cloud.stream.annotation.Input; -import org.springframework.cloud.stream.annotation.Output; -import org.springframework.cloud.stream.annotation.StreamListener; -import org.springframework.messaging.SubscribableChannel; -import org.springframework.messaging.handler.annotation.SendTo; - -import java.util.Arrays; -import java.util.Date; +import org.springframework.context.annotation.Bean; @SpringBootApplication public class KafkaStreamsWordCountApplication { @@ -41,44 +39,26 @@ public class KafkaStreamsWordCountApplication { SpringApplication.run(KafkaStreamsWordCountApplication.class, args); } - @EnableBinding(MultipleProcessor.class) public static class WordCountProcessorApplication { - @StreamListener("binding2") - @SendTo("singleOutput") - public KStream process(KStream input) { + @Bean + public Function, KStream> process() { - return input + return input -> input .flatMapValues(value -> Arrays.asList(value.toLowerCase().split("\\W+"))) .map((key, value) -> new KeyValue<>(value, value)) - .groupByKey(Serialized.with(Serdes.String(), Serdes.String())) - .windowedBy(TimeWindows.of(60_000)) + .groupByKey(Grouped.with(Serdes.String(), Serdes.String())) + .windowedBy(TimeWindows.of(Duration.ofSeconds(60))) .count(Materialized.as("WordCounts-1")) .toStream() - .map((key, value) -> new KeyValue<>(null, new WordCount(key.key(), value, new Date(key.window().start()), new Date(key.window().end())))); + .map((key, value) -> new KeyValue<>(null, + new WordCount(key.key(), value, new Date(key.window().start()), new Date(key.window().end())))); } - @StreamListener("binding1") - public void sink(String input) { - System.out.println("FOOBAR -- " + input); + @Bean + public Consumer sink() { + return s -> System.out.println("FOOBAR -- " + s); } - - } - - interface MultipleProcessor { - - String BINDING_1 = "binding1"; - String BINDING_2 = "binding2"; - String OUTPUT = "singleOutput"; - - @Input(BINDING_1) - SubscribableChannel binding1(); - - @Input(BINDING_2) - KStream binding2(); - - @Output(OUTPUT) - KStream singleOutput(); } static class WordCount { diff --git a/kafka-streams-samples/kafka-streams-message-channel/src/main/resources/application.yml b/kafka-streams-samples/kafka-streams-message-channel/src/main/resources/application.yml index d893d6a..539cd01 100644 --- a/kafka-streams-samples/kafka-streams-message-channel/src/main/resources/application.yml +++ b/kafka-streams-samples/kafka-streams-message-channel/src/main/resources/application.yml @@ -1,15 +1,14 @@ -spring.cloud.stream.kafka.streams.binder.configuration.commit.interval.ms: 1000 -spring.cloud.stream.kafka.streams.binder.configuration: +spring.cloud.stream: + function: + definition: sink;process + kafka.streams.binder.configuration: + commit.interval.ms: 1000 default.key.serde: org.apache.kafka.common.serialization.Serdes$StringSerde default.value.serde: org.apache.kafka.common.serialization.Serdes$StringSerde -spring.cloud.stream.bindings.singleOutput: +spring.cloud.stream.bindings.process-out-0: destination: counts -spring.cloud.stream.bindings.binding2: +spring.cloud.stream.bindings.process-in-0: destination: words -spring.cloud.stream.bindings.binding1: +spring.cloud.stream.bindings.sink-in-0: destination: words -spring.cloud.stream.kafka.streams.binder: - brokers: localhost spring.application.name: kafka-streams-message-channel-sample - -