diff --git a/kafka-streams-samples/kafka-streams-table-join/pom.xml b/kafka-streams-samples/kafka-streams-table-join/pom.xml index 363a1eb..04bce0d 100644 --- a/kafka-streams-samples/kafka-streams-table-join/pom.xml +++ b/kafka-streams-samples/kafka-streams-table-join/pom.xml @@ -11,23 +11,62 @@ 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 spring-cloud-stream-binder-kafka-streams - 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 + @@ -39,4 +78,55 @@ + + + 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-table-join/src/main/java/kafka/streams/table/join/KafkaStreamsTableJoin.java b/kafka-streams-samples/kafka-streams-table-join/src/main/java/kafka/streams/table/join/KafkaStreamsTableJoin.java index 2157453..f9d4a75 100644 --- a/kafka-streams-samples/kafka-streams-table-join/src/main/java/kafka/streams/table/join/KafkaStreamsTableJoin.java +++ b/kafka-streams-samples/kafka-streams-table-join/src/main/java/kafka/streams/table/join/KafkaStreamsTableJoin.java @@ -16,19 +16,17 @@ package kafka.streams.table.join; +import java.util.function.BiFunction; + 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.Joined; import org.apache.kafka.streams.kstream.KStream; import org.apache.kafka.streams.kstream.KTable; -import org.apache.kafka.streams.kstream.Serialized; 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.StreamListener; -import org.springframework.cloud.stream.binder.kafka.streams.annotations.KafkaStreamsProcessor; -import org.springframework.messaging.handler.annotation.SendTo; +import org.springframework.context.annotation.Bean; @SpringBootApplication public class KafkaStreamsTableJoin { @@ -37,33 +35,22 @@ public class KafkaStreamsTableJoin { SpringApplication.run(KafkaStreamsTableJoin.class, args); } - @EnableBinding(KStreamProcessorX.class) public static class KStreamToTableJoinApplication { + @Bean + public BiFunction, KTable, KStream> process() { - @StreamListener - @SendTo("output") - public KStream process(@Input("input") KStream userClicksStream, - @Input("inputTable") KTable userRegionsTable) { - - return userClicksStream + return (userClicksStream, userRegionsTable) -> userClicksStream .leftJoin(userRegionsTable, (clicks, region) -> new RegionWithClicks(region == null ? "UNKNOWN" : region, clicks), Joined.with(Serdes.String(), Serdes.Long(), null)) .map((user, regionWithClicks) -> new KeyValue<>(regionWithClicks.getRegion(), regionWithClicks.getClicks())) - .groupByKey(Serialized.with(Serdes.String(), Serdes.Long())) + .groupByKey(Grouped.with(Serdes.String(), Serdes.Long())) .reduce((firstClicks, secondClicks) -> firstClicks + secondClicks) .toStream(); } } - - interface KStreamProcessorX extends KafkaStreamsProcessor { - - @Input("inputTable") - KTable inputKTable(); - } - private static final class RegionWithClicks { private final String region; diff --git a/kafka-streams-samples/kafka-streams-table-join/src/main/java/kafka/streams/table/join/Producers.java b/kafka-streams-samples/kafka-streams-table-join/src/main/java/kafka/streams/table/join/Producers.java index 9513eeb..cb3f2a1 100644 --- a/kafka-streams-samples/kafka-streams-table-join/src/main/java/kafka/streams/table/join/Producers.java +++ b/kafka-streams-samples/kafka-streams-table-join/src/main/java/kafka/streams/table/join/Producers.java @@ -57,7 +57,7 @@ public class Producers { DefaultKafkaProducerFactory pf = new DefaultKafkaProducerFactory<>(props); KafkaTemplate template = new KafkaTemplate<>(pf, true); - template.setDefaultTopic("user-clicks3"); + template.setDefaultTopic("user-clicks"); for (KeyValue keyValue : userClicks) { template.sendDefault(keyValue.key, keyValue.value); diff --git a/kafka-streams-samples/kafka-streams-table-join/src/main/resources/application.yml b/kafka-streams-samples/kafka-streams-table-join/src/main/resources/application.yml index 7a1670f..a9de737 100644 --- a/kafka-streams-samples/kafka-streams-table-join/src/main/resources/application.yml +++ b/kafka-streams-samples/kafka-streams-table-join/src/main/resources/application.yml @@ -1,12 +1,11 @@ spring.application.name: stream-table-sample -spring.cloud.stream.bindings.input: - destination: user-clicks3 -spring.cloud.stream.bindings.inputTable: +spring.cloud.stream.bindings.process-in-0: + destination: user-clicks +spring.cloud.stream.bindings.process-in-1: destination: user-regions -spring.cloud.stream.bindings.output: +spring.cloud.stream.bindings.process-out-0: destination: output-topic spring.cloud.stream.kafka.streams.binder: - brokers: localhost configuration: default.key.serde: org.apache.kafka.common.serialization.Serdes$StringSerde default.value.serde: org.apache.kafka.common.serialization.Serdes$StringSerde