RxJS中如何实现MergeMap序列在任务完成后立即启动新任务以维持指定并发数?
嘿,这个问题我之前做批量文件上传功能的时候刚好踩过一样的坑!你现在的核心问题是把「攒够一批再执行、等整批完才续下一批」的模式,改成「一个任务完成就立刻补新任务、始终维持指定并发数」的模式,其实根本不用绕弯子,RxJS的mergeMap本身就支持这个能力。
先说说你现有代码的问题:你在mergeMap后面加了bufferCount(batchQtnFiles),这个操作符会先把waitingFiles里的任务攒够5个再一次性发给后续操作,然后必须等这5个上传任务全完成,才会处理下一批攒好的5个,这就导致了“批量等待”的情况,不是真正的“动态维持并发数”。
正确的实现方式
直接在mergeMap的第二个参数里指定最大并发数就搞定了!mergeMap的第二个参数就是专门用来设置「同时允许活跃的内部Observable数量」的,去掉多余的bufferCount,代码改成这样:
from(waitingFiles) .pipe( // 第二个参数就是我们要维持的最大并发数(这里是batchQtnFiles,也就是5) mergeMap((file: UploadQueue) => { return this.doUpload(file, documentId, documentVersion, saveDocument); }, batchQtnFiles), tap(() => { this.updateQtyPending(saveDocument, documentId); }) ) .subscribe();
为什么这样有效?
当你给mergeMap设置了并发参数后,它会同时最多订阅n个由doUpload返回的上传任务Observable。每当其中一个上传任务完成(内部Observable执行完毕),mergeMap会立刻从源Observable(waitingFiles)里取出下一个待上传的文件,创建新的上传任务订阅,自动补位,始终保持最多5个活跃的上传任务在运行——完全不需要手动控制分组和等待。
另外补充个小细节:如果你的waitingFiles不是一次性的数组,而是一个会持续推送新文件的异步Observable(比如用户不断添加新文件到队列),这个逻辑也能完美自适应,只要源有新的文件进来,mergeMap会在并发有空位时立刻启动新的上传。
内容来源于stack exchange

