diff --git a/build.gradle b/build.gradle
index 805079343d..f1c2f5e353 100644
--- a/build.gradle
+++ b/build.gradle
@@ -293,10 +293,7 @@ project('spring-integration-core') {
exclude group: 'org.springframework', module: 'spring-core'
}
// compile ("org.springframework.cloud:spring-cloud-cluster-core:$springCloudClusterVersion", optional)
- compile ("io.projectreactor:reactor-stream:$reactorVersion") {
- optional it
- exclude group: 'org.slf4j', module: 'slf4j-api'
- }
+ compile "io.projectreactor:reactor-core:$reactorVersion"
compile("com.fasterxml.jackson.core:jackson-databind:$jackson2Version", optional)
compile("com.jayway.jsonpath:json-path:$jsonpathVersion") {
optional it
@@ -309,8 +306,7 @@ project('spring-integration-core') {
// testCompile ("net.openhft:chronicle:$chronicleVersion")
// testCompile ("io.projectreactor:reactor-chronicle:$reactorVersion")
- testCompile "org.testng:testng:6.8.21"
- testCompile "org.reactivestreams:reactive-streams-tck:1.0.0"
+ testCompile "io.projectreactor:reactor-stream:$reactorVersion"
}
}
diff --git a/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/config/AmqpInboundGatewayParserTests-context.xml b/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/config/AmqpInboundGatewayParserTests-context.xml
index b7da686bf8..f8c7e18ee1 100644
--- a/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/config/AmqpInboundGatewayParserTests-context.xml
+++ b/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/config/AmqpInboundGatewayParserTests-context.xml
@@ -12,7 +12,9 @@
+ connection-factory="rabbitConnectionFactory" message-converter="testConverter">
+
+
diff --git a/spring-integration-core/src/main/java/org/springframework/integration/channel/ReactiveMessageChannel.java b/spring-integration-core/src/main/java/org/springframework/integration/channel/ReactiveChannel.java
similarity index 82%
rename from spring-integration-core/src/main/java/org/springframework/integration/channel/ReactiveMessageChannel.java
rename to spring-integration-core/src/main/java/org/springframework/integration/channel/ReactiveChannel.java
index ef521cd897..7b44f5c9db 100644
--- a/spring-integration-core/src/main/java/org/springframework/integration/channel/ReactiveMessageChannel.java
+++ b/spring-integration-core/src/main/java/org/springframework/integration/channel/ReactiveChannel.java
@@ -23,24 +23,24 @@ import org.reactivestreams.Subscriber;
import org.springframework.messaging.Message;
import org.springframework.messaging.MessageChannel;
-import reactor.core.subscription.ReactiveSession;
-import reactor.rx.broadcast.Broadcaster;
+import reactor.core.publisher.ProcessorGroup;
+import reactor.core.subscriber.ReactiveSession;
/**
* @author Artem Bilan
* @since 5.0
*/
-public class ReactiveMessageChannel implements MessageChannel, Publisher> {
+public class ReactiveChannel implements MessageChannel, Publisher> {
private final Processor, Message>> processor;
private final ReactiveSession> reactiveSession;
- public ReactiveMessageChannel() {
- this(Broadcaster.passthrough());
+ public ReactiveChannel() {
+ this(ProcessorGroup.>sync().get());
}
- public ReactiveMessageChannel(Processor, Message>> processor) {
+ public ReactiveChannel(Processor, Message>> processor) {
this.processor = processor;
this.reactiveSession = ReactiveSession.create(processor);
}
diff --git a/spring-integration-core/src/main/java/org/springframework/integration/config/ConsumerEndpointFactoryBean.java b/spring-integration-core/src/main/java/org/springframework/integration/config/ConsumerEndpointFactoryBean.java
index fa919b454e..06bbbb20b0 100644
--- a/spring-integration-core/src/main/java/org/springframework/integration/config/ConsumerEndpointFactoryBean.java
+++ b/spring-integration-core/src/main/java/org/springframework/integration/config/ConsumerEndpointFactoryBean.java
@@ -21,7 +21,6 @@ import java.util.List;
import org.aopalliance.aop.Advice;
import org.apache.commons.logging.Log;
import org.apache.commons.logging.LogFactory;
-import org.reactivestreams.Publisher;
import org.reactivestreams.Subscriber;
import org.springframework.aop.framework.Advised;
@@ -56,7 +55,8 @@ import org.springframework.util.Assert;
import org.springframework.util.CollectionUtils;
import org.springframework.util.StringUtils;
-import reactor.core.subscriber.SubscriberFactory;
+import reactor.core.subscriber.Subscribers;
+
/**
* @author Mark Fisher
@@ -289,20 +289,16 @@ public class ConsumerEndpointFactoryBean
pollingConsumer.setBeanFactory(this.beanFactory);
this.endpoint = pollingConsumer;
}
- else if (channel instanceof Publisher) {
- Publisher> publisher = (Publisher>) channel;
+ else {
Subscriber> subscriber;
if (this.handler instanceof Subscriber) {
subscriber = (Subscriber>) this.handler;
}
else {
//TODO errorConsumer, completeConsumer
- subscriber = SubscriberFactory.consumer(this.handler::handleMessage);
+ subscriber = Subscribers.consumer(this.handler::handleMessage);
}
- this.endpoint = new ReactiveEndpoint(publisher, subscriber);
- }
- else {
- throw new IllegalArgumentException("unsupported channel type: [" + channel.getClass() + "]");
+ this.endpoint = new ReactiveEndpoint(channel, subscriber);
}
this.endpoint.setBeanName(this.beanName);
this.endpoint.setBeanFactory(this.beanFactory);
diff --git a/spring-integration-core/src/main/java/org/springframework/integration/config/annotation/AbstractMethodAnnotationPostProcessor.java b/spring-integration-core/src/main/java/org/springframework/integration/config/annotation/AbstractMethodAnnotationPostProcessor.java
index 32edf75988..4e5c7d9900 100644
--- a/spring-integration-core/src/main/java/org/springframework/integration/config/annotation/AbstractMethodAnnotationPostProcessor.java
+++ b/spring-integration-core/src/main/java/org/springframework/integration/config/annotation/AbstractMethodAnnotationPostProcessor.java
@@ -80,7 +80,7 @@ import org.springframework.util.CollectionUtils;
import org.springframework.util.ObjectUtils;
import org.springframework.util.StringUtils;
-import reactor.core.subscriber.SubscriberFactory;
+import reactor.core.subscriber.Subscribers;
/**
* Base class for Method-level annotation post-processors.
@@ -306,7 +306,8 @@ public abstract class AbstractMethodAnnotationPostProcessor annotations) {
+ @SuppressWarnings("unchecked")
+ protected AbstractEndpoint doCreateEndpoint(MessageHandler handler, MessageChannel inputChannel,List annotations) {
AbstractEndpoint endpoint;
if (inputChannel instanceof PollableChannel) {
PollingConsumer pollingConsumer = new PollingConsumer((PollableChannel) inputChannel, handler);
@@ -318,16 +319,15 @@ public abstract class AbstractMethodAnnotationPostProcessor> publisher = (Publisher>) inputChannel;
Subscriber> subscriber;
if (handler instanceof Subscriber) {
subscriber = (Subscriber>) handler;
}
else {
//TODO errorConsumer, completeConsumer
- subscriber = SubscriberFactory.consumer(handler::handleMessage);
+ subscriber = Subscribers.consumer(handler::handleMessage);
}
- endpoint = new ReactiveEndpoint(publisher, subscriber);
+ endpoint = new ReactiveEndpoint(inputChannel, subscriber);
}
else {
endpoint = new EventDrivenConsumer((SubscribableChannel) inputChannel, handler);
diff --git a/spring-integration-core/src/main/java/org/springframework/integration/endpoint/ReactiveEndpoint.java b/spring-integration-core/src/main/java/org/springframework/integration/endpoint/ReactiveEndpoint.java
index e681c17919..b6e5bd994e 100644
--- a/spring-integration-core/src/main/java/org/springframework/integration/endpoint/ReactiveEndpoint.java
+++ b/spring-integration-core/src/main/java/org/springframework/integration/endpoint/ReactiveEndpoint.java
@@ -20,9 +20,11 @@ import org.reactivestreams.Publisher;
import org.reactivestreams.Subscriber;
import org.springframework.messaging.Message;
+import org.springframework.messaging.MessageChannel;
import org.springframework.util.Assert;
-import reactor.core.subscription.ReactiveSession;
+import reactor.core.subscriber.ReactiveSession;
+
/**
* @author Artem Bilan
@@ -36,11 +38,17 @@ public class ReactiveEndpoint extends AbstractEndpoint {
private ReactiveSession> reactiveSession;
- public ReactiveEndpoint(Publisher> inputChannel, Subscriber> subscriber) {
- Assert.isInstanceOf(Publisher.class, inputChannel,
- "The 'inputChannel', must implement org.reactivestreams.Publisher.");
+ @SuppressWarnings("unchecked")
+ public ReactiveEndpoint(MessageChannel inputChannel, Subscriber> subscriber) {
+ Assert.notNull(inputChannel);
Assert.notNull(subscriber);
- this.inputChannel = inputChannel;
+ if (inputChannel instanceof Publisher) {
+ this.inputChannel = (Publisher>) inputChannel;
+ }
+ else {
+ //TODO: Wrap all other channels to the Publisher>
+ this.inputChannel = null;
+ }
this.subscriber = subscriber;
}
diff --git a/spring-integration-core/src/main/java/org/springframework/integration/gateway/GatewayProxyFactoryBean.java b/spring-integration-core/src/main/java/org/springframework/integration/gateway/GatewayProxyFactoryBean.java
index d0d33ce2d4..b14913a17a 100644
--- a/spring-integration-core/src/main/java/org/springframework/integration/gateway/GatewayProxyFactoryBean.java
+++ b/spring-integration-core/src/main/java/org/springframework/integration/gateway/GatewayProxyFactoryBean.java
@@ -66,7 +66,8 @@ import org.springframework.util.ObjectUtils;
import org.springframework.util.ReflectionUtils;
import org.springframework.util.StringUtils;
-import reactor.rx.Promises;
+import reactor.core.publisher.Mono;
+
/**
* Generates a proxy for the provided service interface to enable interaction
@@ -400,7 +401,7 @@ public class GatewayProxyFactoryBean extends AbstractEndpoint
}
}
if (reactorPresent && Publisher.class.isAssignableFrom(returnType)) {
- return Promises.