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

RxJS concat中Observable未执行:任务并发不符合串行预期

RxJS 任务串行与并发问题修复

问题核心

你的场景要求:

  • startTask必须串行执行,同一时间只能有一个实例运行
  • fetchScan和fetchTitle都依赖startTask
  • 不同标题的fetchTitle任务可并发,但各自依赖的startTask仍需串行

原代码用concat构建队列后出现两个问题:fetchTitle任务并发(导致startTask并行)、fetchScan未执行,本质是没有把startTask的调用统一到串行控制流中,而是直接控制了上层任务的执行顺序。

修复方案

1. 构建startTask串行执行队列

创建一个单例的请求流,所有startTask调用都通过这个流处理,用concatMap保证串行执行:

import { Subject, concatMap, of, merge } from 'rxjs';

// 模拟原异步函数
function startTask(taskId: string) {
  return new Promise(resolve => {
    console.log(`startTask ${taskId} 启动`);
    setTimeout(() => {
      console.log(`startTask ${taskId} 完成`);
      resolve(`startTask-${taskId}-result`);
    }, 1000);
  });
}

// 创建startTask请求主题,用于接收所有调用请求
const startTaskQueue$ = new Subject<{
  taskId: string;
  resolve: (value: string) => void;
}>();

// 用concatMap保证startTask串行执行
startTaskQueue$.pipe(
  concatMap((req) => 
    of(null).pipe(concatMap(() => startTask(req.taskId)))
  )
).subscribe(result => {
  // 完成后返回结果给对应请求
  const currentReq = startTaskQueue$.value;
  currentReq.resolve(result as string);
});

// 封装串行版startTask
const serialStartTask = (taskId: string): Promise<string> => {
  return new Promise(resolve => {
    startTaskQueue$.next({ taskId, resolve });
  });
};

2. 改造依赖函数

让fetchScan和fetchTitle使用串行版的startTask:

function fetchScan() {
  return serialStartTask('scan').then(res => {
    console.log('fetchScan 处理完成');
    return res;
  });
}

function fetchTitle(title: string) {
  return serialStartTask(`title-${title}`).then(res => {
    console.log(`fetchTitle ${title} 处理完成`);
    return `${res}-${title}`;
  });
}

3. 执行任务流

用merge允许fetchTitle并发,同时保证所有startTask串行:

// 执行测试:先触发fetchScan,再并发触发两个fetchTitle
merge(
  fetchScan(),
  fetchTitle('标题A'),
  fetchTitle('标题B')
).subscribe();

预期输出

startTask scan 启动
startTask scan 完成
fetchScan 处理完成
startTask title-标题A 启动
startTask title-标题A 完成
fetchTitle 标题A 处理完成
startTask title-标题B 启动
startTask title-标题B 完成
fetchTitle 标题B 处理完成

原问题分析

  • 原代码错误地将fetchTitle和fetchScan整体放入concat队列,要么限制了fetchTitle的并发,要么没有统一控制startTask的调用,导致多个startTask同时执行
  • 未将fetchScan正确加入执行流,导致其未被触发

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.04 00:00:25