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);
这里有两个关键问题:
Flowable.create是冷流,每个订阅都会触发工厂回调:当merge订阅a1时,a会被第一次订阅,此时emitter变量被赋值为a1对应的发射器;紧接着merge订阅a2时,a会被第二次订阅,工厂回调再次执行,emitter变量被覆盖为a2的发射器。- 全局静态变量导致只有最后一个发射器生效:你后续通过定时器调用
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
相关产品推荐
相关产品推荐

