From b7ee269cc9055afa631b9bd5f6a52be0d0a421c0 Mon Sep 17 00:00:00 2001 From: Artem Bilan Date: Fri, 1 Nov 2019 14:49:29 -0400 Subject: [PATCH] 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 --- .../channel/MessageChannelReactiveUtils.java | 34 ++++---------- .../endpoint/ReactiveStreamsConsumer.java | 2 +- .../ReactiveStreamsConsumerTests.java | 46 +++++++++++-------- 3 files changed, 36 insertions(+), 46 deletions(-) diff --git a/spring-integration-core/src/main/java/org/springframework/integration/channel/MessageChannelReactiveUtils.java b/spring-integration-core/src/main/java/org/springframework/integration/channel/MessageChannelReactiveUtils.java index 5f558d9550..216bc1be23 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/channel/MessageChannelReactiveUtils.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/channel/MessageChannelReactiveUtils.java @@ -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 Publisher> adaptSubscribableChannelToPublisher(SubscribableChannel inputChannel) { - return new SubscribableChannelPublisherAdapter<>(inputChannel); + return Flux.defer(() -> { + EmitterProcessor> publisher = EmitterProcessor.create(1); + @SuppressWarnings("unchecked") + MessageHandler messageHandler = (message) -> publisher.onNext((Message) message); + inputChannel.subscribe(messageHandler); + return publisher + .doOnCancel(() -> inputChannel.unsubscribe(messageHandler)); + }); } private static Publisher> adaptPollableChannelToPublisher(PollableChannel inputChannel) { return new PollableChannelPublisherAdapter<>(inputChannel); } - - 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. - >create(emitter -> { - MessageHandler messageHandler = emitter::next; - this.channel.subscribe(messageHandler); - emitter.onCancel(() -> this.channel.unsubscribe(messageHandler)); - }, - FluxSink.OverflowStrategy.IGNORE) - .subscribe((Subscriber>) subscriber); - } - - } - private static final class PollableChannelPublisherAdapter implements Publisher> { private final PollableChannel channel; diff --git a/spring-integration-core/src/main/java/org/springframework/integration/endpoint/ReactiveStreamsConsumer.java b/spring-integration-core/src/main/java/org/springframework/integration/endpoint/ReactiveStreamsConsumer.java index 0591131a4f..9d5e19cae0 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/endpoint/ReactiveStreamsConsumer.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/endpoint/ReactiveStreamsConsumer.java @@ -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 diff --git a/spring-integration-core/src/test/java/org/springframework/integration/channel/reactive/ReactiveStreamsConsumerTests.java b/spring-integration-core/src/test/java/org/springframework/integration/channel/reactive/ReactiveStreamsConsumerTests.java index aa1968e284..94ee2bc7de 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/channel/reactive/ReactiveStreamsConsumerTests.java +++ b/spring-integration-core/src/test/java/org/springframework/integration/channel/reactive/ReactiveStreamsConsumerTests.java @@ -108,35 +108,43 @@ public class ReactiveStreamsConsumerTests { @Test - @SuppressWarnings("unchecked") public void testReactiveStreamsConsumerDirectChannel() throws InterruptedException { DirectChannel testChannel = new DirectChannel(); - Subscriber> testSubscriber = (Subscriber>) Mockito.mock(Subscriber.class); - BlockingQueue> messages = new LinkedBlockingQueue<>(); - willAnswer(i -> { - messages.put(i.getArgument(0)); - return null; - }) - .given(testSubscriber) - .onNext(any(Message.class)); + Subscriber> testSubscriber = Mockito.spy(new Subscriber>() { + + @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 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