RxJava observeOn调度疑问:如何让订阅者按事件发布顺序执行?
这个问题我之前也碰到过,核心原因是你用的Schedulers.single()的调度逻辑和你预期的不一样,咱们一步步拆解:
为什么会出现“先全处理Subscriber1,再处理Subscriber2”的情况?
Schedulers.single()是一个单线程调度器,它会维护一个任务队列,所有通过observeOn提交的任务都会按提交顺序排队,然后逐个执行。
你在循环里先连续调用10次publishSubject1.onNext,等于把10个“执行Subscriber1逻辑”的任务一口气塞进了队列;接着又连续调用10次publishSubject2.onNext,把10个“执行Subscriber2逻辑”的任务加到队列尾部。队列的顺序是「10个Subscriber1任务 → 10个Subscriber2任务」,自然会先跑完所有Subscriber1的任务,再执行Subscriber2的。
两种实现交替执行的方案
方案1:合并为单个Subject(最直接的方式)
把两个Subject的事件打包成带类型标识的事件对象,用同一个Subject发送。这样每次循环先发Subject1的事件,再发Subject2的事件,任务队列里就会是交替的,处理时自然交替输出:
import java.util.concurrent.TimeUnit; import io.reactivex.schedulers.Schedulers; import io.reactivex.subjects.PublishSubject; public class ObserveOnApp { // 自定义事件类,区分不同来源的事件 static class TaskEvent { String source; String data; TaskEvent(String source, String data) { this.source = source; this.data = data; } } public static void main(String[] args) { PublishSubject<TaskEvent> sharedSubject = PublishSubject.create(); sharedSubject .observeOn(Schedulers.single()) .subscribe(event -> { if ("subject1".equals(event.source)) { System.out.println("Subscriber1"); } else if ("subject2".equals(event.source)) { System.out.println("Subscriber2"); } }); for (int i = 0; i < 10; i++) { // 交替发送两个来源的事件 sharedSubject.onNext(new TaskEvent("subject1", "next")); sharedSubject.onNext(new TaskEvent("subject2", "next")); } try { TimeUnit.SECONDS.sleep(2); } catch (InterruptedException e) { e.printStackTrace(); } } }
方案2:用Zip操作符配对两个Subject的事件
如果必须保留两个独立的Subject,可以用Observable.zip把两个Subject的事件按顺序配对,每收到一对事件就连续处理两次:
import java.util.concurrent.TimeUnit; import io.reactivex.Observable; import io.reactivex.schedulers.Schedulers; import io.reactivex.subjects.PublishSubject; public class ObserveOnApp { // 简单的配对类,用来装两个Subject的事件 static class EventPair<T1, T2> { T1 subject1Event; T2 subject2Event; EventPair(T1 subject1Event, T2 subject2Event) { this.subject1Event = subject1Event; this.subject2Event = subject2Event; } } public static void main(String[] args) { PublishSubject<String> publishSubject1 = PublishSubject.create(); PublishSubject<String> publishSubject2 = PublishSubject.create(); // 配对两个Subject的事件,每收到一对就处理一次 Observable.zip(publishSubject1, publishSubject2, EventPair::new) .observeOn(Schedulers.single()) .subscribe(pair -> { System.out.println("Subscriber1"); System.out.println("Subscriber2"); }); for (int i = 0; i < 10; i++) { publishSubject1.onNext("next"); publishSubject2.onNext("next"); } try { TimeUnit.SECONDS.sleep(2); } catch (InterruptedException e) { e.printStackTrace(); } } }
两种方案都能实现你要的交替执行效果:方案1更灵活,适合后续扩展;方案2保留了原有的两个Subject结构,适合不想改动太多的场景。
内容的提问来源于stack exchange,提问作者Przemek Krysztofiak
相关产品推荐
相关产品推荐

