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
相关产品推荐
相关产品推荐

