Use EmitterProcessor for Channels adaptation (#3100)

* Use `EmitterProcessor` for Channels adaptation

Related https://github.com/spring-cloud/spring-cloud-stream/issues/1835

To honor a back-pressure after `MessageChannel` adaptation it is better
to use an `EmitterProcessor.create(1)` instead of `Flux.create()`.
This way whenever an emitter buffer is full, we block upstream producer
and don't allow it to produce more messages

**Cherry-pick to 5.1.x**

* * Wrap every new subscription into a `Flux.defer()`
* Fix `ReactiveStreamsConsumerTests` to use a new `Subscription` after
each `stop()/start()` on the `ReactiveStreamsConsumer`

* * Remove unused imports

# Conflicts:
#	spring-integration-core/src/main/java/org/springframework/integration/endpoint/ReactiveStreamsConsumer.java
#	spring-integration-core/src/test/java/org/springframework/integration/channel/MessageChannelReactiveUtilsTests.java
#	spring-integration-core/src/test/java/org/springframework/integration/channel/reactive/ReactiveStreamsConsumerTests.java

* Fixing conflicts in tests
This commit is contained in:
Artem Bilan
2019-11-01 14:49:29 -04:00
parent 54de7a2209
commit b7ee269cc9
3 changed files with 36 additions and 46 deletions

View File

@@ -25,6 +25,7 @@ import org.springframework.messaging.MessageHandler;
import org.springframework.messaging.PollableChannel;
import org.springframework.messaging.SubscribableChannel;
import reactor.core.publisher.EmitterProcessor;
import reactor.core.publisher.Flux;
import reactor.core.publisher.FluxSink;
import reactor.core.scheduler.Schedulers;
@@ -60,37 +61,20 @@ public final class MessageChannelReactiveUtils {
}
private static <T> Publisher<Message<T>> adaptSubscribableChannelToPublisher(SubscribableChannel inputChannel) {
return new SubscribableChannelPublisherAdapter<>(inputChannel);
return Flux.defer(() -> {
EmitterProcessor<Message<T>> publisher = EmitterProcessor.create(1);
@SuppressWarnings("unchecked")
MessageHandler messageHandler = (message) -> publisher.onNext((Message<T>) message);
inputChannel.subscribe(messageHandler);
return publisher
.doOnCancel(() -> inputChannel.unsubscribe(messageHandler));
});
}
private static <T> Publisher<Message<T>> adaptPollableChannelToPublisher(PollableChannel inputChannel) {
return new PollableChannelPublisherAdapter<>(inputChannel);
}
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<?>>create(emitter -> {
MessageHandler messageHandler = emitter::next;
this.channel.subscribe(messageHandler);
emitter.onCancel(() -> this.channel.unsubscribe(messageHandler));
},
FluxSink.OverflowStrategy.IGNORE)
.subscribe((Subscriber<? super Message<?>>) subscriber);
}
}
private static final class PollableChannelPublisherAdapter<T> implements Publisher<Message<T>> {
private final PollableChannel channel;

View File

@@ -140,8 +140,8 @@ public class ReactiveStreamsConsumer extends AbstractEndpoint implements Integra
@Override
public void hookOnSubscribe(Subscription s) {
this.delegate.onSubscribe(s);
ReactiveStreamsConsumer.this.subscription = s;
this.delegate.onSubscribe(s);
}
@Override

View File

@@ -108,35 +108,43 @@ public class ReactiveStreamsConsumerTests {
@Test
@SuppressWarnings("unchecked")
public void testReactiveStreamsConsumerDirectChannel() throws InterruptedException {
DirectChannel testChannel = new DirectChannel();
Subscriber<Message<?>> testSubscriber = (Subscriber<Message<?>>) Mockito.mock(Subscriber.class);
BlockingQueue<Message<?>> messages = new LinkedBlockingQueue<>();
willAnswer(i -> {
messages.put(i.getArgument(0));
return null;
})
.given(testSubscriber)
.onNext(any(Message.class));
Subscriber<Message<?>> testSubscriber = Mockito.spy(new Subscriber<Message<?>>() {
@Override
public void onSubscribe(Subscription subscription) {
subscription.request(1);
}
@Override
public void onNext(Message<?> message) {
messages.offer(message);
}
@Override
public void onError(Throwable t) {
}
@Override
public void onComplete() {
}
});
ReactiveStreamsConsumer reactiveConsumer = new ReactiveStreamsConsumer(testChannel, testSubscriber);
reactiveConsumer.setBeanFactory(mock(BeanFactory.class));
reactiveConsumer.afterPropertiesSet();
reactiveConsumer.start();
Message<?> testMessage = new GenericMessage<>("test");
final Message<?> testMessage = new GenericMessage<>("test");
testChannel.send(testMessage);
ArgumentCaptor<Subscription> subscriptionArgumentCaptor = ArgumentCaptor.forClass(Subscription.class);
verify(testSubscriber).onSubscribe(subscriptionArgumentCaptor.capture());
Subscription subscription = subscriptionArgumentCaptor.getValue();
subscription.request(1);
Message<?> message = messages.poll(10, TimeUnit.SECONDS);
assertSame(testMessage, message);
@@ -152,10 +160,6 @@ public class ReactiveStreamsConsumerTests {
reactiveConsumer.start();
subscription.request(1);
testMessage = new GenericMessage<>("test2");
testChannel.send(testMessage);
message = messages.poll(10, TimeUnit.SECONDS);
@@ -165,6 +169,8 @@ public class ReactiveStreamsConsumerTests {
verify(testSubscriber, never()).onComplete();
assertTrue(messages.isEmpty());
reactiveConsumer.stop();
}
@Test