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

如何实现队列长度为1的类RxJS concatMap()运算符?

实现队列长度为1的类concatMap运算符

你需要的逻辑很明确:处理当前Observable时,新输入的Observable不会中断当前任务,也不会全部排队,只保留最新的那个待处理项,等当前任务结束后立刻处理这个最新项。下面是具体的RxJS实现方案:

自定义运算符实现

通过组合基础操作符维护状态,实现目标逻辑:

import { Observable, EMPTY, scan, concatMap, map, concatAll } from 'rxjs';

function concatMapLatestPending<T, R>(
  project: (value: T) => Observable<R>
): (source: Observable<T>) => Observable<R> {
  return (source) =>
    source.pipe(
      // 用scan维护状态:是否在处理、待处理的最新Observable
      scan(
        (state, value) => {
          const nextObs = project(value);
          if (state.isProcessing) {
            // 当前忙,替换待处理项为最新的输入
            return { ...state, pending: nextObs };
          }
          // 当前空闲,直接启动任务
          return { isProcessing: true, pending: null, current: nextObs };
        },
        { isProcessing: false, pending: null, current: null } as {
          isProcessing: boolean;
          pending: Observable<R> | null;
          current: Observable<R> | null;
        }
      ),
      // 处理当前任务,完成后带出待处理项
      concatMap((state) => {
        if (!state.current) return EMPTY;
        return state.current.pipe(
          map((res) => ({ res, pending: state.pending }))
        );
      }),
      // 当前任务完成后,若有pending则继续处理
      concatMap(({ res, pending }) => {
        const output = [res];
        if (pending) output.push(pending);
        return output;
      }),
      concatAll()
    );
}

核心逻辑拆解

  • 新输入到来时:
    • 若当前无任务在处理,直接启动该任务并标记「处理中」
    • 若当前正在处理,仅保留最新的输入作为待处理项(覆盖之前的待处理)
  • 当前任务完成后:
    • 若有暂存的待处理任务,立即启动它
    • 若无待处理任务,回到空闲状态

这个实现完全匹配你的需求:既不会像switchMap那样打断当前任务,也不会像concatMap那样堆积所有待处理任务,更不会像exhaustMap那样忽略后续任务,只保留最新的待处理项。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.14 15:20:22