diff --git a/kafka-streams-samples/kafka-streams-aggregate/pom.xml b/kafka-streams-samples/kafka-streams-aggregate/pom.xml
index 16ea1bc..9a5ba7e 100644
--- a/kafka-streams-samples/kafka-streams-aggregate/pom.xml
+++ b/kafka-streams-samples/kafka-streams-aggregate/pom.xml
@@ -11,25 +11,65 @@
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 +79,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-aggregate/src/main/java/kafka/streams/table/join/KafkaStreamsAggregateSample.java b/kafka-streams-samples/kafka-streams-aggregate/src/main/java/kafka/streams/table/join/KafkaStreamsAggregateSample.java
index 1140cc3..15dde85 100644
--- a/kafka-streams-samples/kafka-streams-aggregate/src/main/java/kafka/streams/table/join/KafkaStreamsAggregateSample.java
+++ b/kafka-streams-samples/kafka-streams-aggregate/src/main/java/kafka/streams/table/join/KafkaStreamsAggregateSample.java
@@ -16,23 +16,23 @@
package kafka.streams.table.join;
+import java.util.function.Consumer;
+
import com.fasterxml.jackson.databind.ObjectMapper;
import org.apache.kafka.common.serialization.Serde;
import org.apache.kafka.common.serialization.Serdes;
import org.apache.kafka.common.utils.Bytes;
+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.state.KeyValueStore;
import org.apache.kafka.streams.state.QueryableStoreTypes;
import org.apache.kafka.streams.state.ReadOnlyKeyValueStore;
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.StreamListener;
import org.springframework.cloud.stream.binder.kafka.streams.InteractiveQueryService;
+import org.springframework.context.annotation.Bean;
import org.springframework.kafka.support.serializer.JsonSerde;
import org.springframework.web.bind.annotation.RequestMapping;
import org.springframework.web.bind.annotation.RestController;
@@ -47,22 +47,23 @@ public class KafkaStreamsAggregateSample {
SpringApplication.run(KafkaStreamsAggregateSample.class, args);
}
- @EnableBinding(KafkaStreamsProcessorX.class)
public static class KafkaStreamsAggregateSampleApplication {
- @StreamListener("input")
- public void process(KStream input) {
+ @Bean
+ public Consumer> aggregate() {
+
ObjectMapper mapper = new ObjectMapper();
Serde domainEventSerde = new JsonSerde<>( DomainEvent.class, mapper );
- input
+ return input -> input
.groupBy(
(s, domainEvent) -> domainEvent.boardUuid,
- Serialized.with(null, domainEventSerde))
+ Grouped.with(null, domainEventSerde))
.aggregate(
String::new,
(s, domainEvent, board) -> board.concat(domainEvent.eventType),
- Materialized.>as("test-events-snapshots").withKeySerde(Serdes.String()).
+ Materialized.>as("test-events-snapshots")
+ .withKeySerde(Serdes.String()).
withValueSerde(Serdes.String())
);
}
@@ -80,9 +81,4 @@ public class KafkaStreamsAggregateSample {
}
}
- interface KafkaStreamsProcessorX {
-
- @Input("input")
- KStream, ?> input();
- }
}
diff --git a/kafka-streams-samples/kafka-streams-aggregate/src/main/resources/application.yml b/kafka-streams-samples/kafka-streams-aggregate/src/main/resources/application.yml
index a23fc16..3dfaa7e 100644
--- a/kafka-streams-samples/kafka-streams-aggregate/src/main/resources/application.yml
+++ b/kafka-streams-samples/kafka-streams-aggregate/src/main/resources/application.yml
@@ -1,9 +1,8 @@
spring.application.name: kafka-streams-aggregate-sample
-spring.cloud.stream.bindings.input:
+spring.cloud.stream.bindings.aggregate-in-0:
destination: foobar
spring.cloud.stream.kafka.streams.binder:
- brokers: localhost #192.168.99.100
configuration:
default.key.serde: org.apache.kafka.common.serialization.Serdes$StringSerde
- default.value.serde: org.apache.kafka.common.serialization.Serdes$BytesSerde
+ default.value.serde: org.apache.kafka.common.serialization.Serdes$StringSerde
commit.interval.ms: 1000
\ No newline at end of file
diff --git a/kafka-streams-samples/kafka-streams-aggregate/src/test/java/kafka/streams/table/join/KafkaStreamsAggregateSampleTests.java b/kafka-streams-samples/kafka-streams-aggregate/src/test/java/kafka/streams/table/join/KafkaStreamsAggregateSampleTests.java
index 5d03d67..e6ba52c 100644
--- a/kafka-streams-samples/kafka-streams-aggregate/src/test/java/kafka/streams/table/join/KafkaStreamsAggregateSampleTests.java
+++ b/kafka-streams-samples/kafka-streams-aggregate/src/test/java/kafka/streams/table/join/KafkaStreamsAggregateSampleTests.java
@@ -1,18 +1,97 @@
package kafka.streams.table.join;
-import org.junit.Ignore;
+import java.util.Map;
+
+import com.fasterxml.jackson.databind.ObjectMapper;
+import org.apache.kafka.clients.producer.ProducerConfig;
+import org.apache.kafka.common.serialization.Serde;
+import org.apache.kafka.common.serialization.StringSerializer;
+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.boot.web.server.LocalServerPort;
+import org.springframework.http.ResponseEntity;
+import org.springframework.kafka.config.StreamsBuilderFactoryBean;
+import org.springframework.kafka.core.DefaultKafkaProducerFactory;
+import org.springframework.kafka.core.KafkaTemplate;
+import org.springframework.kafka.support.serializer.JsonSerde;
+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 org.springframework.web.client.RestTemplate;
+
+import static org.assertj.core.api.Assertions.assertThat;
@RunWith(SpringRunner.class)
-@SpringBootTest
+@SpringBootTest(
+ webEnvironment = SpringBootTest.WebEnvironment.RANDOM_PORT)
public class KafkaStreamsAggregateSampleTests {
+ @ClassRule
+ public static EmbeddedKafkaRule embeddedKafkaRule = new EmbeddedKafkaRule(1, true, "foobar");
+
+ private static EmbeddedKafkaBroker embeddedKafka = embeddedKafkaRule.getEmbeddedKafka();
+
+ @Autowired
+ StreamsBuilderFactoryBean streamsBuilderFactoryBean;
+
+ @LocalServerPort
+ int randomServerPort;
+
+ @Before
+ public void before() {
+ streamsBuilderFactoryBean.setCloseTimeout(0);
+ }
+
+ @BeforeClass
+ public static void setUp() {
+ System.setProperty("spring.cloud.stream.kafka.streams.binder.brokers", embeddedKafka.getBrokersAsString());
+ }
+
+ @AfterClass
+ public static void tearDown() {
+ System.clearProperty("spring.cloud.stream.kafka.streams.binder.brokers");
+ }
+
@Test
- @Ignore
- public void contextLoads() {
+ public void testKafkaStreamsWordCountProcessor() throws Exception {
+ Map senderProps = KafkaTestUtils.producerProps(embeddedKafka);
+ ObjectMapper mapper = new ObjectMapper();
+ Serde domainEventSerde = new JsonSerde<>(DomainEvent.class, mapper);
+
+ senderProps.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
+ senderProps.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, domainEventSerde.serializer().getClass());
+
+ DefaultKafkaProducerFactory pf = new DefaultKafkaProducerFactory<>(senderProps);
+ try {
+
+
+ KafkaTemplate template = new KafkaTemplate<>(pf, true);
+ template.setDefaultTopic("foobar");
+
+ DomainEvent ddEvent = new DomainEvent();
+ ddEvent.setBoardUuid("12345");
+ ddEvent.setEventType("create-domain-event");
+
+ template.sendDefault("", ddEvent);
+ Thread.sleep(1000);
+ RestTemplate restTemplate = new RestTemplate();
+ String fooResourceUrl
+ = "http://localhost:" + randomServerPort + "/events";
+ ResponseEntity response
+ = restTemplate.getForEntity(fooResourceUrl, String.class);
+ assertThat(response.getBody()).contains("create-domain-event");
+ }
+ finally {
+ pf.destroy();
+ }
}
}