From c3d7624796d346e82e12dedb78acbff89283ed26 Mon Sep 17 00:00:00 2001 From: Gary Russell Date: Tue, 14 Sep 2021 13:38:21 -0400 Subject: [PATCH] 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. --- build.gradle | 4 ++-- .../producer/RabbitStreamOperations.java | 23 ++++++++++++++++++- .../stream/producer/RabbitStreamTemplate.java | 19 +++++++++++++++ .../stream/listener/RabbitListenerTests.java | 13 +++++++++-- 4 files changed, 54 insertions(+), 5 deletions(-) diff --git a/build.gradle b/build.gradle index 5f6f9c7a..ae85085e 100644 --- a/build.gradle +++ b/build.gradle @@ -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' diff --git a/spring-rabbit-stream/src/main/java/org/springframework/rabbit/stream/producer/RabbitStreamOperations.java b/spring-rabbit-stream/src/main/java/org/springframework/rabbit/stream/producer/RabbitStreamOperations.java index 3481da72..bf3a8ea6 100644 --- a/spring-rabbit-stream/src/main/java/org/springframework/rabbit/stream/producer/RabbitStreamOperations.java +++ b/spring-rabbit-stream/src/main/java/org/springframework/rabbit/stream/producer/RabbitStreamOperations.java @@ -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 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 + } + } diff --git a/spring-rabbit-stream/src/main/java/org/springframework/rabbit/stream/producer/RabbitStreamTemplate.java b/spring-rabbit-stream/src/main/java/org/springframework/rabbit/stream/producer/RabbitStreamTemplate.java index 32079daa..b651832b 100644 --- a/spring-rabbit-stream/src/main/java/org/springframework/rabbit/stream/producer/RabbitStreamTemplate.java +++ b/spring-rabbit-stream/src/main/java/org/springframework/rabbit/stream/producer/RabbitStreamTemplate.java @@ -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 send(Message message) { SettableListenableFuture future = new SettableListenableFuture<>(); diff --git a/spring-rabbit-stream/src/test/java/org/springframework/rabbit/stream/listener/RabbitListenerTests.java b/spring-rabbit-stream/src/test/java/org/springframework/rabbit/stream/listener/RabbitListenerTests.java index 0d53093c..2c7352cb 100644 --- a/spring-rabbit-stream/src/test/java/org/springframework/rabbit/stream/listener/RabbitListenerTests.java +++ b/spring-rabbit-stream/src/test/java/org/springframework/rabbit/stream/listener/RabbitListenerTests.java @@ -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(); }