diff --git a/spring-integration-redis/src/main/java/org/springframework/integration/redis/inbound/ReactiveRedisStreamMessageProducer.java b/spring-integration-redis/src/main/java/org/springframework/integration/redis/inbound/ReactiveRedisStreamMessageProducer.java index b1281b9369..6ab5a15e02 100644 --- a/spring-integration-redis/src/main/java/org/springframework/integration/redis/inbound/ReactiveRedisStreamMessageProducer.java +++ b/spring-integration-redis/src/main/java/org/springframework/integration/redis/inbound/ReactiveRedisStreamMessageProducer.java @@ -44,7 +44,7 @@ import reactor.core.publisher.Mono; * output channel. * By default this adapter reads message as a standalone client {@code XREAD} (Redis command) but can be switched to a * Consumer Group feature {@code XREADGROUP} by setting {@link #consumerName} field. - * By default the Consumer Group name is an id of this bean {@link #getBeanName()}. + * By default the Consumer Group name is the id of this bean {@link #getBeanName()}. * * @author Attoumane Ahamadi * @author Artem Bilan @@ -195,7 +195,7 @@ public class ReactiveRedisStreamMessageProducer extends MessageProducerSupport { Consumer consumer = Consumer.from(this.consumerGroup, this.consumerName); // NOSONAR if (offset.getOffset().equals(ReadOffset.latest())) { - // for consumer group offset id should be equal '>' + // for consumer group offset id should be equal to '>' offset = StreamOffset.create(this.streamKey, ReadOffset.lastConsumed()); } 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 f6b1cbfe9b..2e8c208749 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 @@ -38,6 +38,7 @@ import org.springframework.integration.redis.outbound.ReactiveRedisStreamMessage import org.springframework.integration.redis.rules.RedisAvailable; import org.springframework.integration.redis.rules.RedisAvailableRule; import org.springframework.integration.redis.rules.RedisAvailableTests; +import org.springframework.integration.redis.support.RedisHeaders; import org.springframework.integration.redis.util.Address; import org.springframework.integration.redis.util.Person; import org.springframework.messaging.support.GenericMessage; @@ -118,7 +119,11 @@ public class ReactiveRedisStreamMessageProducerTests extends RedisAvailableTests Flux.from(this.fluxMessageChannel) .as(StepVerifier::create) - .assertNext(message -> assertThat(message.getPayload()).isEqualTo(person)) + .assertNext(message -> { + assertThat(message.getPayload()).isEqualTo(person); + assertThat(message.getHeaders()).containsKeys(RedisHeaders.STREAM_KEY, + RedisHeaders.STREAM_MESSAGE_ID); + }) .thenCancel() .verify(Duration.ofSeconds(10)); } @@ -126,7 +131,27 @@ public class ReactiveRedisStreamMessageProducerTests extends RedisAvailableTests @Test @RedisAvailable public void testReadingMessageAsConsumerInConsumerGroup() { - //TODO find why the test above does not execute before implementing this one + Address address = new Address("Winterfell, Westeros"); + Person person = new Person(address, "John Snow"); + + this.template.opsForStream().createGroup(STREAM_KEY, this.redisStreamMessageProducer.getBeanName()) + .subscribe(); + + this.redisStreamMessageProducer.setCreateConsumerGroup(false); + this.redisStreamMessageProducer.setConsumerName(CONSUMER); + this.redisStreamMessageProducer.afterPropertiesSet(); + this.redisStreamMessageProducer.start(); + + this.messageHandler.handleMessage(new GenericMessage<>(person)); + + Flux.from(this.fluxMessageChannel) + .as(StepVerifier::create) + .assertNext(message -> { + assertThat(message.getPayload()).isEqualTo(person); + assertThat(message.getHeaders()).containsKeys(RedisHeaders.CONSUMER_GROUP, RedisHeaders.CONSUMER); + }) + .thenCancel() + .verify(Duration.ofSeconds(10)); } @Configuration