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

RxJS 7中无法捕获BehaviorSubject等Subject抛出的错误问题

问题原因

RxJS的Subject一旦调用error()方法,就会进入终止状态,后续所有next()、error()调用都不会再向订阅者发射任何数据。你的代码中,3秒时调用source.error()已经让这个Subject彻底终止,6秒时再调用source.error()不会触发任何订阅逻辑,所以控制台看不到这个错误。

同时,retry操作符的逻辑是当源Observable出错时重新订阅源,但你的源是已经终止的Subject,重新订阅也不会有任何发射,重试机制自然无法继续工作。

解决方案

要实现多次触发错误并被捕获的效果,需要避免让源Observable提前终止,以下是两种可行的修改方式:

方式1:每次重试时创建新的Subject

用defer包装Subject的创建逻辑,让每次重试都生成新的Subject,确保流能重新激活:

import {
  timer,
  Subject,
  defer
} from 'rxjs';
import { retry, tap } from 'rxjs/operators';

const example = defer(() => {
  const source = new Subject();
  
  // 发射初始数据
  source.next({ status: 'open' });
  
  // 3秒触发第一次错误
  setTimeout(() => {
    source.error({ message: 'Error here !!!' });
  }, 3000);
  
  // 6秒触发第二次错误
  setTimeout(() => {
    console.log('Next err =================>');
    source.error({ message: 'Next err =================>' });
  }, 6000);
  
  return source.pipe(
    tap((ob) => console.log('OBJ: ', ob))
  );
}).pipe(
  retry({
    delay: (errors) => {
      console.log('Last error : ', errors);
      return timer(3000);
    },
  })
);

const subscribe = example.subscribe();

方式2:用Subject发射错误信号而非直接调用error()

不直接在源Subject上调用error(),而是发射错误对象,再通过mergeMap转换为错误流,这样源Subject不会终止,能持续接收后续信号:

import {
  timer,
  Subject,
  throwError,
  of
} from 'rxjs';
import { mergeMap, retry, tap } from 'rxjs/operators';

const source = new Subject();

const example = source.pipe(
  tap((ob) => console.log('OBJ: ', ob)),
  mergeMap(data => {
    // 判断是否为错误对象,转换为错误流
    if (data.message) {
      return throwError(() => data);
    }
    return of(data);
  }),
  retry({
    delay: (errors) => {
      console.log('Last error : ', errors);
      return timer(3000);
    },
  })
);

// 发射初始正常数据
source.next({ status: 'open' });

// 3秒发射第一个错误对象
setTimeout(() => {
  source.next({ message: 'Error here !!!' });
}, 3000);

// 6秒发射第二个错误对象
setTimeout(() => {
  console.log('Next err =================>');
  source.next({ message: 'Next err =================>' });
}, 6000);

const subscribe = example.subscribe();
关键知识点
  • RxJS的Observable(包括Subject)是一次性的:一旦调用complete()或error(),流就永久终止,无法复用。
  • retry操作符的作用是重新订阅源Observable,如果源已经终止,重试不会产生任何效果。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.17 01:25:27