diff --git a/kafka-streams-samples/kafka-streams-branching/README.adoc b/kafka-streams-samples/kafka-streams-branching/README.adoc index 1bb9692..5029fcc 100644 --- a/kafka-streams-samples/kafka-streams-branching/README.adoc +++ b/kafka-streams-samples/kafka-streams-branching/README.adoc @@ -4,10 +4,7 @@ This is an example of a Spring Cloud Stream processor using Kafka Streams branch The example is based on the word count application from the https://github.com/confluentinc/examples/blob/3.2.x/kafka-streams/src/main/java/io/confluent/examples/streams/WordCountLambdaExample.java[reference documentation]. It uses a single input and 3 output destinations. -In essence, the application receives text messages from an input topic, filter them by language (Englihs, French, Spanish and ignoring the rest), and computes word occurrence counts in a configurable time window and report that in the output topics. -This sample uses lambda expressions and thus requires Java 8+. - -By default native decoding and encoding are disabled and this means that any deserializaion on inbound and serialization on outbound is performed by the Binder using the configured content types. +In essence, the application receives text messages from an input topic, filter them by language (English, French, Spanish and ignoring the rest), and computes word occurrence counts in a configurable time window and report that in the output topics. === Running the app: @@ -17,7 +14,7 @@ Go to the root of the repository and do: `./mvnw clean package` -`java -jar target/kafka-streams-branching-0.0.1-SNAPSHOT.jar --spring.cloud.stream.kafka.streams.timeWindow.length=60000` +`java -jar target/kafka-streams-branching-0.0.1-SNAPSHOT.jar` Issue the following commands: diff --git a/kafka-streams-samples/kafka-streams-branching/pom.xml b/kafka-streams-samples/kafka-streams-branching/pom.xml index 003c24c..96a2c57 100644 --- a/kafka-streams-samples/kafka-streams-branching/pom.xml +++ b/kafka-streams-samples/kafka-streams-branching/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 @@ -27,6 +43,30 @@ 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 + @@ -37,4 +77,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-branching/src/main/java/kafka/streams/branching/KafkaStreamsBranchingSample.java b/kafka-streams-samples/kafka-streams-branching/src/main/java/kafka/streams/branching/KafkaStreamsBranchingSample.java index 816a426..5d3e87f 100644 --- a/kafka-streams-samples/kafka-streams-branching/src/main/java/kafka/streams/branching/KafkaStreamsBranchingSample.java +++ b/kafka-streams-samples/kafka-streams-branching/src/main/java/kafka/streams/branching/KafkaStreamsBranchingSample.java @@ -16,22 +16,19 @@ package kafka.streams.branching; +import java.time.Duration; +import java.util.Arrays; +import java.util.Date; +import java.util.function.Function; + import org.apache.kafka.streams.KeyValue; import org.apache.kafka.streams.kstream.KStream; import org.apache.kafka.streams.kstream.Materialized; import org.apache.kafka.streams.kstream.Predicate; import org.apache.kafka.streams.kstream.TimeWindows; -import org.springframework.beans.factory.annotation.Autowired; 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.handler.annotation.SendTo; - -import java.util.Arrays; -import java.util.Date; +import org.springframework.context.annotation.Bean; @SpringBootApplication public class KafkaStreamsBranchingSample { @@ -40,44 +37,28 @@ public class KafkaStreamsBranchingSample { SpringApplication.run(KafkaStreamsBranchingSample.class, args); } - @EnableBinding(KStreamProcessorX.class) public static class WordCountProcessorApplication { - @StreamListener("input") - @SendTo({"output1","output2","output3"}) + @Bean @SuppressWarnings("unchecked") - public KStream[] process(KStream input) { + public Function, KStream[]> process() { Predicate isEnglish = (k, v) -> v.word.equals("english"); Predicate isFrench = (k, v) -> v.word.equals("french"); Predicate isSpanish = (k, v) -> v.word.equals("spanish"); - return input + return input -> input .flatMapValues(value -> Arrays.asList(value.toLowerCase().split("\\W+"))) .groupBy((key, value) -> value) - .windowedBy(TimeWindows.of(60_000)) + .windowedBy(TimeWindows.of(Duration.ofSeconds(6))) .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())))) .branch(isEnglish, isFrench, isSpanish); } } - interface KStreamProcessorX { - - @Input("input") - KStream input(); - - @Output("output1") - KStream output1(); - - @Output("output2") - KStream output2(); - - @Output("output3") - KStream output3(); - } - static class WordCount { private String word; diff --git a/kafka-streams-samples/kafka-streams-branching/src/main/resources/application.yml b/kafka-streams-samples/kafka-streams-branching/src/main/resources/application.yml index f396270..cea8bf2 100644 --- a/kafka-streams-samples/kafka-streams-branching/src/main/resources/application.yml +++ b/kafka-streams-samples/kafka-streams-branching/src/main/resources/application.yml @@ -2,14 +2,12 @@ spring.cloud.stream.kafka.streams.binder.configuration.commit.interval.ms: 1000 spring.cloud.stream.kafka.streams.binder.configuration: default.key.serde: org.apache.kafka.common.serialization.Serdes$StringSerde default.value.serde: org.apache.kafka.common.serialization.Serdes$StringSerde -spring.cloud.stream.bindings.output1: +spring.cloud.stream.bindings.process-out-0: destination: english-counts -spring.cloud.stream.bindings.output2: +spring.cloud.stream.bindings.process-out-1: destination: french-counts -spring.cloud.stream.bindings.output3: +spring.cloud.stream.bindings.process-out-2: destination: spanish-counts -spring.cloud.stream.bindings.input: +spring.cloud.stream.bindings.process-in-0: destination: words -spring.cloud.stream.kafka.streams.binder: - brokers: localhost #192.168.99.100 #localhost spring.application.name: kafka-streams-branching-sample \ No newline at end of file diff --git a/kafka-streams-samples/kafka-streams-branching/src/test/java/kafka/streams/branching/KafkaStreamsBranchingSampleTests.java b/kafka-streams-samples/kafka-streams-branching/src/test/java/kafka/streams/branching/KafkaStreamsBranchingSampleTests.java index 5eb6987..c9c9b5f 100644 --- a/kafka-streams-samples/kafka-streams-branching/src/test/java/kafka/streams/branching/KafkaStreamsBranchingSampleTests.java +++ b/kafka-streams-samples/kafka-streams-branching/src/test/java/kafka/streams/branching/KafkaStreamsBranchingSampleTests.java @@ -16,19 +16,88 @@ package kafka.streams.branching; -import org.junit.Ignore; +import java.util.Map; + +import org.apache.kafka.clients.consumer.Consumer; +import org.apache.kafka.clients.consumer.ConsumerConfig; +import org.apache.kafka.clients.consumer.ConsumerRecord; +import org.junit.AfterClass; +import org.junit.Before; +import org.junit.BeforeClass; +import org.junit.ClassRule; 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.kafka.config.StreamsBuilderFactoryBean; +import org.springframework.kafka.core.DefaultKafkaConsumerFactory; +import org.springframework.kafka.core.DefaultKafkaProducerFactory; +import org.springframework.kafka.core.KafkaTemplate; +import org.springframework.kafka.test.EmbeddedKafkaBroker; +import org.springframework.kafka.test.rule.EmbeddedKafkaRule; +import org.springframework.kafka.test.utils.KafkaTestUtils; import org.springframework.test.context.junit4.SpringRunner; +import static org.assertj.core.api.Assertions.assertThat; + @RunWith(SpringRunner.class) -@SpringBootTest +@SpringBootTest( + webEnvironment = SpringBootTest.WebEnvironment.NONE) public class KafkaStreamsBranchingSampleTests { - @Test - @Ignore - public void contextLoads() { + @ClassRule + public static EmbeddedKafkaRule embeddedKafkaRule = new EmbeddedKafkaRule(1, true, "words", + "english-counts", "french-counts", "spanish-counts"); + + private static EmbeddedKafkaBroker embeddedKafka = embeddedKafkaRule.getEmbeddedKafka(); + + private static Consumer consumer; + + @Autowired + StreamsBuilderFactoryBean streamsBuilderFactoryBean; + + @Before + public void before() { + streamsBuilderFactoryBean.setCloseTimeout(0); } -} + @BeforeClass + public static void setUp() { + Map consumerProps = KafkaTestUtils.consumerProps("group", "false", embeddedKafka); + consumerProps.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest"); + DefaultKafkaConsumerFactory cf = new DefaultKafkaConsumerFactory<>(consumerProps); + consumer = cf.createConsumer(); + embeddedKafka.consumeFromEmbeddedTopics(consumer, "english-counts", "french-counts", "spanish-counts"); + System.setProperty("spring.cloud.stream.kafka.streams.binder.brokers", embeddedKafka.getBrokersAsString()); + } + + @AfterClass + public static void tearDown() { + consumer.close(); + System.clearProperty("spring.cloud.stream.kafka.streams.binder.brokers"); + } + + @Test + public void testKafkaStreamsWordCountProcessor() throws InterruptedException { + Map senderProps = KafkaTestUtils.producerProps(embeddedKafka); + DefaultKafkaProducerFactory pf = new DefaultKafkaProducerFactory<>(senderProps); + try { + KafkaTemplate template = new KafkaTemplate<>(pf, true); + template.setDefaultTopic("words"); + template.sendDefault("english"); + template.sendDefault("french"); + template.sendDefault("spanish"); + Thread.sleep(2000); + ConsumerRecord cr = KafkaTestUtils.getSingleRecord(consumer, "english-counts", 5000); + assertThat(cr.value().contains("english")).isTrue(); + cr = KafkaTestUtils.getSingleRecord(consumer, "french-counts", 5000); + assertThat(cr.value().contains("french")).isTrue(); + cr = KafkaTestUtils.getSingleRecord(consumer, "spanish-counts", 5000); + assertThat(cr.value().contains("spanish")).isTrue(); + } + finally { + pf.destroy(); + } + } + +} \ No newline at end of file