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

自定义RxJS操作符仅首次触发问题排查

问题

我尝试创建自定义RxJS操作符rxThrowCustomServerError,用来检测API返回的BaseServerResponse是否包含error字段,存在则抛出对应错误,并在ticket-feedback.service.ts中处理该错误。但首次点击表单提交按钮时功能正常(获取API数据、抛出错误、catchError处理),后续点击无任何反应(甚至不发起API请求),移除该自定义操作符后每次点击均正常。请问问题出在哪里?


相关代码片段

BaseServerResponse类型定义

class BaseServerResponse<T> {
    data?: T;
    error?: string;
    message?: string;
}

自定义RxJS操作符代码

export const rxThrowCustomServerError = <T>() => {
  return function (source: Observable<BaseServerResponse<T>>): Observable<BaseServerResponse<T>> {
    return new Observable((subscriber) => {
      const subscription = source.subscribe({
        next(x) {
          if (x.error) {
            subscriber.error(x.error);
          }
          subscriber.next(x);
        },
        error(error) {
          subscriber.error(error);
        },
        complete() {
          subscriber.complete();
        },
      });

      return () => {
        subscription.unsubscribe();
      };
    });
  };
};

ticket-feedback.component.ts代码

@Component({
  selector: 'ticket-feedback-list',
  standalone: true,
  imports: [...ANGULAR_MODULES, ...ANGULAR_MATERIAL_MODULES],
  providers: [],
  templateUrl: './ticket-feedback-list.component.html',
  styleUrls: ['./ticket-feedback-list.component.scss'],
  changeDetection: ChangeDetectionStrategy.OnPush,
})
export default class TicketFeedbackListComponent {
  private _service = inject(TicketFeedbackService);
  form = this._service.getFeedbackForm;
  feedbackResponse = toSignal(this._service.feedbackResponse$);
  public forSubmit() {
    this._service.feedbackRequested$.next({
      fromDate: this.form.controls.fromDate.value,
      toDate: this.form.controls.toDate.value,
      retailerMobileNumber: this.form.controls.retailerMobileNumber.value,
      pageNumber: 1,
      pageSize: 10,
    });
  }
}

ticket-feedback.service.ts代码

feedbackResponse$ = this.feedbackRequested$.pipe(
    filter((feedbackRequest) => feedbackRequest && !isEmpty(feedbackRequest) && this.getFeedbackForm.valid),
    switchMap((request) =>
      this._http.get<BaseServerResponse<Feedback[]>>(appApiResources.feedback, {
        params: {
          pageNumber: request.pageNumber,
          pageSize: request.pageSize,
          fromDate: request.fromDate,
          toDate: request.toDate,
          'api-version': 1,
        },
      })
    ),
    rxThrowCustomServerError(),
    catchError(() => of(new BaseServerResponse<Feedback[]>([] as Feedback[], 0, 0, 0, '', '', 'Unable to fetch feedbacks'))),
    finalize(() => this.feedbackRequested$.next({} as IGetFeedbackRequest))
  );

HTML模板代码

feedbackResponse() | json

问题原因

  1. 自定义操作符违反RxJS契约:在检测到x.error存在时,先调用subscriber.error(x.error)触发错误流,随后又调用subscriber.next(x)。根据RxJS规范,Observable一旦触发error或complete就会进入终止状态,后续的next调用会被忽略,这种违规操作会导致流状态异常,影响后续事件处理。

  2. catchError中实例化错误:catchError里的BaseServerResponse构造调用传入了7个参数,但类型定义仅包含data、error、message三个可选字段,这会导致额外的未捕获错误,最终使整个feedbackResponse$流彻底终止,后续feedbackRequested$发出的事件无法被处理。


解决方案

1. 修正自定义操作符

将操作符改为符合RxJS规范的实现,确保检测到错误时仅抛出错误,不发送next值;或者使用更简洁的map+throwError方式实现:

方式一:修正原Observable实现

export const rxThrowCustomServerError = <T>() => {
  return function (source: Observable<BaseServerResponse<T>>): Observable<BaseServerResponse<T>> {
    return new Observable((subscriber) => {
      const subscription = source.subscribe({
        next(x) {
          if (x.error) {
            subscriber.error(x.error);
          } else {
            // 仅在无错误时发送数据
            subscriber.next(x);
          }
        },
        error(error) {
          subscriber.error(error);
        },
        complete() {
          subscriber.complete();
        },
      });

      return () => {
        subscription.unsubscribe();
      };
    });
  };
};

方式二:使用map+throwError简化实现

import { map, throwError } from 'rxjs';

export const rxThrowCustomServerError = <T>() => 
  map((response: BaseServerResponse<T>) => {
    if (response.error) {
      throw throwError(() => new Error(response.error));
    }
    return response;
  });

2. 修正catchError中的实例化代码

匹配BaseServerResponse的类型定义,正确实例化对象:

catchError(() => of(new BaseServerResponse<Feedback[]>({
  data: [] as Feedback[],
  message: 'Unable to fetch feedbacks'
})));

如果BaseServerResponse有对应构造函数,也可以改为:

// 假设构造函数为 constructor(data?: T, error?: string, message?: string)
catchError(() => of(new BaseServerResponse<Feedback[]>([] as Feedback[], '', 'Unable to fetch feedbacks')));

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.08 22:14:59