GH-325: RabbitMQ Stream to Latest Snapshot
* Fix docs; set container bean name; fix stream name.
This commit is contained in:
@@ -438,7 +438,7 @@ Environment streamEnv() {
|
||||
ListenerContainerCustomizer<MessageListenerContainer> 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<MessageListenerContainer> 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:
|
||||
|
||||
====
|
||||
|
||||
@@ -58,7 +58,7 @@
|
||||
<dependency>
|
||||
<groupId>org.springframework.amqp</groupId>
|
||||
<artifactId>spring-rabbit-stream</artifactId>
|
||||
<version>2.4.0-M1</version>
|
||||
<version>2.4.0-SNAPSHOT</version>
|
||||
<optional>true</optional>
|
||||
</dependency>
|
||||
<dependency>
|
||||
|
||||
@@ -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());
|
||||
|
||||
@@ -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<ConsumerBuilder> 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;
|
||||
}
|
||||
|
||||
@@ -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<RabbitConsumerProperties>(rProps);
|
||||
props.setAutoStartup(false);
|
||||
Binding<MessageChannel> 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<MessageListenerContainer> 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);
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user