From 330610fde9d08b4d9a22a4cca56698bd6ca1256f Mon Sep 17 00:00:00 2001 From: Soby Chacko Date: Thu, 3 Nov 2022 14:46:33 -0400 Subject: [PATCH] Address compile issues in Kafka binder tests --- .../streams/function/KafkaStreamsComponentBeansTests.java | 5 +++-- .../integration/KafkaStreamsBinderHealthIndicatorTests.java | 3 ++- 2 files changed, 5 insertions(+), 3 deletions(-) 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 c80a657b1..c2762b569 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 @@ -16,6 +16,7 @@ package org.springframework.cloud.stream.binder.kafka.streams.function; +import java.time.Duration; import java.util.Map; import java.util.concurrent.CountDownLatch; import java.util.concurrent.TimeUnit; @@ -169,7 +170,7 @@ public class KafkaStreamsComponentBeansTests { template.sendDefault("foobar"); template.setDefaultTopic("testBiFunctionComponent-in-1"); template.sendDefault("foobar"); - final ConsumerRecords records = KafkaTestUtils.getRecords(consumer2, 10_000, 2); + final ConsumerRecords records = KafkaTestUtils.getRecords(consumer2, Duration.ofSeconds(10), 2); assertThat(records.count()).isEqualTo(2); records.forEach(stringStringConsumerRecord -> assertThat(stringStringConsumerRecord.value().contains("foobar")).isTrue()); } @@ -259,7 +260,7 @@ public class KafkaStreamsComponentBeansTests { template.sendDefault("foobar"); template.setDefaultTopic("testCurriedFunctionWithFunctionTerminal-in-2"); template.sendDefault("foobar"); - final ConsumerRecords records = KafkaTestUtils.getRecords(consumer3, 10_000, 3); + final ConsumerRecords records = KafkaTestUtils.getRecords(consumer3, Duration.ofSeconds(10), 3); assertThat(records.count()).isEqualTo(3); records.forEach(stringStringConsumerRecord -> assertThat(stringStringConsumerRecord.value().contains("foobar")).isTrue()); } diff --git a/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/KafkaStreamsBinderHealthIndicatorTests.java b/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/KafkaStreamsBinderHealthIndicatorTests.java index 16308c673..6f30eda93 100644 --- a/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/KafkaStreamsBinderHealthIndicatorTests.java +++ b/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/KafkaStreamsBinderHealthIndicatorTests.java @@ -16,6 +16,7 @@ package org.springframework.cloud.stream.binder.kafka.streams.integration; +import java.time.Duration; import java.util.List; import java.util.Map; import java.util.concurrent.CompletableFuture; @@ -164,7 +165,7 @@ public class KafkaStreamsBinderHealthIndicatorTests { latch.await(5, TimeUnit.SECONDS); embeddedKafka.consumeFromEmbeddedTopics(consumer, topics); - KafkaTestUtils.getRecords(consumer, 1000); + KafkaTestUtils.getRecords(consumer, Duration.ofSeconds(1)); TimeUnit.SECONDS.sleep(5); checkHealth(context, expected);