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

如何在Angular的NgRx Effect中按顺序发送请求

问题描述

我有一个NgRx Effect sendData$,它会根据Store中的状态生成一系列添加/移除产品的Action(比如addCar/removeCar、addBike/removeBike等)。每个Action对应独立的Effect和Service,这些Effect会发起HTTP请求。

当前的问题是,所有生成的Action会被立即分发,导致对应的HTTP请求同时发送。我尝试过用concatMap但没效果,目前只能通过给每个Action加500ms延迟来临时解决,但这个方案不优雅,希望能实现请求按顺序串行执行(上一个请求完成或失败后,再发送下一个)。

原sendData$代码:

public sendData$ = createEffect(() =>
  this.actions$.pipe(
    ofType(actions.sendData),
    concatLatestFrom(() => [
      this.store.select(selectBike),
      this.store.select(selectCar),
    ]),
    concatMap(
      ([
        ,
        selectedCar,
        selectedBike,
      ]: [
        TypedAction<string>,
        boolean, 
        boolean
      ]) => {
        const actionsArray: Action[] = [];

        if (selectedCar) {
          actionsArray.push(actions.addCar());
        } else {
          actionsArray.push(actions.removeCar());
        }

        if (selectedBike) {
          actionsArray.push(actions.addBike());
        } else {
          actionsArray.push(actions.removeBike());
        }

        // ... 更多产品的判断逻辑

        return actionsArray;
      }
    )
  )
)

汽车相关的Effect和Service参考实现:
CarsEffects

@Injectable()
export class CarsEffects {
  public addCar$ = createEffect(() =>
    this.actions$.pipe(
      ofType(addCar),
      concatLatestFrom(() => [this.store.select(selectCar)]),
      concatMap(([, car]: [TypedAction<string>, Car]) =>
        this.carService.selectCar(car).pipe(
          map(() => addCarSuccess()),
          catchError(() => [addCarFailure()])
        )
      )
    )
  );

  public removeCar$ = createEffect(() =>
    this.actions$.pipe(
      ofType(removeCar),
      concatMap(() =>
        this.carService.removeCar().pipe(
          map(() => removeCarSuccess()),
          catchError(() => [removeCarFailure()])
        )
      )
    )
  );

  constructor(
    private actions$: Actions,
    private store: Store,
    private carService: CarService
  ) {}
}

CarsService

@Injectable({
  providedIn: 'root',
})
export class CarsService {
  constructor(private http: HttpClient) {}

  public selectCar(car: Car): Observable<any> {
    const selectCarUrl = `${getApiUrl()}car/select`;

    return this.http.post<ApiResponse>(selectCarUrl, {
      car_id: car.id,
    });
  }

  public removeCar(): Observable<any> {
    const removeCarUrl = `${getApiUrl()}car/remove`;

    return this.http.delete<ApiResponse>(removeCarUrl, {});
  }
}
解决方案

核心问题是:原来的代码把Action数组直接返回,NgRx会同步分发所有Action,导致对应的Effect并行发起请求。要实现串行执行,需要让每个Action的分发等待前一个Action对应的请求完成(成功或失败)。

具体做法是:

  1. 将Action数组转换为Observable流(用from())
  2. 对每个Action,先分发它,然后等待对应的成功/失败Action触发后,再分发下一个Action
修改后的代码
public sendData$ = createEffect(() =>
  this.actions$.pipe(
    ofType(actions.sendData),
    concatLatestFrom(() => [
      this.store.select(selectBike),
      this.store.select(selectCar),
    ]),
    concatMap(
      ([, selectedCar, selectedBike]) => {
        const actionsArray: Action[] = [];

        // 保持原有的Action构建逻辑不变
        if (selectedCar) {
          actionsArray.push(actions.addCar());
        } else {
          actionsArray.push(actions.removeCar());
        }

        if (selectedBike) {
          actionsArray.push(actions.addBike());
        } else {
          actionsArray.push(actions.removeBike());
        }

        // ... 更多产品的判断逻辑

        // 将Action数组转为串行流,每个Action等待对应操作完成
        return from(actionsArray).pipe(
          concatMap(action => {
            // 先发出当前Action
            return of(action).pipe(
              // 等待当前Action对应的成功或失败Action
              concatMap(() => {
                // 根据当前Action类型,匹配对应的完成Action
                let completionActions: TypedAction<string>[];
                switch (action.type) {
                  case actions.addCar.type:
                    completionActions = [actions.addCarSuccess, actions.addCarFailure];
                    break;
                  case actions.removeCar.type:
                    completionActions = [actions.removeCarSuccess, actions.removeCarFailure];
                    break;
                  case actions.addBike.type:
                    completionActions = [actions.addBikeSuccess, actions.addBikeFailure];
                    break;
                  case actions.removeBike.type:
                    completionActions = [actions.removeBikeSuccess, actions.removeBikeFailure];
                    break;
                  // 添加其他产品对应的成功/失败Action
                  default:
                    completionActions = [];
                }

                return this.actions$.pipe(
                  ofType(...completionActions.map(a => a.type)),
                  take(1) // 只取第一个完成的Action,避免无限等待
                );
              }),
              // 只需要发出原Action,后面的完成Action是用来等待的,不需要再次分发
              map(() => action)
            );
          })
        );
      }
    )
  )
);
原理说明
  • from(actionsArray):把Action数组转为逐个发出的Observable流
  • concatMap:确保每个Action的处理(分发+等待完成)完成后,再处理下一个Action
  • this.actions$.pipe(ofType(...), take(1)):监听当前Action对应的成功或失败Action,确保我们只等待一次操作完成,避免内存泄漏

这样修改后,每个产品对应的HTTP请求会严格按顺序执行,上一个请求完成(不管成功还是失败)后,才会触发下一个请求,完全不需要依赖延迟。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.20 08:14:55