温馨提示×

温馨提示×

您好,登录后才能下订单哦!

密码登录×
登录注册×
其他方式登录
点击 登录注册 即表示同意《亿速云用户服务条款》

rxjava怎么处理并发

发布时间:2026-09-20 16:09:24 来源:亿速云 阅读:84 作者:小樊 栏目:编程语言

RxJava 中,“处理并发”的核心就是:用 Scheduler 控制任务在哪个线程执行,用操作符组合/限制并发度。下面按常见场景给你梳理。


一、RxJava 并发的基础:Scheduler

RxJava 默认是同步、当前线程执行的,想并发必须切换线程。

常用 Scheduler

Schedulers.io()          // IO 密集型(网络、数据库)
Schedulers.computation() // CPU 密集型计算
Schedulers.newThread()   // 新线程(不推荐频繁用)
AndroidSchedulers.mainThread() // Android 主线程
Schedulers.single()      // 单一后台线程
Schedulers.trampoline()  // 当前线程排队执行

二、最基础的线程切换

1️⃣ subscribeOn:指定上游执行线程

Observable.create(emitter -> {
    // 在 IO 线程执行
})
.subscribeOn(Schedulers.io())
.subscribe();

2️⃣ observeOn:指定下游执行线程

Observable.just(1)
    .subscribeOn(Schedulers.io())
    .observeOn(AndroidSchedulers.mainThread())
    .subscribe(i -> {
        // 主线程更新 UI
    });

规则

  • subscribeOn 只有第一个有效(上游)
  • observeOn 可以多次切换

三、并发执行多个任务(并行)

方式一:flatMap + subscribeOn(最常用)

Observable.fromIterable(ids)
    .flatMap(id ->
        Observable.just(id)
            .subscribeOn(Schedulers.io())
            .map(this::request)
    )
    .observeOn(AndroidSchedulers.mainThread())
    .subscribe(result -> {});

✅ 每个 id 在不同 IO 线程并发执行
⚠️ 并发度不可控(可能很多线程)


四、控制并发数量(非常重要)

方式二:flatMap(maxConcurrency)

Observable.fromIterable(ids)
    .flatMap(id ->
        Observable.just(id)
            .subscribeOn(Schedulers.io())
            .map(this::request),
        5 // 最多 5 个并发
    )
    .subscribe();

✅ 推荐做法
✅ 防止线程爆炸


五、等所有任务完成(合并结果)

1️⃣ zip(一对一合并)

Observable.zip(
    task1(), task2(), task3(),
    (a, b, c) -> a + b + c
).subscribe();

2️⃣ merge(合并流,不保证顺序)

Observable.merge(task1(), task2())
    .subscribe();

3️⃣ concat(顺序执行)

Observable.concat(task1(), task2())
    .subscribe();

六、RxJava 2/3 中的并发陷阱

❌ 错误:在 map 里做耗时操作

.map(i -> {
    Thread.sleep(1000); // ❌ 阻塞当前线程
    return i;
})

✅ 正确:用 flatMap + subscribeOn


七、Android 中常见并发模型

api.getUser()
    .subscribeOn(Schedulers.io())
    .flatMap(user -> api.getOrders(user.id)
        .subscribeOn(Schedulers.io()))
    .observeOn(AndroidSchedulers.mainThread())
    .subscribe(orders -> {
        // UI
    });

八、什么时候用 RxJava 并发?

✅ 适合:

  • 网络请求并发
  • 数据流变换
  • 多数据源合并

❌ 不适合:

  • 简单异步(用 Coroutine / CompletableFuture 更简洁)

如果你愿意,可以告诉我:

  • 用的是 RxJava 2 还是 3
  • Android 还是后端
  • 具体场景(如:100 个请求并发?限流?)

我可以直接给你可运行示例代码

向AI问一下细节

免责声明:本站发布的内容(图片、视频和文字)以原创、转载和分享为主,文章观点不代表本网站立场,如果涉及侵权请联系站长邮箱:is@yisu.com进行举报,并提供相关证据,一经查实,将立刻删除涉嫌侵权内容。

AI