diff --git a/kafka-streams-samples/kafka-streams-dlq-sample/README.adoc b/kafka-streams-samples/kafka-streams-dlq-sample/README.adoc index ba6b06a..71fb6de 100644 --- a/kafka-streams-samples/kafka-streams-dlq-sample/README.adoc +++ b/kafka-streams-samples/kafka-streams-dlq-sample/README.adoc @@ -36,6 +36,5 @@ On the console producer, enter some text data. You will see that the messages produce deserialization errors and end up in the DLQ topic - words-count-dlq. You will not see any messages coming to the regular destination counts. -There is another yaml file provided (by-framework-decoding.yml). -Use that as application.yml to see how it works when the deserialization done by the framework. -In this case also, the messages on error appear in the DLQ topic. +In the application's configuration we set the value `Serde` as `Integer` and this is the reason why deserialization errors occur. +The application expects numerical data, but the producer sends regular text. diff --git a/kafka-streams-samples/kafka-streams-dlq-sample/pom.xml b/kafka-streams-samples/kafka-streams-dlq-sample/pom.xml index 75e1641..904ac42 100644 --- a/kafka-streams-samples/kafka-streams-dlq-sample/pom.xml +++ b/kafka-streams-samples/kafka-streams-dlq-sample/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-dlq-sample/src/main/java/kafka/streams/dlq/sample/KafkaStreamsDlqSample.java b/kafka-streams-samples/kafka-streams-dlq-sample/src/main/java/kafka/streams/dlq/sample/KafkaStreamsDlqSample.java index 52d3efd..cc75398 100644 --- a/kafka-streams-samples/kafka-streams-dlq-sample/src/main/java/kafka/streams/dlq/sample/KafkaStreamsDlqSample.java +++ b/kafka-streams-samples/kafka-streams-dlq-sample/src/main/java/kafka/streams/dlq/sample/KafkaStreamsDlqSample.java @@ -16,22 +16,19 @@ package kafka.streams.dlq.sample; -import org.apache.kafka.common.serialization.Serdes; -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.Serialized; -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.StreamListener; -import org.springframework.cloud.stream.binder.kafka.streams.annotations.KafkaStreamsProcessor; -import org.springframework.messaging.handler.annotation.SendTo; - import java.util.Arrays; import java.util.Date; +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.TimeWindows; +import org.springframework.boot.SpringApplication; +import org.springframework.boot.autoconfigure.SpringBootApplication; +import org.springframework.context.annotation.Bean; @SpringBootApplication public class KafkaStreamsDlqSample { @@ -40,21 +37,20 @@ public class KafkaStreamsDlqSample { SpringApplication.run(KafkaStreamsDlqSample.class, args); } - @EnableBinding(KafkaStreamsProcessor.class) public static class WordCountProcessorApplication { - @StreamListener("input") - @SendTo("output") - 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())) + .groupByKey(Grouped.with(Serdes.String(), Serdes.String())) .windowedBy(TimeWindows.of(5000)) .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())))); } } diff --git a/kafka-streams-samples/kafka-streams-dlq-sample/src/main/resources/application.yml b/kafka-streams-samples/kafka-streams-dlq-sample/src/main/resources/application.yml index 3d241aa..28f72d1 100644 --- a/kafka-streams-samples/kafka-streams-dlq-sample/src/main/resources/application.yml +++ b/kafka-streams-samples/kafka-streams-dlq-sample/src/main/resources/application.yml @@ -3,15 +3,13 @@ spring.cloud.stream.kafka.streams.binder: commit.interval.ms: 1000 default.key.serde: org.apache.kafka.common.serialization.Serdes$StringSerde default.value.serde: org.apache.kafka.common.serialization.Serdes$IntegerSerde - application.id: dlq-1 - brokers: localhost + application.id: dlq-demo-sample serdeError: sendToDlq -spring.cloud.stream.bindings.output: +spring.cloud.stream.bindings.process-out-0: destination: counts -spring.cloud.stream.bindings.input: +spring.cloud.stream.bindings.process-in-0: destination: words - group: group1 -spring.cloud.stream.kafka.streams.bindings.input.consumer: +spring.cloud.stream.kafka.streams.bindings.process-in-0.consumer: dlqName: words-count-dlq valueSerde: org.apache.kafka.common.serialization.Serdes$IntegerSerde diff --git a/kafka-streams-samples/kafka-streams-dlq-sample/src/main/resources/by-framework-decoding.yml b/kafka-streams-samples/kafka-streams-dlq-sample/src/main/resources/by-framework-decoding.yml deleted file mode 100644 index 1a1e780..0000000 --- a/kafka-streams-samples/kafka-streams-dlq-sample/src/main/resources/by-framework-decoding.yml +++ /dev/null @@ -1,26 +0,0 @@ -spring.cloud.stream.bindings.output.contentType: application/json -spring.cloud.stream.kafka.streams.binder: - brokers: localhost #192.168.99.100 - serdeError: sendToDlq - 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 - application.id: dlq-2 -spring.cloud.stream.bindings.output: - destination: counts - producer: - useNativeEncoding: false -spring.cloud.stream.bindings.input: - contentType: foo/bar - destination: words - group: group1 - consumer: - useNativeDecoding: false -spring.cloud.stream.kafka.streams.bindings.input.consumer.dlqName: words-count-dlq - - - - - - diff --git a/kafka-streams-samples/kafka-streams-dlq-sample/src/test/java/kafka/streams/dlq/sample/KafkaStreamsDlqExampleTests.java b/kafka-streams-samples/kafka-streams-dlq-sample/src/test/java/kafka/streams/dlq/sample/KafkaStreamsDlqExampleTests.java index cb31324..cd193de 100644 --- a/kafka-streams-samples/kafka-streams-dlq-sample/src/test/java/kafka/streams/dlq/sample/KafkaStreamsDlqExampleTests.java +++ b/kafka-streams-samples/kafka-streams-dlq-sample/src/test/java/kafka/streams/dlq/sample/KafkaStreamsDlqExampleTests.java @@ -1,18 +1,67 @@ package kafka.streams.dlq.sample; -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.ConsumerRecords; +import org.junit.AfterClass; +import org.junit.BeforeClass; +import org.junit.ClassRule; import org.junit.Test; import org.junit.runner.RunWith; import org.springframework.boot.test.context.SpringBootTest; +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 public class KafkaStreamsDlqExampleTests { + @ClassRule + public static EmbeddedKafkaRule embeddedKafkaRule = new EmbeddedKafkaRule(1, true, "words", "words-count-dlq"); + + private static EmbeddedKafkaBroker embeddedKafka = embeddedKafkaRule.getEmbeddedKafka(); + + private static Consumer consumer; + + @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.consumeFromAnEmbeddedTopic(consumer, "words-count-dlq"); + 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 - @Ignore - public void contextLoads() { + public void testKafkaStreamsWordCountProcessor() { + Map senderProps = KafkaTestUtils.producerProps(embeddedKafka); + DefaultKafkaProducerFactory pf = new DefaultKafkaProducerFactory<>(senderProps); + try { + KafkaTemplate template = new KafkaTemplate<>(pf, true); + template.setDefaultTopic("words"); + template.sendDefault("foobar"); + ConsumerRecords cr = KafkaTestUtils.getRecords(consumer); + assertThat(cr.count()).isGreaterThanOrEqualTo(1); + } + finally { + pf.destroy(); + } } }