diff --git a/docs/src/main/asciidoc/overview.adoc b/docs/src/main/asciidoc/overview.adoc index 62485c7d1..bcd106ba1 100644 --- a/docs/src/main/asciidoc/overview.adoc +++ b/docs/src/main/asciidoc/overview.adoc @@ -438,7 +438,7 @@ Environment streamEnv() { ListenerContainerCustomizer customizer() { return (cont, dest, group) -> { StreamListenerContainer container = (StreamListenerContainer) cont; - container.setConsumerCustomizer(builder -> { + container.setConsumerCustomizer((name, builder) -> { builder.offset(OffsetSpecification.first()); }); // ... @@ -447,7 +447,10 @@ ListenerContainerCustomizer customizer() { ---- ==== -The stream name (for the purpose of offset tracking) is set to the binding `destination + '.' + group`. +The `name` argument passed to the customizer is `destination + '.' + group + '.container'`. + +The stream `name()` (for the purpose of offset tracking) is set to the binding `destination + '.' + group`. +It can be changed using a `ConsumerCustomizer` shown above. If you decide to use manual offset tracking, the `Context` is available as a message header: ==== diff --git a/spring-cloud-stream-binder-rabbit/pom.xml b/spring-cloud-stream-binder-rabbit/pom.xml index 45875901a..197d5c7bb 100644 --- a/spring-cloud-stream-binder-rabbit/pom.xml +++ b/spring-cloud-stream-binder-rabbit/pom.xml @@ -58,7 +58,7 @@ org.springframework.amqp spring-rabbit-stream - 2.4.0-M1 + 2.4.0-SNAPSHOT true diff --git a/spring-cloud-stream-binder-rabbit/src/main/java/org/springframework/cloud/stream/binder/rabbit/RabbitMessageChannelBinder.java b/spring-cloud-stream-binder-rabbit/src/main/java/org/springframework/cloud/stream/binder/rabbit/RabbitMessageChannelBinder.java index 99810408b..344334763 100644 --- a/spring-cloud-stream-binder-rabbit/src/main/java/org/springframework/cloud/stream/binder/rabbit/RabbitMessageChannelBinder.java +++ b/spring-cloud-stream-binder-rabbit/src/main/java/org/springframework/cloud/stream/binder/rabbit/RabbitMessageChannelBinder.java @@ -511,6 +511,7 @@ public class RabbitMessageChannelBinder extends AbstractMessageListenerContainer listenerContainer = directContainer ? new DirectMessageListenerContainer(this.connectionFactory) : new SimpleMessageListenerContainer(this.connectionFactory); + listenerContainer.setBeanName(consumerDestination.getName() + "." + group + ".container"); listenerContainer .setAcknowledgeMode(extension.getAcknowledgeMode()); listenerContainer.setChannelTransacted(extension.isTransacted()); diff --git a/spring-cloud-stream-binder-rabbit/src/main/java/org/springframework/cloud/stream/binder/rabbit/StreamContainerUtils.java b/spring-cloud-stream-binder-rabbit/src/main/java/org/springframework/cloud/stream/binder/rabbit/StreamContainerUtils.java index d44845f49..1163a07e2 100644 --- a/spring-cloud-stream-binder-rabbit/src/main/java/org/springframework/cloud/stream/binder/rabbit/StreamContainerUtils.java +++ b/spring-cloud-stream-binder-rabbit/src/main/java/org/springframework/cloud/stream/binder/rabbit/StreamContainerUtils.java @@ -20,11 +20,9 @@ import java.nio.charset.Charset; import java.nio.charset.StandardCharsets; import java.util.Map; import java.util.UUID; -import java.util.function.Consumer; import java.util.function.Supplier; import com.rabbitmq.stream.Codec; -import com.rabbitmq.stream.ConsumerBuilder; import com.rabbitmq.stream.Environment; import com.rabbitmq.stream.MessageBuilder; import com.rabbitmq.stream.MessageBuilder.ApplicationPropertiesBuilder; @@ -46,6 +44,7 @@ import org.springframework.integration.amqp.support.AmqpHeaderMapper; import org.springframework.integration.amqp.support.DefaultAmqpHeaderMapper; import org.springframework.lang.Nullable; import org.springframework.messaging.MessageHeaders; +import org.springframework.rabbit.stream.listener.ConsumerCustomizer; import org.springframework.rabbit.stream.listener.StreamListenerContainer; import org.springframework.rabbit.stream.support.StreamMessageProperties; import org.springframework.rabbit.stream.support.converter.StreamMessageConverter; @@ -82,15 +81,16 @@ public final class StreamContainerUtils { StreamListenerContainer container = new StreamListenerContainer(applicationContext.getBean(Environment.class)) { @Override - public synchronized void setConsumerCustomizer(Consumer consumerCustomizer) { - super.setConsumerCustomizer(builder -> { - builder.name(consumerDestination.getName()); - consumerCustomizer.accept(builder); + public synchronized void setConsumerCustomizer(ConsumerCustomizer consumerCustomizer) { + super.setConsumerCustomizer((id, builder) -> { + builder.name(consumerDestination.getName() + "." + group); + consumerCustomizer.accept(id, builder); }); } }; + container.setBeanName(consumerDestination.getName() + "." + group + ".container"); container.setMessageConverter(new DefaultStreamMessageConverter()); return container; } diff --git a/spring-cloud-stream-binder-rabbit/src/test/java/org/springframework/cloud/stream/binder/rabbit/stream/RabbitStreamBinderModuleTests.java b/spring-cloud-stream-binder-rabbit/src/test/java/org/springframework/cloud/stream/binder/rabbit/stream/RabbitStreamBinderModuleTests.java index 9517023e9..ebe7f9a74 100644 --- a/spring-cloud-stream-binder-rabbit/src/test/java/org/springframework/cloud/stream/binder/rabbit/stream/RabbitStreamBinderModuleTests.java +++ b/spring-cloud-stream-binder-rabbit/src/test/java/org/springframework/cloud/stream/binder/rabbit/stream/RabbitStreamBinderModuleTests.java @@ -18,6 +18,7 @@ package org.springframework.cloud.stream.binder.rabbit.stream; import com.rabbitmq.stream.ConsumerBuilder; import com.rabbitmq.stream.Environment; +import com.rabbitmq.stream.OffsetSpecification; import org.junit.jupiter.api.AfterEach; import org.junit.jupiter.api.Test; @@ -42,6 +43,7 @@ import org.springframework.rabbit.stream.listener.StreamListenerContainer; import static org.assertj.core.api.Assertions.assertThat; import static org.mockito.BDDMockito.given; import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.verify; /** * @author Gary Russell @@ -59,7 +61,7 @@ public class RabbitStreamBinderModuleTests { } @Test - public void testExtendedProperties() { + public void testStreamContainer() { context = new SpringApplicationBuilder(SimpleProcessor.class) .web(WebApplicationType.NONE) .run("--server.port=0"); @@ -72,8 +74,11 @@ public class RabbitStreamBinderModuleTests { new ExtendedConsumerProperties(rProps); props.setAutoStartup(false); Binding binding = rabbitBinder.bindConsumer("testStream", "grp", new QueueChannel(), props); - assertThat(TestUtils.getPropertyValue(binding, "lifecycle.messageListenerContainer")) - .isInstanceOf(StreamListenerContainer.class); + Object container = TestUtils.getPropertyValue(binding, "lifecycle.messageListenerContainer"); + assertThat(container).isInstanceOf(StreamListenerContainer.class); + ((StreamListenerContainer) container).start(); + verify(this.context.getBean(ConsumerBuilder.class)).offset(OffsetSpecification.first()); + ((StreamListenerContainer) container).stop(); } @SpringBootApplication @@ -81,17 +86,26 @@ public class RabbitStreamBinderModuleTests { @Bean public ListenerContainerCustomizer containerCustomizer() { - return (c, q, g) -> ((StreamListenerContainer) c).setBeanName( - "setByCustomizerForQueue:" + q + (g == null ? "" : ",andGroup:" + g)); + return (cont, dest, group) -> { + StreamListenerContainer container = (StreamListenerContainer) cont; + container.setConsumerCustomizer((name, builder) -> { + builder.offset(OffsetSpecification.first()); + }); + }; } @Bean - Environment env() { + Environment env(ConsumerBuilder builder) { Environment env = mock(Environment.class); - given(env.consumerBuilder()).willReturn(mock(ConsumerBuilder.class)); + given(env.consumerBuilder()).willReturn(builder); return env; } + @Bean + ConsumerBuilder consumerBuilder() { + return mock(ConsumerBuilder.class); + } + } }