GH-1352: StreamRabbitTemplate Improvements

Support for SI `RabbitStreamMessageHandler`, which will live in
the SCSt RabbitMQ binder until SI 6.0 due to versioning.

- expose converters
- use producer message builder if no stream converter provided

* Remove streams before tests; AfterAll has a timing problem.
This commit is contained in:
Gary Russell
2021-09-14 13:38:21 -04:00
committed by GitHub
parent 2a7375ca94
commit c3d7624796
4 changed files with 54 additions and 5 deletions

View File

@@ -58,10 +58,10 @@ ext {
micrometerVersion = '1.8.0-M2'
mockitoVersion = '3.11.2'
protonJVersion = '0.33.8'
rabbitmqStreamVersion = '0.1.0'
rabbitmqStreamVersion = '0.3.0'
rabbitmqVersion = project.hasProperty('rabbitmqVersion') ? project.rabbitmqVersion : '5.13.0'
rabbitmqHttpClientVersion = '3.11.0'
reactorVersion = '2020.0.10'
reactorVersion = '2020.0.11'
snappyVersion = '1.1.8.4'
springDataCommonsVersion = '2.6.0-M2'
springVersion = project.hasProperty('springVersion') ? project.springVersion : '5.3.9'

View File

@@ -16,9 +16,12 @@
package org.springframework.rabbit.stream.producer;
import org.springframework.amqp.AmqpException;
import org.springframework.amqp.core.Message;
import org.springframework.amqp.core.MessagePostProcessor;
import org.springframework.amqp.support.converter.MessageConverter;
import org.springframework.lang.Nullable;
import org.springframework.rabbit.stream.support.converter.StreamMessageConverter;
import org.springframework.util.concurrent.ListenableFuture;
import com.rabbitmq.stream.MessageBuilder;
@@ -65,10 +68,28 @@ public interface RabbitStreamOperations extends AutoCloseable {
ListenableFuture<Boolean> send(com.rabbitmq.stream.Message message);
/**
* Returns the producer's {@link MessageBuilder} to create native stream messages.
* Return the producer's {@link MessageBuilder} to create native stream messages.
* @return the builder.
* @see #send(com.rabbitmq.stream.Message)
*/
MessageBuilder messageBuilder();
/**
* Return the message converter.
* @return the converter.
*/
MessageConverter messageConverter();
/**
* Return the stream message converter.
* @return the converter;
*/
StreamMessageConverter streamMessageConverter();
@Override
default void close() throws AmqpException {
// narrow exception to avoid compiler warning - see
// https://bugs.openjdk.java.net/browse/JDK-8155591
}
}

View File

@@ -56,6 +56,8 @@ public class RabbitStreamTemplate implements RabbitStreamOperations, BeanNameAwa
private StreamMessageConverter streamConverter = new DefaultStreamMessageConverter();
private boolean streamConverterSet;
private Producer producer;
private String beanName;
@@ -81,6 +83,10 @@ public class RabbitStreamTemplate implements RabbitStreamOperations, BeanNameAwa
builder.stream(this.streamName);
this.producerCustomizer.accept(this.beanName, builder);
this.producer = builder.build();
if (!this.streamConverterSet) {
((DefaultStreamMessageConverter) this.streamConverter).setBuilderSupplier(
() -> this.producer.messageBuilder());
}
}
return this.producer;
}
@@ -107,6 +113,7 @@ public class RabbitStreamTemplate implements RabbitStreamOperations, BeanNameAwa
public void setStreamConverter(StreamMessageConverter streamConverter) {
Assert.notNull(streamConverter, "'streamConverter' cannot be null");
this.streamConverter = streamConverter;
this.streamConverterSet = true;
}
/**
@@ -118,6 +125,18 @@ public class RabbitStreamTemplate implements RabbitStreamOperations, BeanNameAwa
this.producerCustomizer = producerCustomizer;
}
@Override
public MessageConverter messageConverter() {
return this.messageConverter;
}
@Override
public StreamMessageConverter streamMessageConverter() {
return this.streamConverter;
}
@Override
public ListenableFuture<Boolean> send(Message message) {
SettableListenableFuture<Boolean> future = new SettableListenableFuture<>();

View File

@@ -24,7 +24,6 @@ import java.util.concurrent.CountDownLatch;
import java.util.concurrent.Future;
import java.util.concurrent.TimeUnit;
import org.junit.jupiter.api.AfterAll;
import org.junit.jupiter.api.Test;
import org.springframework.amqp.core.Queue;
@@ -65,7 +64,7 @@ public class RabbitListenerTests extends AbstractIntegrationTests {
@Autowired
Config config;
@AfterAll
// @AfterAll - causes test to throw errors - need to investigate
static void deleteQueues() {
try (Environment environment = Config.environment()) {
environment.deleteStream("test.stream.queue1");
@@ -140,6 +139,16 @@ public class RabbitListenerTests extends AbstractIntegrationTests {
@Override
public void start() {
try {
env.deleteStream("test.stream.queue1");
}
catch (Exception e) {
}
try {
env.deleteStream("test.stream.queue2");
}
catch (Exception e) {
}
env.streamCreator().stream("test.stream.queue1").create();
env.streamCreator().stream("test.stream.queue2").create();
}