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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.04 13:36:04