diff --git a/binders/rabbit-binder/spring-cloud-stream-binder-rabbit-core/pom.xml b/binders/rabbit-binder/spring-cloud-stream-binder-rabbit-core/pom.xml
index b3d24ec9d..0b3e0c5de 100644
--- a/binders/rabbit-binder/spring-cloud-stream-binder-rabbit-core/pom.xml
+++ b/binders/rabbit-binder/spring-cloud-stream-binder-rabbit-core/pom.xml
@@ -22,12 +22,6 @@
org.springframework.boot
spring-boot-starter-amqp
-
- org.springframework.amqp
- spring-rabbit-stream
- 3.0.0-SNAPSHOT
- true
-
org.springframework.integration
spring-integration-amqp
diff --git a/binders/rabbit-binder/spring-cloud-stream-binder-rabbit-core/src/main/java/org/springframework/cloud/stream/binder/rabbit/properties/RabbitConsumerProperties.java b/binders/rabbit-binder/spring-cloud-stream-binder-rabbit-core/src/main/java/org/springframework/cloud/stream/binder/rabbit/properties/RabbitConsumerProperties.java
index 25ae65a84..639d2bcf2 100644
--- a/binders/rabbit-binder/spring-cloud-stream-binder-rabbit-core/src/main/java/org/springframework/cloud/stream/binder/rabbit/properties/RabbitConsumerProperties.java
+++ b/binders/rabbit-binder/spring-cloud-stream-binder-rabbit-core/src/main/java/org/springframework/cloud/stream/binder/rabbit/properties/RabbitConsumerProperties.java
@@ -115,7 +115,7 @@ public class RabbitConsumerProperties extends RabbitCommonProperties {
private int frameMaxHeadroom = 20_000;
/**
- * The container type, SIMPLE, DIRECT, or STREAM.
+ * The container type, SIMPLE or DIRECT.
*/
private ContainerType containerType = ContainerType.SIMPLE;
@@ -139,12 +139,6 @@ public class RabbitConsumerProperties extends RabbitCommonProperties {
*/
private Long receiveTimeout;
- /**
- * When the container type is STREAM, set this to true to create a super stream with
- * competing consumers.
- */
- private boolean superStream;
-
public boolean isTransacted() {
return transacted;
}
@@ -353,14 +347,6 @@ public class RabbitConsumerProperties extends RabbitCommonProperties {
this.receiveTimeout = receiveTimeout;
}
- public boolean isSuperStream() {
- return this.superStream;
- }
-
- public void setSuperStream(boolean superStream) {
- this.superStream = superStream;
- }
-
/**
* Container type.
* @author Gary Russell
diff --git a/binders/rabbit-binder/spring-cloud-stream-binder-rabbit-core/src/main/java/org/springframework/cloud/stream/binder/rabbit/properties/RabbitProducerProperties.java b/binders/rabbit-binder/spring-cloud-stream-binder-rabbit-core/src/main/java/org/springframework/cloud/stream/binder/rabbit/properties/RabbitProducerProperties.java
index 678eb7c9d..46d789657 100644
--- a/binders/rabbit-binder/spring-cloud-stream-binder-rabbit-core/src/main/java/org/springframework/cloud/stream/binder/rabbit/properties/RabbitProducerProperties.java
+++ b/binders/rabbit-binder/spring-cloud-stream-binder-rabbit-core/src/main/java/org/springframework/cloud/stream/binder/rabbit/properties/RabbitProducerProperties.java
@@ -157,12 +157,6 @@ public class RabbitProducerProperties extends RabbitCommonProperties {
*/
private AlternateExchange alternateExchange;
- /**
- * When the producer type is STREAM_*, set this to true to publish to a super stream.
- * Also requires a partition key.
- */
- private boolean superStream;
-
/**
* @param requestHeaderPatterns the patterns.
* @deprecated - use {@link #setHeaderPatterns(String[])}.
@@ -325,14 +319,6 @@ public class RabbitProducerProperties extends RabbitCommonProperties {
this.alternateExchange = alternate;
}
- public boolean isSuperStream() {
- return this.superStream;
- }
-
- public void setSuperStream(boolean superStream) {
- this.superStream = superStream;
- }
-
public static class AlternateExchange {
/**
diff --git a/binders/rabbit-binder/spring-cloud-stream-binder-rabbit-core/src/main/java/org/springframework/cloud/stream/binder/rabbit/provisioning/RabbitExchangeQueueProvisioner.java b/binders/rabbit-binder/spring-cloud-stream-binder-rabbit-core/src/main/java/org/springframework/cloud/stream/binder/rabbit/provisioning/RabbitExchangeQueueProvisioner.java
index c93e87467..f6f85b44d 100644
--- a/binders/rabbit-binder/spring-cloud-stream-binder-rabbit-core/src/main/java/org/springframework/cloud/stream/binder/rabbit/provisioning/RabbitExchangeQueueProvisioner.java
+++ b/binders/rabbit-binder/spring-cloud-stream-binder-rabbit-core/src/main/java/org/springframework/cloud/stream/binder/rabbit/provisioning/RabbitExchangeQueueProvisioner.java
@@ -22,8 +22,6 @@ import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.concurrent.atomic.AtomicInteger;
-import java.util.stream.Collectors;
-import java.util.stream.IntStream;
import java.util.stream.Stream;
import org.apache.commons.logging.Log;
@@ -66,7 +64,6 @@ import org.springframework.cloud.stream.provisioning.ProvisioningProvider;
import org.springframework.context.ApplicationListener;
import org.springframework.context.support.GenericApplicationContext;
import org.springframework.lang.Nullable;
-import org.springframework.rabbit.stream.config.SuperStream;
import org.springframework.util.Assert;
import org.springframework.util.ObjectUtils;
import org.springframework.util.StringUtils;
@@ -248,11 +245,9 @@ public class RabbitExchangeQueueProvisioner
+ ", bound to: " + name);
}
String prefix = properties.getExtension().getPrefix();
- String exchangeName = applyPrefix(prefix, name);
- ContainerType containerType = properties.getExtension().getContainerType();
- boolean superStream = containerType.equals(ContainerType.STREAM) && properties.getExtension().isSuperStream();
+ final String exchangeName = applyPrefix(prefix, name);
Exchange exchange = buildExchange(properties.getExtension(), exchangeName, null, null);
- if (!superStream && properties.getExtension().isDeclareExchange()) {
+ if (properties.getExtension().isDeclareExchange()) {
declareExchange(exchangeName, anonymous ? anonymousGroup : group, exchange);
}
String queueName = applyPrefix(prefix, baseQueueName);
@@ -263,7 +258,6 @@ public class RabbitExchangeQueueProvisioner
String anonQueueName = queueName;
queue = new AnonymousQueue((org.springframework.amqp.core.NamingStrategy) () -> anonQueueName,
queueArgs(queueName, properties.getExtension(), false));
- queueName = queue.getName();
}
else {
if (partitioned) {
@@ -281,68 +275,25 @@ public class RabbitExchangeQueueProvisioner
}
Binding binding = null;
if (properties.getExtension().isBindQueue()) {
- if (superStream) {
- provisionSuperStream(properties, name);
+ if (properties.getExtension().getContainerType().equals(ContainerType.STREAM)) {
+ queue.getArguments().put("x-queue-type", "stream");
+ }
+ declareQueue(queueName, queue);
+ String[] routingKeys = bindingRoutingKeys(properties.getExtension());
+ if (ObjectUtils.isEmpty(routingKeys)) {
+ binding = declareConsumerBindings(name, null, properties, exchange, partitioned, queue);
}
else {
- if (containerType.equals(ContainerType.STREAM)) {
- queue.getArguments().put("x-queue-type", "stream");
- }
- declareQueue(queueName, queue);
- String[] routingKeys = bindingRoutingKeys(properties.getExtension());
- if (ObjectUtils.isEmpty(routingKeys)) {
- binding = declareConsumerBindings(name, null, properties, exchange, partitioned, queue);
- }
- else {
- for (String routingKey : routingKeys) {
- binding = declareConsumerBindings(name, routingKey, properties, exchange, partitioned, queue);
- }
+ for (String routingKey : routingKeys) {
+ binding = declareConsumerBindings(name, routingKey, properties, exchange, partitioned, queue);
}
}
}
- if (durable && !superStream) {
+ if (durable) {
autoBindDLQ(applyPrefix(properties.getExtension().getPrefix(), baseQueueName),
queueName, group, properties.getExtension());
}
- if (superStream) {
- queueName = name; // group is used in the consumer for super streams so not part of the name.
- }
- return new RabbitConsumerDestination(queueName, binding, anonymous ? baseQueueName : group, name);
- }
-
- private void provisionSuperStream(ExtendedConsumerProperties properties,
- String name) {
-
- String routingKey = properties.getExtension().getBindingRoutingKey();
- String rk = routingKey == null ? name : routingKey;
- SuperStream ss = new SuperStream(name, properties.getInstanceCount(), (q, i) -> IntStream.range(0, i)
- .mapToObj(j -> rk + "-" + j)
- .collect(Collectors.toList()));
- synchronized (this.autoDeclareContext) {
- if (!this.autoDeclareContext.containsBean(name + ".superStream")) {
- this.autoDeclareContext.getBeanFactory().registerSingleton(name + ".superStream", ss);
- }
- }
- try {
- ss.getDeclarables().forEach(dec -> {
- if (dec instanceof Exchange exch) {
- this.rabbitAdmin.declareExchange(exch);
- }
- else if (dec instanceof Queue queue) {
- this.rabbitAdmin.declareQueue(queue);
- }
- else if (dec instanceof Binding binding) {
- this.rabbitAdmin.declareBinding(binding);
- }
- });
- }
- catch (AmqpConnectException e) {
- if (this.logger.isDebugEnabled()) {
- this.logger.debug("Declaration of super stream: " + name
- + " deferred - connection not available");
- }
- }
-
+ return new RabbitConsumerDestination(queue.getName(), binding, anonymous ? baseQueueName : group, name);
}
/**
diff --git a/binders/rabbit-binder/spring-cloud-stream-binder-rabbit/pom.xml b/binders/rabbit-binder/spring-cloud-stream-binder-rabbit/pom.xml
index cd1d30493..28b88b37b 100644
--- a/binders/rabbit-binder/spring-cloud-stream-binder-rabbit/pom.xml
+++ b/binders/rabbit-binder/spring-cloud-stream-binder-rabbit/pom.xml
@@ -64,23 +64,11 @@
org.springframework.amqp
spring-rabbit-stream
- 3.0.0-SNAPSHOT
true
-
- org.springframework.integration
- spring-integration-amqp
- 6.0.0-SNAPSHOT
-
-
- org.springframework.integration
- spring-integration-core
- 6.0.0-SNAPSHOT
-
org.springframework.integration
spring-integration-jmx
- 6.0.0-SNAPSHOT
org.springframework.integration
diff --git a/binders/rabbit-binder/spring-cloud-stream-binder-rabbit/src/main/java/org/springframework/cloud/stream/binder/rabbit/RabbitMessageChannelBinder.java b/binders/rabbit-binder/spring-cloud-stream-binder-rabbit/src/main/java/org/springframework/cloud/stream/binder/rabbit/RabbitMessageChannelBinder.java
index 8db8fddb1..6f33e2dae 100644
--- a/binders/rabbit-binder/spring-cloud-stream-binder-rabbit/src/main/java/org/springframework/cloud/stream/binder/rabbit/RabbitMessageChannelBinder.java
+++ b/binders/rabbit-binder/spring-cloud-stream-binder-rabbit/src/main/java/org/springframework/cloud/stream/binder/rabbit/RabbitMessageChannelBinder.java
@@ -422,10 +422,7 @@ public class RabbitMessageChannelBinder extends
private AmqpHeaderMapper configureHeaderMapper(RabbitProducerProperties extendedProperties) {
DefaultAmqpHeaderMapper mapper = DefaultAmqpHeaderMapper.outboundMapper();
List headerPatterns = new ArrayList<>(extendedProperties.getHeaderPatterns().length + 3);
- if (!extendedProperties.isSuperStream()) {
- // need to keep this header until later
- headerPatterns.add("!" + BinderHeaders.PARTITION_HEADER);
- }
+ headerPatterns.add("!" + BinderHeaders.PARTITION_HEADER);
headerPatterns.add("!" + IntegrationMessageHeaderAccessor.SOURCE_DATA);
headerPatterns.add("!" + IntegrationMessageHeaderAccessor.DELIVERY_ATTEMPT);
headerPatterns.add("!rabbitmq_streamContext");
@@ -503,10 +500,7 @@ public class RabbitMessageChannelBinder extends
MessageListenerContainer listenerContainer = createAndConfigureContainer(consumerDestination, group,
properties, destination, extension);
String[] queues = StringUtils.tokenizeToStringArray(destination, ",", true, true);
- if (properties.getExtension().getContainerType() != ContainerType.STREAM
- || !properties.getExtension().isSuperStream()) {
- listenerContainer.setQueueNames(queues);
- }
+ listenerContainer.setQueueNames(queues);
getContainerCustomizer().configure(listenerContainer,
consumerDestination.getName(), group);
listenerContainer.afterPropertiesSet();
@@ -529,12 +523,11 @@ public class RabbitMessageChannelBinder extends
adapter.setErrorChannel(errorInfrastructure.getErrorChannel());
}
adapter.setMessageConverter(passThoughConverter);
- ContainerType containerType = extension.getContainerType();
if (properties.isBatchMode() && extension.isEnableBatching()
- && ContainerType.SIMPLE.equals(containerType)) {
+ && ContainerType.SIMPLE.equals(extension.getContainerType())) {
adapter.setBatchMode(BatchMode.EXTRACT_PAYLOADS_WITH_HEADERS);
}
- if (containerType.equals(ContainerType.STREAM)) {
+ if (extension.getContainerType().equals(ContainerType.STREAM)) {
StreamUtils.configureAdapter(adapter);
}
return adapter;
diff --git a/binders/rabbit-binder/spring-cloud-stream-binder-rabbit/src/main/java/org/springframework/cloud/stream/binder/rabbit/RabbitStreamMessageHandler.java b/binders/rabbit-binder/spring-cloud-stream-binder-rabbit/src/main/java/org/springframework/cloud/stream/binder/rabbit/RabbitStreamMessageHandler.java
new file mode 100644
index 000000000..36f587d73
--- /dev/null
+++ b/binders/rabbit-binder/spring-cloud-stream-binder-rabbit/src/main/java/org/springframework/cloud/stream/binder/rabbit/RabbitStreamMessageHandler.java
@@ -0,0 +1,275 @@
+/*
+ * Copyright 2021-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.cloud.stream.binder.rabbit;
+
+import java.util.concurrent.CompletableFuture;
+import java.util.concurrent.ExecutionException;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.TimeoutException;
+
+import org.springframework.amqp.core.MessageProperties;
+import org.springframework.amqp.support.AmqpHeaders;
+import org.springframework.amqp.support.converter.ContentTypeDelegatingMessageConverter;
+import org.springframework.amqp.support.converter.MessageConverter;
+import org.springframework.context.Lifecycle;
+import org.springframework.integration.amqp.support.AmqpHeaderMapper;
+import org.springframework.integration.amqp.support.DefaultAmqpHeaderMapper;
+import org.springframework.integration.handler.AbstractMessageHandler;
+import org.springframework.lang.Nullable;
+import org.springframework.messaging.Message;
+import org.springframework.messaging.MessageHandler;
+import org.springframework.messaging.MessageHandlingException;
+import org.springframework.messaging.MessageHeaders;
+import org.springframework.rabbit.stream.producer.RabbitStreamOperations;
+import org.springframework.rabbit.stream.support.StreamMessageProperties;
+import org.springframework.util.Assert;
+import org.springframework.util.MimeType;
+
+/**
+ * {@link MessageHandler} based on {@link RabbitStreamOperations}.
+ *
+ * TODO: This class will move to Spring Integration in 6.0.
+ *
+ * @author Gary Russell
+ * @author Chris Bono
+ * @since 3.2
+ *
+ */
+public class RabbitStreamMessageHandler extends AbstractMessageHandler implements Lifecycle {
+
+ private static final int DEFAULT_CONFIRM_TIMEOUT = 10_000;
+
+ private final RabbitStreamOperations streamOperations;
+
+ private boolean sync;
+
+ private long confirmTimeout = DEFAULT_CONFIRM_TIMEOUT;
+
+ private SuccessCallback> successCallback = msg -> { };
+
+ private FailureCallback failureCallback = (msg, ex) -> { };
+
+ private AmqpHeaderMapper headerMapper = DefaultAmqpHeaderMapper.outboundMapper();
+
+ private boolean headersMappedLast;
+
+ /**
+ * Create an instance with the provided {@link RabbitStreamOperations}.
+ * @param streamOperations the operations.
+ */
+ public RabbitStreamMessageHandler(RabbitStreamOperations streamOperations) {
+ Assert.notNull(streamOperations, "'streamOperations' cannot be null");
+ this.streamOperations = streamOperations;
+ }
+
+ /**
+ * Set a callback to be invoked when a send is successful.
+ * @param successCallback the callback.
+ */
+ public void setSuccessCallback(SuccessCallback> successCallback) {
+ Assert.notNull(successCallback, "'successCallback' cannot be null");
+ this.successCallback = successCallback;
+ }
+
+ /**
+ * Set a callback to be invoked when a send fails.
+ * @param failureCallback the callback.
+ */
+ public void setFailureCallback(FailureCallback failureCallback) {
+ Assert.notNull(failureCallback, "'failureCallback' cannot be null");
+ this.failureCallback = failureCallback;
+ }
+
+ /**
+ * Set to true to wait for a confirmation.
+ * @param sync true to wait.
+ * @see #setConfirmTimeout(long)
+ */
+ public void setSync(boolean sync) {
+ this.sync = sync;
+ }
+
+ /**
+ * Set the confirm timeout.
+ * @param confirmTimeout the timeout.
+ * @see #setSync(boolean)
+ */
+ public void setConfirmTimeout(long confirmTimeout) {
+ this.confirmTimeout = confirmTimeout;
+ }
+
+ /**
+ * Set a custom {@link AmqpHeaderMapper} for mapping request and reply headers.
+ * Defaults to {@link DefaultAmqpHeaderMapper#outboundMapper()}.
+ * @param headerMapper the {@link AmqpHeaderMapper} to use.
+ */
+ public void setHeaderMapper(AmqpHeaderMapper headerMapper) {
+ Assert.notNull(headerMapper, "headerMapper must not be null");
+ this.headerMapper = headerMapper;
+ }
+
+ /**
+ * When mapping headers for the outbound message, determine whether the headers are
+ * mapped before the message is converted, or afterwards. This only affects headers
+ * that might be added by the message converter. When false, the converter's headers
+ * win; when true, any headers added by the converter will be overridden (if the
+ * source message has a header that maps to those headers). You might wish to set this
+ * to true, for example, when using a
+ * {@link org.springframework.amqp.support.converter.SimpleMessageConverter} with a
+ * String payload that contains json; the converter will set the content type to
+ * {@code text/plain} which can be overridden to {@code application/json} by setting
+ * the {@link AmqpHeaders#CONTENT_TYPE} message header. Default: false.
+ * @param headersMappedLast true if headers are mapped after conversion.
+ */
+ public void setHeadersMappedLast(boolean headersMappedLast) {
+ this.headersMappedLast = headersMappedLast;
+ }
+
+ /**
+ * Return the {@link RabbitStreamOperations}.
+ * @return the operations.
+ */
+ public RabbitStreamOperations getStreamOperations() {
+ return this.streamOperations;
+ }
+
+ @Override
+ protected void handleMessageInternal(Message> requestMessage) {
+ CompletableFuture future;
+ com.rabbitmq.stream.Message streamMessage;
+ if (requestMessage.getPayload() instanceof com.rabbitmq.stream.Message) {
+ streamMessage = (com.rabbitmq.stream.Message) requestMessage.getPayload();
+ }
+ else {
+ MessageConverter converter = streamOperations.messageConverter();
+ org.springframework.amqp.core.Message amqpMessage = mapMessage(requestMessage, converter,
+ this.headerMapper, this.headersMappedLast);
+ streamMessage = this.streamOperations.streamMessageConverter().fromMessage(amqpMessage);
+ }
+ future = this.streamOperations.send(streamMessage);
+ handleConfirms(requestMessage, future);
+ }
+
+ private void handleConfirms(Message> message, CompletableFuture future) {
+ future.whenComplete((bool, ex) -> {
+ if (ex != null) {
+ this.failureCallback.failure(message, ex);
+ }
+ else {
+ this.successCallback.onSuccess(message);
+ }
+ });
+ if (this.sync) {
+ try {
+ future.get(this.confirmTimeout, TimeUnit.MILLISECONDS);
+ }
+ catch (InterruptedException ex) {
+ Thread.currentThread().interrupt();
+ throw new MessageHandlingException(message, ex);
+ }
+ catch (ExecutionException | TimeoutException ex) {
+ throw new MessageHandlingException(message, ex);
+ }
+ }
+ }
+
+ /*
+ * TODO Copied/modified from MapppingUtils until SI 6.0
+ */
+ private static org.springframework.amqp.core.Message mapMessage(Message> message,
+ MessageConverter converter, AmqpHeaderMapper headerMapper, boolean headersMappedLast) {
+
+ MessageProperties amqpMessageProperties = new StreamMessageProperties();
+ org.springframework.amqp.core.Message amqpMessage;
+ if (!headersMappedLast) {
+ mapHeaders(message.getHeaders(), amqpMessageProperties, headerMapper);
+ }
+ if (converter instanceof ContentTypeDelegatingMessageConverter && headersMappedLast) {
+ String contentType = contentTypeAsString(message.getHeaders());
+ if (contentType != null) {
+ amqpMessageProperties.setContentType(contentType);
+ }
+ }
+ amqpMessage = converter.toMessage(message.getPayload(), amqpMessageProperties);
+ if (headersMappedLast) {
+ mapHeaders(message.getHeaders(), amqpMessageProperties, headerMapper);
+ }
+ return amqpMessage;
+ }
+
+ private static void mapHeaders(MessageHeaders messageHeaders, MessageProperties amqpMessageProperties,
+ AmqpHeaderMapper headerMapper) {
+
+ headerMapper.fromHeadersToRequest(messageHeaders, amqpMessageProperties);
+ }
+
+ private static String contentTypeAsString(MessageHeaders headers) {
+ Object contentType = headers.get(AmqpHeaders.CONTENT_TYPE);
+ if (contentType instanceof MimeType) {
+ contentType = contentType.toString();
+ }
+ if (contentType instanceof String) {
+ return (String) contentType;
+ }
+ else if (contentType != null) {
+ throw new IllegalArgumentException(AmqpHeaders.CONTENT_TYPE
+ + " header must be a MimeType or String, found: " + contentType.getClass().getName());
+ }
+ return null;
+ }
+ /*
+ * End copied/modified from MappingUtils
+ */
+
+ @Override
+ public void start() {
+ }
+
+ @Override
+ public void stop() {
+ this.streamOperations.close();
+ }
+
+ @Override
+ public boolean isRunning() {
+ return true;
+ }
+
+ /**
+ * Callback for when publishing succeeds.
+ */
+ interface SuccessCallback {
+ /**
+ * Called when the future completes with success.
+ * Note that Exceptions raised by this method are ignored.
+ * @param result the result of the future
+ */
+ void onSuccess(@Nullable T result);
+ }
+
+ /**
+ * Callback for when publishing fails.
+ */
+ interface FailureCallback {
+ /**
+ * Message publish failure.
+ * @param message the message.
+ * @param throwable the throwable.
+ */
+ void failure(Message> message, Throwable throwable);
+ }
+}
diff --git a/binders/rabbit-binder/spring-cloud-stream-binder-rabbit/src/main/java/org/springframework/cloud/stream/binder/rabbit/StreamUtils.java b/binders/rabbit-binder/spring-cloud-stream-binder-rabbit/src/main/java/org/springframework/cloud/stream/binder/rabbit/StreamUtils.java
index 70a432509..aee1466ac 100644
--- a/binders/rabbit-binder/spring-cloud-stream-binder-rabbit/src/main/java/org/springframework/cloud/stream/binder/rabbit/StreamUtils.java
+++ b/binders/rabbit-binder/spring-cloud-stream-binder-rabbit/src/main/java/org/springframework/cloud/stream/binder/rabbit/StreamUtils.java
@@ -24,7 +24,6 @@ import com.rabbitmq.stream.Environment;
import org.springframework.amqp.core.MessageProperties;
import org.springframework.amqp.rabbit.listener.MessageListenerContainer;
import org.springframework.amqp.support.converter.MessageConverter;
-import org.springframework.cloud.stream.binder.BinderHeaders;
import org.springframework.cloud.stream.binder.ExtendedConsumerProperties;
import org.springframework.cloud.stream.binder.ExtendedProducerProperties;
import org.springframework.cloud.stream.binder.rabbit.properties.RabbitConsumerProperties;
@@ -34,12 +33,13 @@ import org.springframework.cloud.stream.provisioning.ConsumerDestination;
import org.springframework.cloud.stream.provisioning.ProducerDestination;
import org.springframework.context.support.AbstractApplicationContext;
import org.springframework.integration.amqp.inbound.AmqpInboundChannelAdapter;
-import org.springframework.integration.amqp.outbound.RabbitStreamMessageHandler;
import org.springframework.integration.amqp.support.AmqpHeaderMapper;
import org.springframework.integration.amqp.support.DefaultAmqpHeaderMapper;
import org.springframework.messaging.MessageChannel;
import org.springframework.messaging.MessageHandler;
+import org.springframework.messaging.MessageHandlingException;
import org.springframework.messaging.MessageHeaders;
+import org.springframework.messaging.support.ErrorMessage;
import org.springframework.rabbit.stream.listener.ConsumerCustomizer;
import org.springframework.rabbit.stream.listener.StreamListenerContainer;
import org.springframework.rabbit.stream.producer.RabbitStreamTemplate;
@@ -79,9 +79,7 @@ public final class StreamUtils {
@Override
public synchronized void setConsumerCustomizer(ConsumerCustomizer consumerCustomizer) {
super.setConsumerCustomizer((id, builder) -> {
- if (!properties.getExtension().isSuperStream()) {
- builder.name(consumerDestination.getName() + "." + group);
- }
+ builder.name(consumerDestination.getName() + "." + group);
consumerCustomizer.accept(id, builder);
});
}
@@ -93,9 +91,6 @@ public final class StreamUtils {
if (beanName != null) {
container.setStreamConverter(applicationContext.getBean(beanName, StreamMessageConverter.class));
}
- if (properties.getExtension().isSuperStream()) {
- container.superStream(consumerDestination.getName(), consumerDestination.getName() + "." + group);
- }
return container;
}
@@ -151,13 +146,6 @@ public final class StreamUtils {
RabbitStreamTemplate template = new RabbitStreamTemplate(applicationContext.getBean(Environment.class),
producerDestination.getName());
- if (extendedProperties.isSuperStream()) {
- template.setSuperStreamRouting(message -> {
- Object property = message.getApplicationProperties().getOrDefault(BinderHeaders.PARTITION_HEADER, "0");
- message.getApplicationProperties().remove(BinderHeaders.PARTITION_HEADER);
- return "" + property;
- });
- }
String beanName = extendedProperties.getStreamMessageConverterBeanName();
if (beanName != null) {
template.setMessageConverter(applicationContext.getBean(beanName, MessageConverter.class));
@@ -168,11 +156,9 @@ public final class StreamUtils {
}
RabbitStreamMessageHandler handler = new RabbitStreamMessageHandler(template);
if (errorChannel != null) {
- handler.setSendFailureChannel(errorChannel);
- }
- beanName = extendedProperties.getConfirmAckChannel();
- if (beanName != null) {
- handler.setSendSuccessChannelName(beanName);
+ handler.setFailureCallback((msg, ex) -> {
+ errorChannel.send(new ErrorMessage(new MessageHandlingException(msg, ex)));
+ });
}
handler.setHeaderMapper(headerMapperFunction.apply(extendedProperties));
handler.setSync(ProducerType.STREAM_SYNC.equals(producerProperties.getExtension().getProducerType()));
diff --git a/binders/rabbit-binder/spring-cloud-stream-binder-rabbit/src/test/java/org/springframework/cloud/stream/binder/rabbit/stream/RabbitStreamBinderModuleTests.java b/binders/rabbit-binder/spring-cloud-stream-binder-rabbit/src/test/java/org/springframework/cloud/stream/binder/rabbit/stream/RabbitStreamBinderModuleTests.java
index 5a7b6be7d..18f0f015f 100644
--- a/binders/rabbit-binder/spring-cloud-stream-binder-rabbit/src/test/java/org/springframework/cloud/stream/binder/rabbit/stream/RabbitStreamBinderModuleTests.java
+++ b/binders/rabbit-binder/spring-cloud-stream-binder-rabbit/src/test/java/org/springframework/cloud/stream/binder/rabbit/stream/RabbitStreamBinderModuleTests.java
@@ -22,10 +22,7 @@ import com.rabbitmq.stream.OffsetSpecification;
import com.rabbitmq.stream.ProducerBuilder;
import org.junit.jupiter.api.AfterEach;
import org.junit.jupiter.api.Test;
-import org.testcontainers.containers.RabbitMQContainer;
-import org.springframework.amqp.rabbit.connection.CachingConnectionFactory;
-import org.springframework.amqp.rabbit.connection.ConnectionFactory;
import org.springframework.amqp.rabbit.listener.MessageListenerContainer;
import org.springframework.boot.WebApplicationType;
import org.springframework.boot.autoconfigure.SpringBootApplication;
@@ -36,7 +33,7 @@ import org.springframework.cloud.stream.binder.Binding;
import org.springframework.cloud.stream.binder.ExtendedConsumerProperties;
import org.springframework.cloud.stream.binder.ExtendedProducerProperties;
import org.springframework.cloud.stream.binder.rabbit.RabbitMessageChannelBinder;
-import org.springframework.cloud.stream.binder.rabbit.RabbitTestContainer;
+import org.springframework.cloud.stream.binder.rabbit.RabbitStreamMessageHandler;
import org.springframework.cloud.stream.binder.rabbit.properties.RabbitConsumerProperties;
import org.springframework.cloud.stream.binder.rabbit.properties.RabbitConsumerProperties.ContainerType;
import org.springframework.cloud.stream.binder.rabbit.properties.RabbitProducerProperties;
@@ -45,7 +42,6 @@ import org.springframework.cloud.stream.config.ListenerContainerCustomizer;
import org.springframework.cloud.stream.config.ProducerMessageHandlerCustomizer;
import org.springframework.context.ConfigurableApplicationContext;
import org.springframework.context.annotation.Bean;
-import org.springframework.integration.amqp.outbound.RabbitStreamMessageHandler;
import org.springframework.integration.channel.DirectChannel;
import org.springframework.integration.channel.QueueChannel;
import org.springframework.integration.test.util.TestUtils;
@@ -54,7 +50,6 @@ import org.springframework.messaging.MessageHandler;
import org.springframework.rabbit.stream.listener.StreamListenerContainer;
import static org.assertj.core.api.Assertions.assertThat;
-import static org.mockito.ArgumentMatchers.anyString;
import static org.mockito.BDDMockito.given;
import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.verify;
@@ -64,8 +59,6 @@ import static org.mockito.Mockito.verify;
*/
public class RabbitStreamBinderModuleTests {
- private static final RabbitMQContainer RABBITMQ = RabbitTestContainer.sharedInstance();
-
private ConfigurableApplicationContext context;
@AfterEach
@@ -97,32 +90,6 @@ public class RabbitStreamBinderModuleTests {
((StreamListenerContainer) container).stop();
}
- @Test
- void testSuperStreamContainer() {
- context = new SpringApplicationBuilder(SimpleProcessor.class)
- .web(WebApplicationType.NONE)
- .run("--server.port=0");
- BinderFactory binderFactory = context.getBean(BinderFactory.class);
- RabbitMessageChannelBinder rabbitBinder = (RabbitMessageChannelBinder) binderFactory.getBinder(null,
- MessageChannel.class);
- RabbitConsumerProperties rProps = new RabbitConsumerProperties();
- rProps.setContainerType(ContainerType.STREAM);
- rProps.setSuperStream(true);
- ExtendedConsumerProperties props =
- new ExtendedConsumerProperties(rProps);
- props.setAutoStartup(false);
- props.setInstanceCount(1);
- Binding binding = rabbitBinder.bindConsumer("testSuperStream", "grp", new QueueChannel(), props);
- Object container = TestUtils.getPropertyValue(binding, "lifecycle.messageListenerContainer");
- assertThat(container).isInstanceOf(StreamListenerContainer.class);
- ((StreamListenerContainer) container).start();
- ConsumerBuilder builder = this.context.getBean(ConsumerBuilder.class);
- verify(builder).singleActiveConsumer();
- verify(builder).superStream("testSuperStream");
- verify(builder).name("testSuperStream.grp");
- ((StreamListenerContainer) container).stop();
- }
-
@Test
void testStreamHandler() {
context = new SpringApplicationBuilder(SimpleProcessor.class)
@@ -136,18 +103,13 @@ public class RabbitStreamBinderModuleTests {
ExtendedProducerProperties props =
new ExtendedProducerProperties(rProps);
Binding binding = rabbitBinder.bindProducer("testStream", new DirectChannel(), props);
- Object handler = TestUtils.getPropertyValue(binding, "val$producerMessageHandler");
+ Object handler = TestUtils.getPropertyValue(binding, "lifecycle");
assertThat(handler).isInstanceOf(RabbitStreamMessageHandler.class);
}
@SpringBootApplication(proxyBeanMethods = false)
public static class SimpleProcessor {
- @Bean
- ConnectionFactory cf() {
- return new CachingConnectionFactory(RABBITMQ.getMappedPort(5672));
- }
-
@Bean
ProducerMessageHandlerCustomizer handlerCustomizer() {
return (hand, dest) -> {
@@ -177,10 +139,7 @@ public class RabbitStreamBinderModuleTests {
@Bean
ConsumerBuilder consumerBuilder() {
- ConsumerBuilder mock = mock(ConsumerBuilder.class);
- given(mock.superStream(anyString())).willReturn(mock);
- given(mock.singleActiveConsumer()).willReturn(mock);
- return mock;
+ return mock(ConsumerBuilder.class);
}
@Bean
diff --git a/binders/rabbit-binder/spring-cloud-stream-binder-rabbit/src/test/java/org/springframework/cloud/stream/binder/rabbit/stream/RabbitStreamMessageHandlerTests.java b/binders/rabbit-binder/spring-cloud-stream-binder-rabbit/src/test/java/org/springframework/cloud/stream/binder/rabbit/stream/RabbitStreamMessageHandlerTests.java
new file mode 100644
index 000000000..575de2ccd
--- /dev/null
+++ b/binders/rabbit-binder/spring-cloud-stream-binder-rabbit/src/test/java/org/springframework/cloud/stream/binder/rabbit/stream/RabbitStreamMessageHandlerTests.java
@@ -0,0 +1,120 @@
+/*
+ * Copyright 2021-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.cloud.stream.binder.rabbit.stream;
+
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicReference;
+
+import com.rabbitmq.stream.Address;
+import com.rabbitmq.stream.Consumer;
+import com.rabbitmq.stream.Environment;
+import com.rabbitmq.stream.OffsetSpecification;
+import org.junit.jupiter.api.Disabled;
+import org.junit.jupiter.api.Test;
+import org.testcontainers.containers.RabbitMQContainer;
+
+import org.springframework.cloud.stream.binder.rabbit.RabbitStreamMessageHandler;
+import org.springframework.cloud.stream.binder.rabbit.RabbitTestContainer;
+import org.springframework.integration.support.MessageBuilder;
+import org.springframework.rabbit.stream.producer.RabbitStreamTemplate;
+
+import static org.assertj.core.api.Assertions.assertThat;
+
+/**
+ * @author Gary Russell
+ * @author Chris Bono
+ * @since 3.2
+ */
+public class RabbitStreamMessageHandlerTests {
+
+ private static final RabbitMQContainer RABBITMQ = RabbitTestContainer.sharedInstance();
+
+ @Test
+ void convertAndSend() throws InterruptedException {
+ Environment env = Environment.builder()
+ .lazyInitialization(true)
+ .addressResolver(add -> new Address("localhost", RABBITMQ.getMappedPort(5552)))
+ .build();
+ try {
+ env.deleteStream("stream.stream");
+ }
+ catch (Exception e) {
+ }
+ env.streamCreator().stream("stream.stream").create();
+ RabbitStreamTemplate streamTemplate = new RabbitStreamTemplate(env, "stream.stream");
+ RabbitStreamMessageHandler handler = new RabbitStreamMessageHandler(streamTemplate);
+ handler.setSync(true);
+ handler.handleMessage(MessageBuilder.withPayload("foo")
+ .setHeader("bar", "baz")
+ .build());
+ CountDownLatch latch = new CountDownLatch(1);
+ AtomicReference received = new AtomicReference<>();
+ Consumer consumer = env.consumerBuilder().stream("stream.stream")
+ .offset(OffsetSpecification.first())
+ .messageHandler((context, msg) -> {
+ received.set(msg);
+ latch.countDown();
+ })
+ .build();
+ assertThat(latch.await(10, TimeUnit.SECONDS)).isTrue();
+ assertThat(received.get()).isNotNull();
+ assertThat(received.get().getBodyAsBinary()).isEqualTo("foo".getBytes());
+ assertThat((String) received.get().getApplicationProperties().get("bar")).isEqualTo("baz");
+ consumer.close();
+ handler.stop();
+ }
+
+ @Test
+ @Disabled
+ void sendNative() throws InterruptedException {
+ Environment env = Environment.builder()
+ .lazyInitialization(true)
+ .build();
+ try {
+ env.deleteStream("stream.stream");
+ }
+ catch (Exception e) {
+ }
+ env.streamCreator().stream("stream.stream").create();
+ RabbitStreamTemplate streamTemplate = new RabbitStreamTemplate(env, "stream.stream");
+ RabbitStreamMessageHandler handler = new RabbitStreamMessageHandler(streamTemplate);
+ handler.setSync(true);
+ handler.handleMessage(MessageBuilder.withPayload(streamTemplate.messageBuilder()
+ .addData("foo".getBytes())
+ .applicationProperties().entry("bar", "baz")
+ .messageBuilder()
+ .build())
+ .build());
+ CountDownLatch latch = new CountDownLatch(1);
+ AtomicReference received = new AtomicReference<>();
+ Consumer consumer = env.consumerBuilder().stream("stream.stream")
+ .offset(OffsetSpecification.first())
+ .messageHandler((context, msg) -> {
+ received.set(msg);
+ latch.countDown();
+ })
+ .build();
+ assertThat(latch.await(10, TimeUnit.SECONDS)).isTrue();
+ assertThat(received.get()).isNotNull();
+ assertThat(received.get().getBodyAsBinary()).isEqualTo("foo".getBytes());
+ assertThat((String) received.get().getApplicationProperties().get("bar")).isEqualTo("baz");
+ consumer.close();
+ handler.stop();
+ }
+
+}
diff --git a/docs/src/main/asciidoc/rabbit/rabbit_overview.adoc b/docs/src/main/asciidoc/rabbit/rabbit_overview.adoc
index 3bcaa2854..a5a952ba0 100644
--- a/docs/src/main/asciidoc/rabbit/rabbit_overview.adoc
+++ b/docs/src/main/asciidoc/rabbit/rabbit_overview.adoc
@@ -423,12 +423,17 @@ To enable this feature, you must add the `spring-rabbit-stream` jar to the class
IMPORTANT: The consumer properties described above are not supported when you set the `containerType` property to `stream`; `concurrency` is also not supported at this time.
Only a single stream queue can be consumed by each binding.
-To configure the binder to use `containerType=stream`, Spring Boot will automatically configure an `Environment` `@Bean` from the application properties.
-You can, optionally, add a customizer to customize the listener container.
+To configure the binder to use `containerType=stream`, you must add an `Environment` `@Bean` and, optionally, a customizer to customize the listener container.
====
[source, java]
----
+@Bean
+Environment streamEnv() {
+ return Environment.builder()
+ .build();
+}
+
@Bean
ListenerContainerCustomizer customizer() {
return (cont, dest, group) -> {
@@ -468,39 +473,6 @@ public Consumer> input() {
Refer to the https://rabbitmq.github.io/rabbitmq-stream-java-client/stable/htmlsingle/[RabbitMQ Stream Java Client documentation] for information about configuring the environment and consumer builder.
-[[rabbitmq-super-stream-consumer]]
-==== Consumer Support for the RabbitMQ Super Streams
-
-See https://blog.rabbitmq.com/posts/2022/07/rabbitmq-3-11-feature-preview-super-streams[Super Streams] for information about super streams.
-
-Use of super streams allows for automatic scale-up scale-down with a single active consumer on each partition of a super stream.
-
-Configuration example:
-
-====
-[source, java]
-----
-@Bean
-public Consumer input() {
- ...
-}
-----
-====
-
-====
-[source, properties]
-----
-spring.cloud.stream.bindings.input-in-0.destination=super
-spring.cloud.stream.bindings.input-in-0.group=test
-spring.cloud.stream.bindings.input-in-0.consumer.instance-count=3
-spring.cloud.stream.rabbit.bindings.input-in-0.consumer.container-type=STREAM
-spring.cloud.stream.rabbit.bindings.input-in-0.consumer.super-stream=true
-----
-====
-
-The framework will create a super stream named `super`, with 3 partitions.
-Up to 3 instances of this application can be deployed.
-
=== Advanced Listener Container Configuration
To set listener container properties that are not exposed as binder or binding properties, add a single bean of type `ListenerContainerCustomizer` to the application context.
@@ -1129,74 +1101,6 @@ IMPORTANT: The correlation data must be provided with a unique `id` so that the
You cannot set both `useConfirmHeader` and `confirmAckChannel` but you can still receive returned messages in the error channel when `useConfirmHeader` is true, but using the correlation header is more convenient.
-[[rabbitmq-stream-producer]]
-=== Initial Producer Support for the RabbitMQ Stream Plugin
-
-Basic support for the https://rabbitmq.com/stream.html[RabbitMQ Stream Plugin] is now provided.
-To enable this feature, you must add the `spring-rabbit-stream` jar to the class path - it must be the same version as `spring-amqp` and `spring-rabbit`.
-
-IMPORTANT: The producer properties described above are not supported when you set the `producerType` property to `STREAM_SYNC` or `STREAM_ASYNC`.
-
-To configure the binder to use a stream `ProducerType`, Spring Boot will configure an `Environment` `@Bean` from the applicaation properties.
-You can, optionally, add a customizer to customize the message handler.
-
-====
-[source, java]
-----
-@Bean
-ProducerMessageHandlerCustomizer handlerCustomizer() {
- return (hand, dest) -> {
- RabbitStreamMessageHandler handler = (RabbitStreamMessageHandler) hand;
- handler.setConfirmTimeout(5000);
- ((RabbitStreamTemplate) handler.getStreamOperations()).setProducerCustomizer(
- (name, builder) -> {
- ...
- });
- };
-}
-----
-====
-
-Refer to the https://rabbitmq.github.io/rabbitmq-stream-java-client/stable/htmlsingle/[RabbitMQ Stream Java Client documentation] for information about configuring the environment and producer builder.
-
-[[rabbitmq-super-stream-producer]]
-==== Producer Support for the RabbitMQ Super Streams
-
-See https://blog.rabbitmq.com/posts/2022/07/rabbitmq-3-11-feature-preview-super-streams[Super Streams] for information about super streams.
-
-Use of super streams allows for automatic scale-up scale-down with a single active consumer on each partition of a super stream.
-Using Spring Cloud Stream, you can publish to a super stream either over AMQP, or using the stream client.
-
-IMPORTANT: The super stream must already exist; creating a super stream is not supported by producer bindings.
-
-Publishing to a super stream over AMQP:
-
-====
-[source, properties]
-----
-spring.cloud.stream.bindings.output.destination=super
-spring.cloud.stream.bindings.output.producer.partition-count=3
-spring.cloud.stream.bindings.output.producer.partition-key-expression=headers['cust-no']
-spring.cloud.stream.rabbit.bindings.output.producer.declare-exchange=false
-----
-====
-
-Publishing to a super stream using the stream client:
-
-====
-[source, properties]
-----
-spring.cloud.stream.bindings.output.destination=super
-spring.cloud.stream.bindings.output.producer.partition-count=3
-spring.cloud.stream.bindings.output.producer.partition-key-expression=headers['cust-no']
-spring.cloud.stream.rabbit.bindings.output.producer.producer-type=stream-async
-spring.cloud.stream.rabbit.bindings.output.producer.super-stream=true
-spring.cloud.stream.rabbit.bindings.output.producer.declare-exchange=false
-----
-====
-
-When using the stream client, if you set a `confirmAckChannel`, a copy of a successfully sent message will be sent to that channel.
-
== Using Existing Queues/Exchanges
By default, the binder will automatically provision a topic exchange with the name being derived from the value of the destination binding property ``.
@@ -1338,6 +1242,41 @@ For negatively acknowledged confirmations, the payload is a `NackedAmqpMessageEx
There is no automatic handling of these exceptions (such as sending to a <>).
You can consume these exceptions with your own Spring Integration flow.
+[[rabbitmq-stream-producer]]
+=== Initial Producer Support for the RabbitMQ Stream Plugin
+
+Basic support for the https://rabbitmq.com/stream.html[RabbitMQ Stream Plugin] is now provided.
+To enable this feature, you must add the `spring-rabbit-stream` jar to the class path - it must be the same version as `spring-amqp` and `spring-rabbit`.
+
+IMPORTANT: The producer properties described above are not supported when you set the `producerType` property to `STREAM_SYNC` or `STREAM_ASYNC`.
+
+To configure the binder to use a stream `ProducerType`, you must add an `Environment` `@Bean` and, optionally, a customizer to customize the message handler.
+
+====
+[source, java]
+----
+@Bean
+Environment streamEnv() {
+ return Environment.builder()
+ .build();
+}
+
+@Bean
+ProducerMessageHandlerCustomizer handlerCustomizer() {
+ return (hand, dest) -> {
+ RabbitStreamMessageHandler handler = (RabbitStreamMessageHandler) hand;
+ handler.setConfirmTimeout(5000);
+ ((RabbitStreamTemplate) handler.getStreamOperations()).setProducerCustomizer(
+ (name, builder) -> {
+ ...
+ });
+ };
+}
+----
+====
+
+Refer to the https://rabbitmq.github.io/rabbitmq-stream-java-client/stable/htmlsingle/[RabbitMQ Stream Java Client documentation] for information about configuring the environment and producer builder.
+=======
[[rabbit-binder-health-indicator]]
== Rabbit Binder Health Indicator