From 95775b5eba74e0de52a935ba505d2808826d5ded Mon Sep 17 00:00:00 2001 From: Soby Chacko Date: Fri, 8 Dec 2023 17:32:14 -0500 Subject: [PATCH] Kafka Streams binder test cleanup --- ...kaStreamsBinderWordCountFunctionTests.java | 2 +- .../KafkaStreamsComponentBeansTests.java | 1 + .../KafkaStreamsFunctionStateStoreTests.java | 21 +++++++++---------- .../function/KafkaStreamsRetryTests.java | 21 +++++++++++-------- .../StreamToTableJoinFunctionTests.java | 2 +- 5 files changed, 25 insertions(+), 22 deletions(-) diff --git a/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsBinderWordCountFunctionTests.java b/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsBinderWordCountFunctionTests.java index a23aca685..a9e1d3977 100644 --- a/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsBinderWordCountFunctionTests.java +++ b/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsBinderWordCountFunctionTests.java @@ -441,7 +441,7 @@ class KafkaStreamsBinderWordCountFunctionTests { value -> Arrays.asList(value.toLowerCase().split("\\W+"))) .map((key, value) -> new KeyValue<>(value, value)) .groupByKey(Grouped.with(Serdes.String(), Serdes.String())) - .windowedBy(TimeWindows.of(Duration.ofSeconds(5))).count(Materialized.as("foobar-WordCounts")) + .windowedBy(TimeWindows.ofSizeWithNoGrace(Duration.ofSeconds(5))).count(Materialized.as("foobar-WordCounts")) .toStream() .map((key, value) -> new KeyValue<>(null, null)); } diff --git a/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsComponentBeansTests.java b/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsComponentBeansTests.java index bd7192376..aa9595dcc 100644 --- a/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsComponentBeansTests.java +++ b/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsComponentBeansTests.java @@ -328,6 +328,7 @@ class KafkaStreamsComponentBeansTests { KStream[]> { @Override + @SuppressWarnings("unchecked") public KStream[] apply(KStream stringIntegerKStream) { return stringIntegerKStream.map((key, value) -> new KeyValue<>(key.toString(), value)) .split() diff --git a/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsFunctionStateStoreTests.java b/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsFunctionStateStoreTests.java index 2c8935b3e..3ebd883b3 100644 --- a/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsFunctionStateStoreTests.java +++ b/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsFunctionStateStoreTests.java @@ -22,9 +22,8 @@ import java.util.Map; import org.apache.kafka.common.serialization.Serdes; import org.apache.kafka.streams.kstream.KStream; import org.apache.kafka.streams.kstream.KTable; -import org.apache.kafka.streams.processor.Processor; -import org.apache.kafka.streams.processor.ProcessorContext; -import org.apache.kafka.streams.processor.ProcessorSupplier; +import org.apache.kafka.streams.processor.api.Processor; +import org.apache.kafka.streams.processor.api.Record; import org.apache.kafka.streams.state.KeyValueStore; import org.apache.kafka.streams.state.StoreBuilder; import org.apache.kafka.streams.state.Stores; @@ -121,16 +120,16 @@ class KafkaStreamsFunctionStateStoreTests { @Bean(name = "biConsumerBean") public java.util.function.BiConsumer, KStream> process() { return (input0, input1) -> - input0.process((ProcessorSupplier) () -> new Processor() { + input0.process(() -> new Processor() { + @Override - @SuppressWarnings("unchecked") - public void init(ProcessorContext context) { + public void init(org.apache.kafka.streams.processor.api.ProcessorContext context) { state1 = (KeyValueStore) context.getStateStore("my-store"); state2 = (WindowStore) context.getStateStore("other-store"); } @Override - public void process(Object key, String value) { + public void process(Record record) { processed1 = true; } @@ -144,16 +143,16 @@ class KafkaStreamsFunctionStateStoreTests { @Bean public java.util.function.Consumer> hello() { return input -> { - input.toStream().process(() -> new Processor() { + input.toStream().process(() -> new Processor() { + @Override - @SuppressWarnings("unchecked") - public void init(ProcessorContext context) { + public void init(org.apache.kafka.streams.processor.api.ProcessorContext context) { state3 = (KeyValueStore) context.getStateStore("my-store"); state4 = (WindowStore) context.getStateStore("other-store"); } @Override - public void process(Object key, String value) { + public void process(Record record) { processed2 = true; } diff --git a/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsRetryTests.java b/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsRetryTests.java index a7836fbb1..42794bed0 100644 --- a/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsRetryTests.java +++ b/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsRetryTests.java @@ -24,8 +24,8 @@ import java.util.function.BiConsumer; import org.apache.kafka.streams.kstream.GlobalKTable; import org.apache.kafka.streams.kstream.KStream; import org.apache.kafka.streams.kstream.KTable; -import org.apache.kafka.streams.processor.Processor; -import org.apache.kafka.streams.processor.ProcessorContext; +import org.apache.kafka.streams.processor.api.Processor; +import org.apache.kafka.streams.processor.api.Record; import org.junit.jupiter.api.Test; import org.springframework.beans.factory.NoSuchBeanDefinitionException; @@ -143,13 +143,15 @@ class KafkaStreamsRetryTests { public java.util.function.Consumer> process(@Lazy @Qualifier("process-in-0-RetryTemplate") RetryTemplate retryTemplate) { return input -> input - .process(() -> new Processor() { + .process(() -> new Processor() { + @Override - public void init(ProcessorContext processorContext) { + public void init(org.apache.kafka.streams.processor.api.ProcessorContext context) { + Processor.super.init(context); } @Override - public void process(Object o, String s) { + public void process(Record record) { retryTemplate.execute(context -> { LATCH1.countDown(); throw new RuntimeException(); @@ -191,18 +193,19 @@ class KafkaStreamsRetryTests { public java.util.function.Consumer> process() { return input -> input - .process(() -> new Processor() { + .process(() -> new Processor() { + @Override - public void init(ProcessorContext processorContext) { + public void init(org.apache.kafka.streams.processor.api.ProcessorContext context) { + Processor.super.init(context); } @Override - public void process(Object o, String s) { + public void process(Record record) { fooRetryTemplate().execute(context -> { LATCH2.countDown(); throw new RuntimeException(); }); - } @Override diff --git a/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/function/StreamToTableJoinFunctionTests.java b/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/function/StreamToTableJoinFunctionTests.java index 5c8186e96..f3f9f9021 100644 --- a/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/function/StreamToTableJoinFunctionTests.java +++ b/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/function/StreamToTableJoinFunctionTests.java @@ -562,7 +562,7 @@ class StreamToTableJoinFunctionTests { return (input1Stream, input2Stream) -> input1Stream .join(input2Stream, (event1, event2) -> null, - JoinWindows.of(Duration.ofMillis(5)), + JoinWindows.ofTimeDifferenceWithNoGrace(Duration.ofMillis(5)), StreamJoined.with( Serdes.String(), Serdes.String(),