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

使用RXJS如何实现双源Observable发射计数相等时触发的流?

双Observable同步触发实现方案

以下基于RxJS实现两种需求逻辑,可直接根据场景选择使用:

逻辑(a):双方已发射元素数量相等时触发

实现思路:合并两个源的发射事件,分别维护两个源的发射值缓存队列,每次新事件到达后判断两个队列长度是否相等,相等时触发目标Observable发射最新的一对值。

import { merge, Observable } from 'rxjs';
import { scan, filter, map } from 'rxjs/operators';

type SourceValue<T1, T2> = { source: 'obs1', value: T1 } | { source: 'obs2', value: T2 };

function buildEqualCountObservable<T1, T2>(obs1: Observable<T1>, obs2: Observable<T2>): Observable<[T1, T2]> {
  return merge(
    obs1.pipe(map(val => ({ source: 'obs1', value: val }))),
    obs2.pipe(map(val => ({ source: 'obs2', value: val })))
  ).pipe(
    scan((state, curr) => {
      curr.source === 'obs1' ? state.queue1.push(curr.value) : state.queue2.push(curr.value);
      return state;
    }, { queue1: [] as T1[], queue2: [] as T2[] }),
    filter(state => state.queue1.length === state.queue2.length),
    map(state => [state.queue1[state.queue1.length - 1], state.queue2[state.queue2.length - 1]] as [T1, T2])
  );
}

运行示例符合需求:obs1先发射第1个元素时无触发,obs2发射第1个元素时两队列长度均为1,触发第一次发射;obs2连续发射第2、3个元素时队列长度不一致不触发,直到obs1发射第2个元素时两队列长度均为2,触发第二次发射。

逻辑(b):当前发射数最少的源产生新发射时触发

实现思路:维护两个源的发射计数,每次新事件到达后判断触发事件的源是否为本次发射前计数最小的源,是则触发目标Observable发射。

function buildMinCountIncreaseObservable<T1, T2>(obs1: Observable<T1>, obs2: Observable<T2>): Observable<{ latestObs1: T1, latestObs2: T2, currentMinCount: number }> {
  return merge(
    obs1.pipe(map(val => ({ source: 'obs1', value: val }))),
    obs2.pipe(map(val => ({ source: 'obs2', value: val })))
  ).pipe(
    scan((state, curr) => {
      // 计算发射前的最小计数
      const prevMin = Math.min(state.count1, state.count2);
      // 更新计数和最新值
      if (curr.source === 'obs1') {
        state.count1++;
        state.latestObs1 = curr.value;
      } else {
        state.count2++;
        state.latestObs2 = curr.value;
      }
      // 判断本次发射是否来自之前的最小计数源
      state.needTrigger = (curr.source === 'obs1' ? state.count1 - 1 : state.count2 - 1) === prevMin;
      return state;
    }, { count1: 0, count2: 0, latestObs1: undefined as T1 | undefined, latestObs2: undefined as T2 | undefined, needTrigger: false }),
    filter(state => state.needTrigger),
    map(state => ({
      latestObs1: state.latestObs1!,
      latestObs2: state.latestObs2!,
      currentMinCount: Math.min(state.count1, state.count2)
    }))
  );
}

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.03 20:39:01