RxJS中concatMap未等待前序Promise执行,如何保证队列逻辑同步执行?
问题根因
你遇到的异常核心是 Promise是立即执行(eager) 的:当你在concatMap的回调函数中直接new Promise()时,回调内的同步代码会立刻执行,并不会等到前一个任务完成后再触发。concatMap的串行逻辑只有在处理「惰性执行的Observable」时才会符合预期,Promise不属于这类。
方案1:基于concatMap调整逻辑,最小改动兼容原有写法
用RxJS的defer操作符包裹执行逻辑,保证只有concatMap订阅当前任务时,才会执行更新代码,完美实现串行效果:
import { defer } from 'rxjs'; this.finalize$.pipe( concatMap((field: string) => { // defer包裹后逻辑为惰性执行,仅订阅时触发 return defer(() => { console.log('START', field); this.alreadyExecutedFields.push(field); const remainingFields = this.remainingFieldsToExecute$.value; const targetIndex = remainingFields.indexOf(field); // 新增边界判断,避免字段不存在时的splice异常 if (targetIndex > -1) { remainingFields.splice(targetIndex, 1); } this.remainingFieldsToExecute$.next(remainingFields); console.log('END', field); return void 0; }); }) ).subscribe();
如果需要加延迟测试串行逻辑,可以在defer返回值后拼接delay(1000)操作符,能明显看到任务按顺序执行的效果。
方案2:基于状态流重构,从根源避免并行问题
你的需求本质是基于完结事件更新两个状态,完全可以用RxJS原生的状态操作符实现,不需要依赖串行锁,也不会出现索引异常:
import { scan, tap } from 'rxjs'; // 初始化状态 this.alreadyExecutedFields = []; this.remainingFieldsToExecute$ = new BehaviorSubject<string[]>([/* 你的初始待执行字段 */]); this.finalize$.pipe( // scan累积所有已完结字段,自动去重,内置原子性保证不会出现并发状态冲突 scan((executedFields, currentField) => { if (!executedFields.includes(currentField)) { executedFields.push(currentField); } return executedFields; }, [] as string[]), // 同步更新两个状态列表 tap(latestExecutedFields => { this.alreadyExecutedFields = latestExecutedFields; // 直接过滤生成新的待执行列表,不修改原数组,彻底避免索引操作异常 const remainingFields = this.remainingFieldsToExecute$.value.filter( field => !latestExecutedFields.includes(field) ); this.remainingFieldsToExecute$.next(remainingFields); }) ).subscribe();
这种方案不需要考虑串行逻辑,scan本身会按事件触发顺序处理,每次计算都是基于上一次的最新状态,完全不会出现并发修改的问题,代码可维护性更高。
内容的提问来源于stack exchange,提问作者jmeire
相关产品推荐
相关产品推荐

