NgRx Effects合并多个Action并逐个执行失败,求解决方案
问题分析
你遇到的报错核心原因是:forkJoin和merge是用来合并Observable流的,但你直接传入了Action实例(普通JavaScript对象)。RxJS会把这些对象当作Observable处理,最终返回的不是有效的Action,而是一个不符合NgRx要求的Observable对象,所以才会抛出"Actions must have a type property"和无效action的错误。
正确实现方案
根据你的需求(逐个执行Action),分两种场景给出解决方案:
场景1:所有Action都是同步执行,不需要等待异步操作
如果OpenUsedDroneUpdateChannelRequest和GetUsedDronesSuccess都是同步Action(不需要等待后端请求或其他异步操作完成),你可以直接返回Action数组,或者用of()把Action包装成Observable——NgRx Effects支持直接返回Action数组,会自动逐个dispatch每个Action:
map(drones => { const actions = []; // 逐个添加打开通道的Action drones.forEach((drone) => actions.push( new featureActions.OpenUsedDroneUpdateChannelRequest({ droneId: drone.id, projectId : environment.projectId }) )); // 添加最终的成功Action actions.push(new featureActions.GetUsedDronesSuccess({drones})); // 直接返回Action数组,NgRx会自动逐个dispatch return actions; })
或者用rxjs的of操作符包装(效果完全一致,写法不同):
import { of } from 'rxjs'; // ... map(drones => { const actions = []; drones.forEach((drone) => actions.push( new featureActions.OpenUsedDroneUpdateChannelRequest({ droneId: drone.id, projectId : environment.projectId }) )); actions.push(new featureActions.GetUsedDronesSuccess({drones})); // 将Action数组展开为Observable,逐个发出每个Action return of(...actions); })
场景2:需要等待所有OpenUsedDroneUpdateChannelRequest的异步操作完成后,再执行GetUsedDronesSuccess
如果OpenUsedDroneUpdateChannelRequest是触发异步操作的Action(比如对应一个Effect发起WebSocket连接或后端请求),你需要等待所有异步操作完成后再dispatch成功Action。这时需要把每个异步操作的完成事件转化为Observable,用forkJoin等待全部完成:
假设你的OpenUsedDroneUpdateChannelRequest对应的Effect会在操作完成后dispatch一个OpenUsedDroneUpdateChannelSuccessAction(包含对应的droneId),可以这样实现:
import { forkJoin, of } from 'rxjs'; import { switchMap, concat } from 'rxjs/operators'; import { Actions } from '@ngrx/effects'; // 确保你已经注入了Actions服务 constructor(private actions$: Actions) {} // ... switchMap(drones => { // 创建所有需要dispatch的打开通道Action const openChannelActions = drones.map(drone => new featureActions.OpenUsedDroneUpdateChannelRequest({ droneId: drone.id, projectId : environment.projectId }) ); // 为每个打开通道Action创建一个Observable,监听对应的成功Action const openChannelCompletion$ = forkJoin( openChannelActions.map(action => this.actions$.pipe( ofType(featureActions.OpenUsedDroneUpdateChannelSuccess), // 确保只监听当前drone对应的成功Action filter(success => success.payload.droneId === action.payload.droneId), take(1) // 只取第一个完成信号 ) ) ); // 先逐个dispatch打开通道的Action,等待全部完成后dispatch成功Action return of(...openChannelActions).pipe( concat( openChannelCompletion$.pipe( map(() => new featureActions.GetUsedDronesSuccess({drones})) ) ) ); })
为什么之前的方式失败?
再回头看你的两种尝试:
forkJoin(actions):forkJoin要求传入Observable数组,但你传的是Action实例(普通对象)。RxJS会尝试把这些对象转换成Observable,但结果是一个无效的Observable(带有_isScalar: false标识),不是NgRx认可的Action,所以报错。merge(actions):同理,merge也是用来合并Observable流的,传入普通Action对象会导致同样的类型错误。
内容的提问来源于stack exchange,提问作者Crocsx

