RXJS中如何对单来源数据流实现可暂停的缓冲区功能?
问题根因
- 你对
source$做了两次独立订阅,错误检测分支和输出缓存分支各自接收事件,时序无法对齐,开关状态切换和缓存触发释放逻辑不匹配 bufferToggle与windowToggle的组合逻辑错误:windowToggle返回高阶Observable,和bufferToggle返回的缓冲数组合并后直接用concatMap(from)处理,会导致事件时序错乱、缓冲事件丢失- 触发异常的当前事件没有被存入缓存,直接被错误处理分支吞掉
- 业务函数
myFunction被两个订阅分支重复执行,不符合预期
修复方案
import { Subject, BehaviorSubject, Observable, fromEventPattern, concatMap, tap, filter, take, retry } from 'rxjs'; interface IPsEvent {} declare function addEventHandler(handler: (evt: IPsEvent) => void): void; declare function myFunction(evt: IPsEvent): void; const source$ = new Subject<IPsEvent>(); // 状态开关:true=正常执行 false=异常中缓存事件 const isWorking$ = new BehaviorSubject(true); // 事件处理主逻辑 source$.pipe( concatMap((evt) => { return isWorking$.pipe( filter(working => working), take(1), tap(() => { myFunction(evt); }), retry({ delay: () => { // 执行报错就切换为异常状态 isWorking$.next(false); // 每50ms轮询直到业务逻辑恢复 return new Observable(sub => { const timer = setInterval(() => { try { myFunction(evt); // 恢复后切回正常状态 isWorking$.next(true); clearInterval(timer); sub.next(null); sub.complete(); } catch {} }, 50); }) } }) ) }) ).subscribe({ next: (evt) => console.log('处理完成', evt), error: (err) => console.error('致命错误', err) }); // 接入原始事件流 fromEventPattern<IPsEvent>(addEventHandler).subscribe(source$);
逻辑说明
- 仅对
source$做一次订阅处理,所有事件按进入顺序用concatMap串行处理,天然保证时序不乱 - 每次处理事件前先判断当前运行状态,仅正常状态下才执行业务逻辑
- 执行抛出异常时自动切换为异常状态,启动50ms间隔的轮询重试,重试成功后切回正常状态,当前事件处理完成后自动处理后续排队的事件
- 依赖RxJS原生
concatMap的入队能力做事件缓存,不需要手动维护bufferToggle相关的复杂缓存逻辑,不会出现事件丢失、顺序错乱问题
内容的提问来源于stack exchange,提问作者Dave Stein
相关产品推荐
相关产品推荐

