Upgrade to Kafka 2.6.0
Closes gh-22731
This commit is contained in:
@@ -18,9 +18,14 @@ package org.springframework.boot.autoconfigure.kafka;
|
||||
|
||||
import java.util.concurrent.CountDownLatch;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
import java.util.regex.Pattern;
|
||||
|
||||
import org.apache.kafka.clients.admin.NewTopic;
|
||||
import org.apache.kafka.clients.producer.Producer;
|
||||
import org.apache.kafka.streams.StreamsBuilder;
|
||||
import org.apache.kafka.streams.kstream.KStream;
|
||||
import org.apache.kafka.streams.kstream.KTable;
|
||||
import org.apache.kafka.streams.kstream.Materialized;
|
||||
import org.junit.jupiter.api.AfterEach;
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
@@ -124,6 +129,12 @@ class KafkaAutoConfigurationIntegrationTests {
|
||||
@EnableKafkaStreams
|
||||
static class KafkaStreamsConfig {
|
||||
|
||||
@Bean
|
||||
public KTable<?, ?> table(StreamsBuilder builder) {
|
||||
KStream<Object, Object> stream = builder.stream(Pattern.compile("test"));
|
||||
return stream.groupByKey().count(Materialized.as("store"));
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
static class Listener {
|
||||
|
||||
Reference in New Issue
Block a user