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

如何合并两个Observable并实现符合特定条件的完成逻辑?

实现自定义Observable完成逻辑

需求明确

现有两个Observable,每个都会先发射一个值,延迟一段时间后完成(弹珠图类似---1---|)。需要基于这两个Observable创建新的obsC$,满足以下任一条件时立即完成:

  • 条件1:其中一个Observable完成,且另一个从未发射过任何值;
  • 条件2:两个Observable都已发射过值,且都已完成。

实现方案

核心思路是跟踪两个Observable的发射状态和完成状态,在状态满足终止条件时主动结束流。以下是具体代码实现:

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

// 定义状态结构,跟踪两个Observable的发射/完成状态
interface TrackState {
  aEmitted: boolean;
  bEmitted: boolean;
  aCompleted: boolean;
  bCompleted: boolean;
}

const initialState: TrackState = {
  aEmitted: false,
  bEmitted: false,
  aCompleted: false,
  bCompleted: false
};

// 工具函数:将源Observable转换为包含发射/完成事件的流
const trackObservable = (source$: Observable<any>, key: 'a' | 'b') => 
  merge(
    // 发射事件:携带值和所属Observable的标识
    source$.pipe(map(value => ({ type: 'emitted', key, value }))),
    // 完成事件:仅携带所属Observable的标识
    source$.pipe(ignoreElements(), map(() => ({ type: 'completed', key })))
  );

// 创建目标Observable obsC$
const obsC$ = merge(
  trackObservable(obsA$, 'a'),
  trackObservable(obsB$, 'b')
).pipe(
  // 累积状态,同时保留原始发射值
  scan((acc, event) => {
    const newState = { ...acc.state };
    if (event.type === 'emitted') {
      newState[`${event.key}Emitted`] = true;
      return { state: newState, value: event.value };
    } else {
      newState[`${event.key}Completed`] = true;
      return { state: newState, value: undefined };
    }
  }, { state: initialState, value: undefined }),
  // 过滤掉完成事件,只保留原始发射值
  filter(({ value }) => value !== undefined),
  // 判断是否继续发射:满足终止条件时停止
  takeWhile(({ state }) => {
    const condition1 = (state.aCompleted && !state.bEmitted) || (state.bCompleted && !state.aEmitted);
    const condition2 = state.aCompleted && state.bCompleted && state.aEmitted && state.bEmitted;
    // 返回false时,流立即完成
    return !(condition1 || condition2);
  }, true), // 确保最后一个符合条件的值会被发射
  // 映射回原始发射值
  map(({ value }) => value)
);

代码说明

  1. 状态跟踪:通过trackObservable将每个源Observable拆分为「发射事件」和「完成事件」,确保能准确区分每个事件所属的Observable。
  2. 状态累积:使用scan操作符累积状态,记录每个Observable是否已发射值、是否已完成,同时保留原始的发射值。
  3. 终止条件判断:通过takeWhile检查当前状态,当满足条件1或条件2时,立即终止流,完成obsC$。
  4. 值传递:过滤掉完成事件,仅保留源Observable的发射值,保证obsC$的输出与merge(obsA$, obsB$)的输出一致,直到流完成。

内容的提问来源于stack exchange,提问作者Chabbey François

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.17 14:14:58