From 128cff6e27fe1dd8e0f344014b0bb5fa6c550e77 Mon Sep 17 00:00:00 2001 From: Christian Tzolov Date: Wed, 23 Sep 2020 22:25:01 +0200 Subject: [PATCH] Enable Kafka Producer/Consumer/Stream metrics - Change bindings input output to allows execution in SCDF. - Update to Boot 2.3.4 to ensure micrometer 1.5.5. --- .../kafka-streams-word-count/pom.xml | 34 +++++++++++++++++++ .../KafkaStreamsWordCountApplication.java | 3 +- .../src/main/resources/application.yml | 33 +++++++++++++----- .../WordCountProcessorApplicationTests.java | 4 ++- 4 files changed, 62 insertions(+), 12 deletions(-) diff --git a/kafka-streams-samples/kafka-streams-word-count/pom.xml b/kafka-streams-samples/kafka-streams-word-count/pom.xml index 5ab2f36..ecd5b10 100644 --- a/kafka-streams-samples/kafka-streams-word-count/pom.xml +++ b/kafka-streams-samples/kafka-streams-word-count/pom.xml @@ -33,6 +33,19 @@ + + io.micrometer + micrometer-registry-prometheus + + + io.micrometer + micrometer-registry-wavefront + + + io.micrometer.prometheus + prometheus-rsocket-spring + 1.2.1 + org.springframework.cloud spring-cloud-stream-binder-kafka-streams @@ -74,6 +87,27 @@ org.springframework.boot spring-boot-maven-plugin + + com.google.cloud.tools + jib-maven-plugin + 2.5.2 + + + springcloud/baseimage:1.0.0 + + + springcloudstream/${project.artifactId} + + 3.0.0-SNAPSHOT + + + + USE_CURRENT_TIMESTAMP + Docker + + + + 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 1369b79..d8c4b1d 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 @@ -27,7 +27,6 @@ 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; @@ -45,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 = 30000; + public static final int WINDOW_SIZE_MS = 1000; @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 848296b..cb34018 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 @@ -1,9 +1,15 @@ spring.cloud.stream: - bindings: - process-in-0: - destination: words - process-out-0: - destination: counts + function: + definition: process + bindings: + process-in-0: input + process-out-0: output +# bindings: +# input: +# destination: input +# output: +# destination: output + kafka: streams: binder: @@ -15,14 +21,23 @@ spring.cloud.stream: value.serde: org.apache.kafka.common.serialization.Serdes$StringSerde #Enable metrics management: + metrics: + export: + wavefront: + enabled: false + prometheus: + enabled: false + rsocket: + enabled: false endpoint: health: show-details: ALWAYS endpoints: web: exposure: - include: metrics,health + include: health,info,bindings +# include: metrics,health,info,bindings #Enable logging to debug for spring kafka config -logging: - level: - org.springframework.kafka.config: debug +#logging: +# level: +# org.springframework.kafka.config: debug 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 86d5faa..17a19ea 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 @@ -76,7 +76,9 @@ public class WordCountProcessorApplicationTests { final StreamsBuilder builder = new StreamsBuilder(); KStream input = builder.stream(INPUT_TOPIC, Consumed.with(nullSerde, stringSerde)); KafkaStreamsWordCountApplication.WordCountProcessorApplication app = new KafkaStreamsWordCountApplication.WordCountProcessorApplication(); - final Function, KStream> process = app.process(); + + //final Function, KStream> process = app.process(); + final Function, KStream> process = null; final KStream output = process.apply(input);