From 3874620d2a98857f14423bc94924232ebefa1695 Mon Sep 17 00:00:00 2001 From: Marius Bogoevici Date: Sun, 14 Feb 2016 23:30:05 -0500 Subject: [PATCH] Fix RxJava processor lifecycle - Make SubjectMessageHandler a SmartLifecycle and assign it to phase 0, allowing it to be started after the outputs are bound and before the inputs are bound; - Setup Rx in the start() method; - Set the ServiceActivator associated with the handler to phase 0 as well; --- .../rxjava/RxJavaProcessorConfiguration.java | 2 +- .../rxjava/SubjectMessageHandler.java | 98 +++++++++++++------ 2 files changed, 70 insertions(+), 30 deletions(-) diff --git a/spring-cloud-stream-rxjava/src/main/java/org/springframework/cloud/stream/annotation/rxjava/RxJavaProcessorConfiguration.java b/spring-cloud-stream-rxjava/src/main/java/org/springframework/cloud/stream/annotation/rxjava/RxJavaProcessorConfiguration.java index 7524839d3..03f798e18 100644 --- a/spring-cloud-stream-rxjava/src/main/java/org/springframework/cloud/stream/annotation/rxjava/RxJavaProcessorConfiguration.java +++ b/spring-cloud-stream-rxjava/src/main/java/org/springframework/cloud/stream/annotation/rxjava/RxJavaProcessorConfiguration.java @@ -33,7 +33,7 @@ public class RxJavaProcessorConfiguration { @Autowired RxJavaProcessor processor; - @ServiceActivator(inputChannel = Processor.INPUT) + @ServiceActivator(inputChannel = Processor.INPUT, phase = "0") @Bean public MessageHandler subjectMessageHandler() { SubjectMessageHandler messageHandler = new SubjectMessageHandler(processor); diff --git a/spring-cloud-stream-rxjava/src/main/java/org/springframework/cloud/stream/annotation/rxjava/SubjectMessageHandler.java b/spring-cloud-stream-rxjava/src/main/java/org/springframework/cloud/stream/annotation/rxjava/SubjectMessageHandler.java index 2c3a90c11..8e62d767f 100644 --- a/spring-cloud-stream-rxjava/src/main/java/org/springframework/cloud/stream/annotation/rxjava/SubjectMessageHandler.java +++ b/spring-cloud-stream-rxjava/src/main/java/org/springframework/cloud/stream/annotation/rxjava/SubjectMessageHandler.java @@ -19,7 +19,7 @@ package org.springframework.cloud.stream.annotation.rxjava; import org.slf4j.Logger; import org.slf4j.LoggerFactory; -import org.springframework.beans.factory.DisposableBean; +import org.springframework.context.SmartLifecycle; import org.springframework.integration.handler.AbstractMessageProducingHandler; import org.springframework.messaging.Message; import org.springframework.messaging.support.MessageBuilder; @@ -62,46 +62,80 @@ import rx.subjects.Subject; * @author Ilayaperumal Gopinathan */ @SuppressWarnings({"unchecked", "rawtypes"}) -public class SubjectMessageHandler extends AbstractMessageProducingHandler implements DisposableBean { +public class SubjectMessageHandler extends AbstractMessageProducingHandler implements SmartLifecycle { private final Logger logger = LoggerFactory.getLogger(getClass()); @SuppressWarnings("rawtypes") private final RxJavaProcessor processor; - private final Subject subject; + private volatile Subject subject; - private final Subscription subscription; + private volatile Subscription subscription; + + private volatile boolean running = false; @SuppressWarnings({"unchecked", "rawtypes"}) public SubjectMessageHandler(RxJavaProcessor processor) { Assert.notNull(processor, "RxJava processor must not be null."); this.processor = processor; - subject = new SerializedSubject(PublishSubject.create()); - Observable outputStream = processor.process(subject); - subscription = outputStream.subscribe(new Action1() { - @Override - public void call(Object outputObject) { - if (ClassUtils.isAssignable(Message.class, outputObject.getClass())) { - getOutputChannel().send((Message) outputObject); - } - else { - getOutputChannel().send(MessageBuilder.withPayload(outputObject).build()); - } - } - }, new Action1() { - @Override - public void call(Throwable throwable) { - logger.error(throwable.getMessage(), throwable); - } - }, new Action0() { - @Override - public void call() { - logger.error("Subscription close for [" + subscription + "]"); - } - }); } + @Override + public synchronized void start() { + if (!running) { + subject = new SerializedSubject(PublishSubject.create()); + Observable outputStream = processor.process(subject); + subscription = outputStream.subscribe(new Action1() { + + @Override + public void call(Object outputObject) { + if (ClassUtils.isAssignable(Message.class, outputObject.getClass())) { + getOutputChannel().send((Message) outputObject); + } + else { + getOutputChannel().send(MessageBuilder.withPayload(outputObject).build()); + } + } + }, new Action1() { + + @Override + public void call(Throwable throwable) { + logger.error(throwable.getMessage(), throwable); + } + }, new Action0() { + + @Override + public void call() { + logger.error("Subscription close for [" + subscription + "]"); + } + }); + running = true; + } + } + + @Override + public synchronized boolean isRunning() { + return running; + } + + @Override + public boolean isAutoStartup() { + return false; + } + + @Override + public void stop(Runnable callback) { + stop(); + if (callback != null) { + callback.run(); + } + } + + @Override + public int getPhase() { + return 0; + } @Override //todo: support module input type @@ -110,7 +144,13 @@ public class SubjectMessageHandler extends AbstractMessageProducingHandler imple } @Override - public void destroy() throws Exception { - subscription.unsubscribe(); + public synchronized void stop() { + if (running) { + subject.onCompleted(); + subscription.unsubscribe(); + subscription = null; + subject = null; + running = false; + } } }