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

