如何在不同时间合并多个stream且不丢失监听器实现多对多流转
动态多流合并实现方案
你需要的能力可以通过高阶Subject+mergeAll操作符结合BehaviorSubject实现,不需要复杂的自定义组件,主流Rx系库(RxJS、RxJava、RxDart等)都支持这套逻辑,以下是核心实现逻辑:
核心思路
- 声明一个专门接收流的高阶
Subject,所有后续需要加入合并队列的新流都发送到这个Subject中 - 通过
mergeAll操作符把高阶Subject中所有的流打平,合并为一个统一的输出流 - 若需要新订阅者能拿到历史最新值,把合并后的输出流喂给
BehaviorSubject即可
代码示例(RxJS为例)
import { Subject, BehaviorSubject, mergeAll } from 'rxjs'; // 高阶流:用于动态接收待合并的子流 const streamCollector = new Subject(); // 合并所有子流得到原始合并流 const rawMergedStream = streamCollector.pipe(mergeAll()); // 接入BehaviorSubject实现新订阅者获取最新值的能力 const mergedStream = new BehaviorSubject(null /* 可自定义初始值 */); rawMergedStream.subscribe(mergedStream); // 暴露添加新流的方法 function addStreamToMerge(newStream) { streamCollector.next(newStream); }
场景验证
完全匹配你给出的示例逻辑:
- 调用
addStreamToMerge(streamA)、addStreamToMerge(streamB)将A、B加入合并队列 - E、F订阅
mergedStream,可以收到A、B发送的所有数据 - 后续G订阅
mergedStream,可以先拿到最近一次推送的数据,再接收后续新数据 - 调用
addStreamToMerge(streamC)将C加入合并队列,已订阅的E、F、G都能同步收到C发送的所有数据
扩展能力
如果需要支持动态移除已加入的子流,只需要在添加子流时搭配takeUntil操作符控制子流的销毁时机即可,不会影响其他已合并的流和已有订阅者。
内容的提问来源于stack exchange,提问作者Guy Luz
相关产品推荐
相关产品推荐

