From 04296664e3da57b5e54278aaed4b1c667fe4beaf Mon Sep 17 00:00:00 2001 From: Artem Bilan Date: Thu, 8 Oct 2020 13:07:48 -0400 Subject: [PATCH] One more attempt for Redis Streams test --- .../inbound/ReactiveRedisStreamMessageProducerTests.java | 9 ++++++++- 1 file changed, 8 insertions(+), 1 deletion(-) diff --git a/spring-integration-redis/src/test/java/org/springframework/integration/redis/inbound/ReactiveRedisStreamMessageProducerTests.java b/spring-integration-redis/src/test/java/org/springframework/integration/redis/inbound/ReactiveRedisStreamMessageProducerTests.java index d9df4f975b..ae5af17703 100644 --- a/spring-integration-redis/src/test/java/org/springframework/integration/redis/inbound/ReactiveRedisStreamMessageProducerTests.java +++ b/spring-integration-redis/src/test/java/org/springframework/integration/redis/inbound/ReactiveRedisStreamMessageProducerTests.java @@ -186,6 +186,13 @@ public class ReactiveRedisStreamMessageProducerTests extends RedisAvailableTests Address address = new Address("Winterfell, Westeros"); Person person = new Person(address, "John Snow"); + this.template.opsForStream() + .createGroup(STREAM_KEY, this.redisStreamMessageProducer.getBeanName()) + .as(StepVerifier::create) + .assertNext(message -> assertThat(message).isEqualTo("OK")) + .thenCancel() + .verify(Duration.ofSeconds(10)); + this.redisStreamMessageProducer.setCreateConsumerGroup(true); this.redisStreamMessageProducer.setAutoAck(false); this.redisStreamMessageProducer.setConsumerName(CONSUMER); @@ -209,7 +216,7 @@ public class ReactiveRedisStreamMessageProducerTests extends RedisAvailableTests this.messageHandler.handleMessage(new GenericMessage<>(person)); - stepVerifier.verify(Duration.ofSeconds(20)); + stepVerifier.verify(Duration.ofSeconds(10)); await().until(() -> template.opsForStream()