diff --git a/kafka-streams-samples/kafka-streams-word-count/src/main/java/kafka/streams/word/count/KafkaStreamsWordCountApplication.java b/kafka-streams-samples/kafka-streams-word-count/src/main/java/kafka/streams/word/count/KafkaStreamsWordCountApplication.java index d8c4b1d..9d93ade 100644 --- a/kafka-streams-samples/kafka-streams-word-count/src/main/java/kafka/streams/word/count/KafkaStreamsWordCountApplication.java +++ b/kafka-streams-samples/kafka-streams-word-count/src/main/java/kafka/streams/word/count/KafkaStreamsWordCountApplication.java @@ -44,7 +44,7 @@ public class KafkaStreamsWordCountApplication { public static final String INPUT_TOPIC = "input"; public static final String OUTPUT_TOPIC = "output"; - public static final int WINDOW_SIZE_MS = 1000; + public static final int WINDOW_SIZE_MS = 30_000; @Bean public Function, KStream> process() { diff --git a/kafka-streams-samples/kafka-streams-word-count/src/main/resources/application.yml b/kafka-streams-samples/kafka-streams-word-count/src/main/resources/application.yml index cb34018..bbd9da2 100644 --- a/kafka-streams-samples/kafka-streams-word-count/src/main/resources/application.yml +++ b/kafka-streams-samples/kafka-streams-word-count/src/main/resources/application.yml @@ -2,8 +2,8 @@ spring.cloud.stream: function: definition: process bindings: - process-in-0: input - process-out-0: output + process-in-0: words + process-out-0: counts # bindings: # input: # destination: input diff --git a/kafka-streams-samples/kafka-streams-word-count/src/test/java/kafka/streams/word/count/WordCountProcessorApplicationTests.java b/kafka-streams-samples/kafka-streams-word-count/src/test/java/kafka/streams/word/count/WordCountProcessorApplicationTests.java index 17a19ea..4d79823 100644 --- a/kafka-streams-samples/kafka-streams-word-count/src/test/java/kafka/streams/word/count/WordCountProcessorApplicationTests.java +++ b/kafka-streams-samples/kafka-streams-word-count/src/test/java/kafka/streams/word/count/WordCountProcessorApplicationTests.java @@ -78,7 +78,7 @@ public class WordCountProcessorApplicationTests { KafkaStreamsWordCountApplication.WordCountProcessorApplication app = new KafkaStreamsWordCountApplication.WordCountProcessorApplication(); //final Function, KStream> process = app.process(); - final Function, KStream> process = null; + final Function, KStream> process = app.process(); final KStream output = process.apply(input);