如何在RxJS中实现类似zip但逻辑不同的Subject组合?
自定义RxJS组合操作符实现需求
你要的这种行为,核心是每一轮需要两个流都至少各更新一次,才触发一次最新值的组合——和zip的严格顺序配对、combineLatest的任意更新触发逻辑都不同。可以通过RxJS基础操作符组合实现:
实现思路
- 跟踪两个流的「待更新」状态:初始时需要两个流各发射一次;
- 任一流发射时更新对应状态,只有当两个流的待更新状态都变为「已完成」时,输出当前最新值,随后重置状态进入下一轮。
代码实现
import { Subject, combineLatest, scan, filter, map } from 'rxjs'; function pairedLatest<T, U>(source1: Subject<T>, source2: Subject<U>) { return combineLatest([source1, source2]).pipe( scan((state, [val1, val2]) => { // 更新两个流的最新值 const newState = { latest1: val1, latest2: val2, needs1: state.needs1, needs2: state.needs2 }; // 标记当前轮次的待更新状态 if (state.needs1) newState.needs1 = false; if (state.needs2) newState.needs2 = false; // 当两个流都完成当前轮次更新时,标记可发射并重置状态 if (!newState.needs1 && !newState.needs2) { newState.canEmit = true; newState.needs1 = true; newState.needs2 = true; } else { newState.canEmit = false; } return newState; }, { latest1: null as T, latest2: null as U, needs1: true, needs2: true, canEmit: false }), // 仅输出标记为可发射的状态 filter(state => state.canEmit), // 映射为最终的组合值格式 map(state => `${state.latest1}${state.latest2}`) ); } // 测试示例 const sub1 = new Subject<number>(); const sub2 = new Subject<string>(); pairedLatest(sub1, sub2).subscribe(console.log); // 模拟你给出的发射序列 sub1.next(1); setTimeout(() => sub2.next('a'), 1000); // 输出1a setTimeout(() => sub1.next(2), 2000); setTimeout(() => sub1.next(3), 3000); setTimeout(() => sub1.next(4), 4000); setTimeout(() => sub2.next('b'), 5000); // 输出4b setTimeout(() => sub2.next('c'), 6000); setTimeout(() => sub2.next('d'), 7000); setTimeout(() => sub1.next(5), 8000); // 输出5d // 测试你的简单示例:Subject1先发射5次,Subject2再发射1次 // setTimeout(() => { // sub1.next(1); // sub1.next(2); // sub1.next(3); // sub1.next(4); // sub1.next(5); // setTimeout(() => sub2.next('x'), 1000); // 输出5x // }, 2000);
逻辑说明
combineLatest始终同步两个流的最新值,确保能拿到当前最准确的配对数据;scan维护状态机:跟踪两个流是否还需要当前轮次的更新,以及是否可以发射组合值;- 完全匹配你的需求:不管单个流中间发射多少次,只有当两个流都至少完成一次更新后,才会输出最新的组合值。
内容的提问来源于stack exchange,提问作者lukasniessen
相关产品推荐
相关产品推荐

