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

如何在RxJS中实现类似zip但逻辑不同的Subject组合?

自定义RxJS组合操作符实现需求

你要的这种行为,核心是每一轮需要两个流都至少各更新一次,才触发一次最新值的组合——和zip的严格顺序配对、combineLatest的任意更新触发逻辑都不同。可以通过RxJS基础操作符组合实现:

实现思路

  1. 跟踪两个流的「待更新」状态:初始时需要两个流各发射一次;
  2. 任一流发射时更新对应状态,只有当两个流的待更新状态都变为「已完成」时,输出当前最新值,随后重置状态进入下一轮。

代码实现

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.05 06:43:21