在 RxJava 中,“处理并发”的核心就是:用 Scheduler 控制任务在哪个线程执行,用操作符组合/限制并发度。下面按常见场景给你梳理。
RxJava 默认是同步、当前线程执行的,想并发必须切换线程。
Schedulers.io() // IO 密集型(网络、数据库)
Schedulers.computation() // CPU 密集型计算
Schedulers.newThread() // 新线程(不推荐频繁用)
AndroidSchedulers.mainThread() // Android 主线程
Schedulers.single() // 单一后台线程
Schedulers.trampoline() // 当前线程排队执行
Observable.create(emitter -> {
// 在 IO 线程执行
})
.subscribeOn(Schedulers.io())
.subscribe();
Observable.just(1)
.subscribeOn(Schedulers.io())
.observeOn(AndroidSchedulers.mainThread())
.subscribe(i -> {
// 主线程更新 UI
});
✅ 规则
subscribeOn 只有第一个有效(上游)observeOn 可以多次切换Observable.fromIterable(ids)
.flatMap(id ->
Observable.just(id)
.subscribeOn(Schedulers.io())
.map(this::request)
)
.observeOn(AndroidSchedulers.mainThread())
.subscribe(result -> {});
✅ 每个 id 在不同 IO 线程并发执行
⚠️ 并发度不可控(可能很多线程)
Observable.fromIterable(ids)
.flatMap(id ->
Observable.just(id)
.subscribeOn(Schedulers.io())
.map(this::request),
5 // 最多 5 个并发
)
.subscribe();
✅ 推荐做法
✅ 防止线程爆炸
Observable.zip(
task1(), task2(), task3(),
(a, b, c) -> a + b + c
).subscribe();
Observable.merge(task1(), task2())
.subscribe();
Observable.concat(task1(), task2())
.subscribe();
.map(i -> {
Thread.sleep(1000); // ❌ 阻塞当前线程
return i;
})
✅ 正确:用 flatMap + subscribeOn
api.getUser()
.subscribeOn(Schedulers.io())
.flatMap(user -> api.getOrders(user.id)
.subscribeOn(Schedulers.io()))
.observeOn(AndroidSchedulers.mainThread())
.subscribe(orders -> {
// UI
});
✅ 适合:
❌ 不适合:
Coroutine / CompletableFuture 更简洁)如果你愿意,可以告诉我:
我可以直接给你可运行示例代码。
免责声明:本站发布的内容(图片、视频和文字)以原创、转载和分享为主,文章观点不代表本网站立场,如果涉及侵权请联系站长邮箱:is@yisu.com进行举报,并提供相关证据,一经查实,将立刻删除涉嫌侵权内容。