GH-1835 Changed channel-to-publisher adapter logic to use EmitterProcessor
This effectively transforms Publisher to back-pressure honoring publisher polishing
This commit is contained in:
@@ -34,9 +34,8 @@ import java.util.stream.Stream;
|
|||||||
import org.apache.commons.logging.Log;
|
import org.apache.commons.logging.Log;
|
||||||
import org.apache.commons.logging.LogFactory;
|
import org.apache.commons.logging.LogFactory;
|
||||||
import org.reactivestreams.Publisher;
|
import org.reactivestreams.Publisher;
|
||||||
import org.reactivestreams.Subscriber;
|
import reactor.core.publisher.EmitterProcessor;
|
||||||
import reactor.core.publisher.Flux;
|
import reactor.core.publisher.Flux;
|
||||||
import reactor.core.publisher.FluxSink;
|
|
||||||
import reactor.core.publisher.Mono;
|
import reactor.core.publisher.Mono;
|
||||||
import reactor.core.publisher.MonoSink;
|
import reactor.core.publisher.MonoSink;
|
||||||
import reactor.util.function.Tuples;
|
import reactor.util.function.Tuples;
|
||||||
@@ -92,7 +91,6 @@ import org.springframework.integration.support.MessageBuilder;
|
|||||||
import org.springframework.lang.Nullable;
|
import org.springframework.lang.Nullable;
|
||||||
import org.springframework.messaging.Message;
|
import org.springframework.messaging.Message;
|
||||||
import org.springframework.messaging.MessageChannel;
|
import org.springframework.messaging.MessageChannel;
|
||||||
import org.springframework.messaging.MessageHandler;
|
|
||||||
import org.springframework.messaging.MessageHeaders;
|
import org.springframework.messaging.MessageHeaders;
|
||||||
import org.springframework.messaging.SubscribableChannel;
|
import org.springframework.messaging.SubscribableChannel;
|
||||||
import org.springframework.util.Assert;
|
import org.springframework.util.Assert;
|
||||||
@@ -366,7 +364,7 @@ public class FunctionConfiguration {
|
|||||||
|
|
||||||
Publisher[] inputPublishers = inputBindingNames.stream().map(inputBindingName -> {
|
Publisher[] inputPublishers = inputBindingNames.stream().map(inputBindingName -> {
|
||||||
SubscribableChannel inputChannel = this.context.getBean(inputBindingName, SubscribableChannel.class);
|
SubscribableChannel inputChannel = this.context.getBean(inputBindingName, SubscribableChannel.class);
|
||||||
return this.enhancePublisher(new SubscribableChannelPublisherAdapter(inputChannel), inputBindingName);
|
return this.enhancePublisher(this.convertToPublisher(inputChannel), inputBindingName);
|
||||||
}).toArray(Publisher[]::new);
|
}).toArray(Publisher[]::new);
|
||||||
|
|
||||||
|
|
||||||
@@ -427,7 +425,7 @@ public class FunctionConfiguration {
|
|||||||
if (FunctionTypeUtils.isReactive(FunctionTypeUtils.getInputType(functionType, 0)) && StringUtils.hasText(outputChannelName)) {
|
if (FunctionTypeUtils.isReactive(FunctionTypeUtils.getInputType(functionType, 0)) && StringUtils.hasText(outputChannelName)) {
|
||||||
MessageChannel outputChannel = context.getBean(outputChannelName, MessageChannel.class);
|
MessageChannel outputChannel = context.getBean(outputChannelName, MessageChannel.class);
|
||||||
SubscribableChannel subscribeChannel = (SubscribableChannel) inputChannel;
|
SubscribableChannel subscribeChannel = (SubscribableChannel) inputChannel;
|
||||||
Publisher<?> publisher = this.enhancePublisher(new SubscribableChannelPublisherAdapter<>(subscribeChannel),
|
Publisher<?> publisher = this.enhancePublisher(this.convertToPublisher(inputChannel),
|
||||||
((DirectWithAttributesChannel) inputChannel).getBeanName());
|
((DirectWithAttributesChannel) inputChannel).getBeanName());
|
||||||
this.subscribeToInput(function, publisher, message -> outputChannel.send((Message<?>) message));
|
this.subscribeToInput(function, publisher, message -> outputChannel.send((Message<?>) message));
|
||||||
}
|
}
|
||||||
@@ -545,6 +543,14 @@ public class FunctionConfiguration {
|
|||||||
return bindableProxyFactory instanceof BindableFunctionProxyFactory
|
return bindableProxyFactory instanceof BindableFunctionProxyFactory
|
||||||
&& ((BindableFunctionProxyFactory) bindableProxyFactory).isMultiple();
|
&& ((BindableFunctionProxyFactory) bindableProxyFactory).isMultiple();
|
||||||
}
|
}
|
||||||
|
|
||||||
|
private Publisher<Message<?>> convertToPublisher(SubscribableChannel inputChannel) {
|
||||||
|
EmitterProcessor<Message<?>> publisher = EmitterProcessor.create(1);
|
||||||
|
inputChannel.subscribe(message -> {
|
||||||
|
publisher.onNext(message);
|
||||||
|
});
|
||||||
|
return publisher;
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
@@ -679,27 +685,4 @@ public class FunctionConfiguration {
|
|||||||
}
|
}
|
||||||
|
|
||||||
}
|
}
|
||||||
|
|
||||||
private static final class SubscribableChannelPublisherAdapter<T> implements Publisher<Message<T>> {
|
|
||||||
|
|
||||||
private final SubscribableChannel channel;
|
|
||||||
|
|
||||||
SubscribableChannelPublisherAdapter(SubscribableChannel channel) {
|
|
||||||
this.channel = channel;
|
|
||||||
}
|
|
||||||
|
|
||||||
@Override
|
|
||||||
@SuppressWarnings("unchecked")
|
|
||||||
public void subscribe(Subscriber<? super Message<T>> subscriber) {
|
|
||||||
Flux.
|
|
||||||
<Message<?>>push(emitter -> {
|
|
||||||
MessageHandler messageHandler = emitter::next;
|
|
||||||
this.channel.subscribe(messageHandler);
|
|
||||||
emitter.onCancel(() -> this.channel.unsubscribe(messageHandler));
|
|
||||||
},
|
|
||||||
FluxSink.OverflowStrategy.BUFFER)
|
|
||||||
.subscribe((Subscriber<? super Message<?>>) subscriber);
|
|
||||||
}
|
|
||||||
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user