diff --git a/build.gradle b/build.gradle index a74a9fe3..91ede6bd 100644 --- a/build.gradle +++ b/build.gradle @@ -58,7 +58,7 @@ ext { micrometerVersion = '1.10.0-SNAPSHOT' micrometerTracingVersion = '1.0.0-SNAPSHOT' mockitoVersion = '4.5.1' - rabbitmqStreamVersion = '0.4.0' + rabbitmqStreamVersion = '0.7.0' rabbitmqVersion = project.hasProperty('rabbitmqVersion') ? project.rabbitmqVersion : '5.13.1' rabbitmqHttpClientVersion = '3.12.1' reactorVersion = '2020.0.18' @@ -105,7 +105,7 @@ allprojects { maven { url 'https://repo.spring.io/libs-milestone' } if (version.endsWith('-SNAPSHOT')) { maven { url 'https://repo.spring.io/libs-snapshot' } - maven { url 'https://oss.sonatype.org/content/repositories/snapshots' } + // maven { url 'https://oss.sonatype.org/content/repositories/snapshots' } } // maven { url 'https://repo.spring.io/libs-staging-local' } } diff --git a/spring-rabbit-stream/src/main/java/org/springframework/rabbit/stream/listener/StreamListenerContainer.java b/spring-rabbit-stream/src/main/java/org/springframework/rabbit/stream/listener/StreamListenerContainer.java index 2a9684cc..64cd4ccd 100644 --- a/spring-rabbit-stream/src/main/java/org/springframework/rabbit/stream/listener/StreamListenerContainer.java +++ b/spring-rabbit-stream/src/main/java/org/springframework/rabbit/stream/listener/StreamListenerContainer.java @@ -96,6 +96,19 @@ public class StreamListenerContainer implements MessageListenerContainer, BeanNa this.builder.stream(queueNames[0]); } + /** + * Enable Single Active Consumer on a Super Stream. + * @param superStream the stream. + * @param name the consumer name. + * @since 3.0 + */ + public void superStream(String superStream, String name) { + Assert.notNull(superStream, "'superStream' cannot be null"); + this.builder.superStream(superStream) + .singleActiveConsumer() + .name(name); + } + /** * Get a {@link StreamMessageConverter} used to convert a * {@link com.rabbitmq.stream.Message} to a diff --git a/spring-rabbit-stream/src/test/java/org/springframework/rabbit/stream/listener/SuperStreamSACTests.java b/spring-rabbit-stream/src/test/java/org/springframework/rabbit/stream/listener/SuperStreamSACTests.java new file mode 100644 index 00000000..e9d3a28c --- /dev/null +++ b/spring-rabbit-stream/src/test/java/org/springframework/rabbit/stream/listener/SuperStreamSACTests.java @@ -0,0 +1,142 @@ +/* + * Copyright 2022 the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * https://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.springframework.rabbit.stream.listener; + +import static org.assertj.core.api.Assertions.assertThat; + +import java.util.ArrayList; +import java.util.List; +import java.util.Map; +import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.TimeUnit; + +import org.junit.jupiter.api.Test; + +import org.springframework.amqp.core.Declarables; +import org.springframework.amqp.core.DirectExchange; +import org.springframework.amqp.core.Message; +import org.springframework.amqp.core.Queue; +import org.springframework.amqp.rabbit.connection.CachingConnectionFactory; +import org.springframework.amqp.rabbit.connection.ConnectionFactory; +import org.springframework.amqp.rabbit.core.RabbitAdmin; +import org.springframework.amqp.rabbit.core.RabbitTemplate; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.beans.factory.config.ConfigurableBeanFactory; +import org.springframework.context.ApplicationContext; +import org.springframework.context.annotation.Bean; +import org.springframework.context.annotation.Configuration; +import org.springframework.context.annotation.Scope; +import org.springframework.rabbit.stream.config.SuperStream; +import org.springframework.rabbit.stream.support.AbstractIntegrationTests; +import org.springframework.test.context.junit.jupiter.SpringJUnitConfig; + +import com.rabbitmq.stream.Address; +import com.rabbitmq.stream.Environment; +import com.rabbitmq.stream.OffsetSpecification; + +/** + * @author Gary Russell + * @since 3.0 + * + */ +@SpringJUnitConfig +public class SuperStreamSACTests extends AbstractIntegrationTests { + + @Test + void superStream(@Autowired ApplicationContext context, @Autowired RabbitTemplate template, + @Autowired Environment env, @Autowired Config config, @Autowired RabbitAdmin admin, + @Autowired Declarables declarables) throws InterruptedException { + + template.getConnectionFactory().createConnection(); + StreamListenerContainer container1 = context.getBean(StreamListenerContainer.class, env, "one"); + container1.start(); + StreamListenerContainer container2 = context.getBean(StreamListenerContainer.class, env, "two"); + container2.start(); + StreamListenerContainer container3 = context.getBean(StreamListenerContainer.class, env, "three"); + container3.start(); + template.convertAndSend("ss.sac.test", "0", "foo"); + template.convertAndSend("ss.sac.test", "1", "bar"); + template.convertAndSend("ss.sac.test", "2", "baz"); + assertThat(config.latch.await(10, TimeUnit.SECONDS)).isTrue(); + assertThat(config.messages.keySet()).contains("one", "two", "three"); + assertThat(config.info).contains("one:foo", "two:bar", "three:baz"); + container1.stop(); + container2.stop(); + container3.stop(); + clean(admin, declarables); + } + + private void clean(RabbitAdmin admin, Declarables declarables) { + declarables.getDeclarablesByType(Queue.class).forEach(queue -> admin.deleteQueue(queue.getName())); + declarables.getDeclarablesByType(DirectExchange.class).forEach(ex -> admin.deleteExchange(ex.getName())); + } + + @Configuration + public static class Config { + + final List info = new ArrayList<>(); + + final Map messages = new ConcurrentHashMap<>(); + + final CountDownLatch latch = new CountDownLatch(3); + + @Bean + CachingConnectionFactory cf() { + return new CachingConnectionFactory("localhost", amqpPort()); + } + + @Bean + RabbitAdmin admin(ConnectionFactory cf) { + return new RabbitAdmin(cf); + } + + @Bean + RabbitTemplate template(ConnectionFactory cf) { + return new RabbitTemplate(cf); + } + + @Bean + SuperStream superStream() { + return new SuperStream("ss.sac.test", 3); + } + + @Bean + static Environment environment() { + return Environment.builder() + .addressResolver(add -> new Address("localhost", streamPort())) + .build(); + } + + @Bean + @Scope(ConfigurableBeanFactory.SCOPE_PROTOTYPE) + StreamListenerContainer container(Environment env, String name) { + StreamListenerContainer container = new StreamListenerContainer(env); + container.superStream("ss.sac.test", "test"); + container.setupMessageListener(msg -> { + this.messages.put(name, msg); + this.info.add(name + ":" + new String(msg.getBody())); + this.latch.countDown(); + }); + container.setConsumerCustomizer((id, builder) -> builder.offset(OffsetSpecification.last())); + container.setAutoStartup(false); + return container; + } + + } + +} diff --git a/src/reference/asciidoc/quick-tour.adoc b/src/reference/asciidoc/quick-tour.adoc index af2aa94e..28a4a6bc 100644 --- a/src/reference/asciidoc/quick-tour.adoc +++ b/src/reference/asciidoc/quick-tour.adoc @@ -36,6 +36,8 @@ The minimum Spring Framework version dependency is 5.2.0. The minimum `amqp-client` Java client library version is 5.7.0. +The minimum `stream-client` Java client library for stream queues is 0.7.0. + ===== Very, Very Quick This section offers the fastest introduction. diff --git a/src/reference/asciidoc/stream.adoc b/src/reference/asciidoc/stream.adoc index bc18e091..776eb5f3 100644 --- a/src/reference/asciidoc/stream.adoc +++ b/src/reference/asciidoc/stream.adoc @@ -182,4 +182,21 @@ The `RabbitAdmin` detects this bean and will declare the exchange (`my.super.str ===== Consuming Super Streams with Single Active Consumers -TBD. +Invoke the `superStream` method on the listener container to enable a single active consumer on a super stream. + +==== +[source, java] +---- +@Bean +@Scope(ConfigurableBeanFactory.SCOPE_PROTOTYPE) +StreamListenerContainer container(Environment env, String name) { + StreamListenerContainer container = new StreamListenerContainer(env); + container.superStream("ss.sac", "myConsumer"); + container.setupMessageListener(msg -> { + ... + }); + container.setConsumerCustomizer((id, builder) -> builder.offset(OffsetSpecification.last())); + return container; +} +---- +==== diff --git a/src/reference/asciidoc/whats-new.adoc b/src/reference/asciidoc/whats-new.adoc index 6825e4b5..2b43baca 100644 --- a/src/reference/asciidoc/whats-new.adoc +++ b/src/reference/asciidoc/whats-new.adoc @@ -20,6 +20,9 @@ See <> for more information. ==== Stream Support Changes `RabbitStreamOperations2` and `RabbitStreamTemplate2` have been deprecated in favor of `RabbitStreamOperations` and `RabbitStreamTemplate` respectively. + +Super streams and single active consumers thereon are now supported. + See <> for more information. ==== `@RabbitListener` Changes @@ -35,4 +38,4 @@ See <> for more information. ==== Connection Factory Changes The default `addressShuffleMode` in `AbstractConnectionFactory` is now `RANDOM`. This results in connecting to a random host when multiple addresses are provided. -See <> for more information. \ No newline at end of file +See <> for more information.