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

RxJava2中mergeWith合并Flowable无输出问题求助

RxJava2 mergeWith 后第二个 Flowable 无输出的原因分析与解决办法

嘿,这个问题我之前也碰到过类似的情况,核心原因其实是线程被阻塞导致第二个 Flowable 的上游逻辑根本没机会执行,咱们一步步拆解来看:

问题根源

先对比你的两段代码的核心差异:

  • 在 Snippet1 中,你给两个 Flowable 分别调用了 subscribeOn(Schedulers.computation()),这意味着它们的上游创建逻辑(也就是 Flowable.create 里的 while 循环)会分别运行在 Computation 线程池的独立线程中,两个阻塞的无限循环互不干扰,所以都能正常输出数据。
  • 但在 Snippet2 中,你只在合并后的流上调用了一次 subscribeOn(Schedulers.computation()),RxJava 会把这个线程调度作用于所有被合并的 Flowable。按照 merge 的执行顺序,它会先订阅第一个 Flowable(createQ2Flowable()),而这个 Flowable 的 create 里是一个无限阻塞循环(只要 running() 返回 true,就会一直卡在 sp.in("rxLoggingKey") 这里),直接占满了这唯一的线程,导致第二个 Flowable(createMetricsFlowable())的订阅逻辑完全没机会启动,自然看不到它的输出。

再看你的 createMetricsFlowable() 和 createQ2Flowable() 实现:里面的 sp.in() 是阻塞式方法,加上外层的 while 循环,本质上是一个会一直占用线程的同步逻辑——如果不给它单独分配线程,就会阻塞后续所有依赖该线程的操作。

解决办法

最直接的修复方式是给每个被合并的 Flowable 单独指定 subscribeOn,让它们的阻塞逻辑运行在独立线程中:

createQ2Flowable()
    .subscribeOn(Schedulers.computation())
    .mergeWith(createMetricsFlowable().subscribeOn(Schedulers.computation()))
    .subscribe(onNext -> System.out.println(onNext));

另外补充一点:如果你后续需要处理背压问题,当前用的 BackpressureStrategy.BUFFER 要注意内存溢出风险,但这次的问题和背压无关,纯粹是线程调度导致的阻塞。

内容的提问来源于stack exchange,提问作者chhil

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.15 07:51:43