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

RxJava中Flowable.merge合并create/interval生成流结果不同的原因

问题根源分析:为什么Flowable.merge对自定义创建的流和interval流表现不同?

核心问题出在你自定义的Flowable a的订阅行为和Flowable.interval的冷流特性差异上,咱们一步步拆解:

一、为什么a1没有发射事件?

看你定义Flowable a的代码:

private static FlowableEmitter<Integer> emitter;
Flowable<Integer> a = Flowable.create(emitter -> MergeProblem.emitter = emitter, BackpressureStrategy.BUFFER);

这里有两个关键问题:

  1. Flowable.create是冷流,每个订阅都会触发工厂回调:当merge订阅a1时,a会被第一次订阅,此时emitter变量被赋值为a1对应的发射器;紧接着merge订阅a2时,a会被第二次订阅,工厂回调再次执行,emitter变量被覆盖为a2的发射器。
  2. 全局静态变量导致只有最后一个发射器生效:你后续通过定时器调用emitter.onNext()时,这个emitter已经是a2的了,a1的发射器早就被覆盖,根本接收不到任何事件——这就是为什么你只看到a2的输出。

简单说:你以为a1和a2共享同一个流的事件,但实际上它们各自触发了a的订阅,而你用全局变量把发射器替换成了最后一个,自然只有a2能收到事件。

二、为什么b1和b2能正常工作?

再看Flowable.interval的特性:

Flowable<Long> b = Flowable.interval(1, TimeUnit.SECONDS);
Flowable<String> b1 = b.map(x -> "b1 " + x);
Flowable<String> b2 = b.map(x -> "b2 " + x);

interval是冷流,每个订阅都会启动一个独立的定时器。当merge订阅b1和b2时:

  • b1订阅b,启动第一个定时器,每秒发射递增的Long值;
  • b2订阅b,启动第二个独立的定时器,同样每秒发射递增的Long值;
    两个定时器互不干扰,各自向b1和b2发射事件,所以merge能同时收到两者的输出。

解决方案:让自定义流支持多订阅者

如果你想让a1和a2都能收到同一个流的事件,需要把a转换成热流,让所有订阅者共享同一个发射器。可以用publish()或share()操作符:

// 将冷流转成热流,所有订阅者共享同一个发射源
Flowable<Integer> a = Flowable.create(emitter -> MergeProblem.emitter = emitter, BackpressureStrategy.BUFFER)
    .publish()
    .autoConnect(); // autoConnect()会在第一个订阅者订阅时触发流的创建

这样a1和a2订阅的是同一个热流,emitter只会被赋值一次,后续调用onNext()时,a1和a2都会收到事件。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.15 07:44:37