如何交错/合并异步可迭代对象?
并行消费异步迭代器,优先输出就绪值
嘿,我懂你的需求——原来的串行拼接会让你浪费时间等待慢的迭代器,现在要让所有异步迭代器同时跑,谁先产出值就先输出,而且不能把所有元素都缓冲起来,毕竟还可能是无限序列对吧?
核心思路
我们需要同时启动所有异步迭代器的遍历流程,然后用Promise.race每次等待当前所有未完成的迭代步骤中最快完成的那个,产出值之后立刻继续该迭代器的下一个步骤,直到拿到足够的元素或者所有迭代器结束。
完整实现代码
下面是包含原有测试对象和新合并逻辑的完整代码:
// Promisified sleep function const sleep = ms => new Promise((resolve, reject) => { setTimeout(() => resolve(ms), ms); }); const a = { [Symbol.asyncIterator]: async function * () { yield 'a'; await sleep(1000); yield 'b'; await sleep(2000); yield 'c'; }, }; const b = { [Symbol.asyncIterator]: async function * () { await sleep(6000); yield 'i'; yield 'j'; await sleep(2000); yield 'k'; }, }; const c = { [Symbol.asyncIterator]: async function * () { yield 'x'; await sleep(2000); yield 'y'; await sleep(8000); yield 'z'; await sleep(10000); throw new Error('You have gone too far! '); }, }; // 合并多个异步迭代器,并行产出值 async function* mergeAsyncIterables(...iterables) { const pending = new Set(); // 初始化:启动所有迭代器的第一个next请求 for (const iterable of iterables) { const iterator = iterable[Symbol.asyncIterator](); pending.add((async () => { try { const result = await iterator.next(); return { iterator, result, error: null }; } catch (err) { return { iterator, result: null, error: err }; } })()); } while (pending.size > 0) { // 等待最快完成的迭代步骤 const winner = await Promise.race(pending); pending.delete(winner); // 处理错误情况 if (winner.error) { throw winner.error; // 如果不想因为单个迭代器错误终止全部,可以替换成: // console.error('Iterator error:', winner.error); // continue; } if (!winner.result.done) { // 产出当前就绪的值 yield winner.result.value; // 继续该迭代器的下一个步骤 pending.add((async () => { try { const result = await winner.iterator.next(); return { iterator: winner.iterator, result, error: null }; } catch (err) { return { iterator: winner.iterator, result: null, error: err }; } })()); } // 如果result.done为true,说明该迭代器已完成,无需后续处理 } } // 测试代码,保留你原来的取值逻辑 (async () => { const limit = 9; let i = 0; const xs = []; for await (const x of mergeAsyncIterables(a, b, c)) { xs.push(x); i++; if (i === limit) { break; } } console.log(xs); // 预期输出顺序(可能因定时器精度略有差异): ['a', 'x', 'b', 'y', 'c', 'i', 'j', 'k', 'z'] })().catch(error => console.error(error));
工作原理拆解
- 初始化并行任务:为每个输入的异步迭代器创建实例,立刻调用
next()并把返回的Promise加入pending集合,这样所有迭代器会同时开始执行各自的逻辑。 - 竞赛式等待:每次用
Promise.race找到pending中最先完成的Promise——不管是产出了新值、迭代器完成,还是抛出了错误。 - 产出与续跑:如果是正常产出值,先把值
yield出去,然后立刻给该迭代器发起下一个next()请求,把新的Promise放回pending,让它继续参与下一轮竞赛。 - 错误处理:默认情况下,单个迭代器抛出的错误会终止整个合并过程;如果你希望单个错误不影响其他迭代器,可以把
throw winner.error换成错误日志加continue。
这种方式完全符合你的需求:不需要缓冲所有元素,逐个产出;不关心顺序,最快就绪的值先输出;也能处理无限序列——只要不终止循环,它会一直消费所有迭代器的产出。
内容的提问来源于stack exchange,提问作者sdgfsdh
相关产品推荐
相关产品推荐

