From a2595a26e3b785031d46cc6bae4e02fb6b92eb81 Mon Sep 17 00:00:00 2001 From: Oleg Zhurakousky Date: Fri, 1 Nov 2019 12:44:16 +0100 Subject: [PATCH] GH-1835 Fixed channel-to-publisher adapter - Needed to copy one from SI to address the issue but the actual fix should go into SI Resolves #1835 --- .../function/FunctionConfiguration.java | 31 +++++++++++++++++-- 1 file changed, 28 insertions(+), 3 deletions(-) diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/FunctionConfiguration.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/FunctionConfiguration.java index 1ba3a8985..3acfc06e2 100644 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/FunctionConfiguration.java +++ b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/FunctionConfiguration.java @@ -34,7 +34,9 @@ import java.util.stream.Stream; import org.apache.commons.logging.Log; import org.apache.commons.logging.LogFactory; import org.reactivestreams.Publisher; +import org.reactivestreams.Subscriber; import reactor.core.publisher.Flux; +import reactor.core.publisher.FluxSink; import reactor.core.publisher.Mono; import reactor.core.publisher.MonoSink; import reactor.util.function.Tuples; @@ -82,7 +84,6 @@ import org.springframework.core.annotation.AnnotationUtils; import org.springframework.core.env.Environment; import org.springframework.core.type.MethodMetadata; import org.springframework.integration.channel.AbstractMessageChannel; -import org.springframework.integration.channel.MessageChannelReactiveUtils; import org.springframework.integration.dsl.IntegrationFlow; import org.springframework.integration.dsl.IntegrationFlowBuilder; import org.springframework.integration.dsl.IntegrationFlows; @@ -91,6 +92,7 @@ import org.springframework.integration.support.MessageBuilder; import org.springframework.lang.Nullable; import org.springframework.messaging.Message; import org.springframework.messaging.MessageChannel; +import org.springframework.messaging.MessageHandler; import org.springframework.messaging.MessageHeaders; import org.springframework.messaging.SubscribableChannel; import org.springframework.util.Assert; @@ -364,7 +366,7 @@ public class FunctionConfiguration { Publisher[] inputPublishers = inputBindingNames.stream().map(inputBindingName -> { SubscribableChannel inputChannel = this.context.getBean(inputBindingName, SubscribableChannel.class); - return this.enhancePublisher(MessageChannelReactiveUtils.toPublisher(inputChannel), inputBindingName); + return this.enhancePublisher(new SubscribableChannelPublisherAdapter(inputChannel), inputBindingName); }).toArray(Publisher[]::new); @@ -425,7 +427,7 @@ public class FunctionConfiguration { if (FunctionTypeUtils.isReactive(FunctionTypeUtils.getInputType(functionType, 0)) && StringUtils.hasText(outputChannelName)) { MessageChannel outputChannel = context.getBean(outputChannelName, MessageChannel.class); SubscribableChannel subscribeChannel = (SubscribableChannel) inputChannel; - Publisher publisher = this.enhancePublisher(MessageChannelReactiveUtils.toPublisher(subscribeChannel), + Publisher publisher = this.enhancePublisher(new SubscribableChannelPublisherAdapter<>(subscribeChannel), ((DirectWithAttributesChannel) inputChannel).getBeanName()); this.subscribeToInput(function, publisher, message -> outputChannel.send((Message) message)); } @@ -677,4 +679,27 @@ public class FunctionConfiguration { } } + + private static final class SubscribableChannelPublisherAdapter implements Publisher> { + + private final SubscribableChannel channel; + + SubscribableChannelPublisherAdapter(SubscribableChannel channel) { + this.channel = channel; + } + + @Override + @SuppressWarnings("unchecked") + public void subscribe(Subscriber> subscriber) { + Flux. + >push(emitter -> { + MessageHandler messageHandler = emitter::next; + this.channel.subscribe(messageHandler); + emitter.onCancel(() -> this.channel.unsubscribe(messageHandler)); + }, + FluxSink.OverflowStrategy.BUFFER) + .subscribe((Subscriber>) subscriber); + } + + } }