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

RXJS合并多Observable时如何实现互斥执行(忽略非活跃流触发)

解决方案

要实现任意一个Observable处于活跃状态时,所有其他触发请求都被忽略的效果,核心是把所有按钮的触发事件合并后,用**全局的exhaustMap**统一管控并发,而非给每个Subject单独设置exhaustMap。

基础实现代码

import { merge, timer, exhaustMap, map, Subject } from 'rxjs';

// 假设你的三个Subject定义如下
const one$ = new Subject();
const two$ = new Subject();
const three$ = new Subject();

// 1. 合并所有按钮的触发事件,标记来源(可选,方便区分不同按钮的逻辑)
const allButtonTriggers$ = merge(
  one$.pipe(map(() => 'btn1')),
  two$.pipe(map(() => 'btn2')),
  three$.pipe(map(() => 'btn3'))
);

// 2. 用全局exhaustMap处理,只要有一个内部Observable活跃,新触发就会被忽略
const globalProcess$ = allButtonTriggers$.pipe(
  exhaustMap((triggerSource) => {
    console.log(`开始处理${triggerSource}的请求`);
    return timer(0, 1000); // 示例统一使用timer逻辑
  })
);

// 订阅结果
globalProcess$.subscribe(console.log);

原理说明

  • 原写法中,每个Subject单独绑定exhaustMap,只能忽略同一个按钮的重复触发,不同按钮的触发会各自启动timer,互不干扰。
  • 现在先通过merge把所有按钮的触发合并为一个Observable,再套一层全局exhaustMap:只要当前有一个内部Observable(比如timer)在执行,所有新触发(无论来自哪个按钮)都会被直接忽略,直到当前活跃的Observable完成。

扩展:不同按钮对应不同处理逻辑

如果每个按钮需要执行差异化业务逻辑,可在exhaustMap中根据标记的来源分支处理:

import { merge, timer, interval, of, exhaustMap, map, delay, take, EMPTY, Subject } from 'rxjs';

const globalProcess$ = allButtonTriggers$.pipe(
  exhaustMap((triggerSource) => {
    switch(triggerSource) {
      case 'btn1':
        return timer(0, 1000).pipe(take(5)); // 每秒触发,共执行5次
      case 'btn2':
        return interval(500).pipe(take(10)); // 每500ms触发,共执行10次
      case 'btn3':
        return of('btn3处理完成').pipe(delay(3000)); // 延迟3秒返回结果
      default:
        return EMPTY;
    }
  })
);

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.09 10:36:31