You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.14 07:08:44