如何在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对应的请求完成(成功或失败)。
具体做法是:
- 将Action数组转换为Observable流(用
from()) - 对每个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的处理(分发+等待完成)完成后,再处理下一个Actionthis.actions$.pipe(ofType(...), take(1)):监听当前Action对应的成功或失败Action,确保我们只等待一次操作完成,避免内存泄漏
这样修改后,每个产品对应的HTTP请求会严格按顺序执行,上一个请求完成(不管成功还是失败)后,才会触发下一个请求,完全不需要依赖延迟。
内容的提问来源于stack exchange,提问作者Marek
相关产品推荐
相关产品推荐

