You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何交错/合并异步可迭代对象?

并行消费异步迭代器,优先输出就绪值

嘿,我懂你的需求——原来的串行拼接会让你浪费时间等待慢的迭代器,现在要让所有异步迭代器同时跑,谁先产出值就先输出,而且不能把所有元素都缓冲起来,毕竟还可能是无限序列对吧?

核心思路

我们需要同时启动所有异步迭代器的遍历流程,然后用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));

工作原理拆解

  1. 初始化并行任务:为每个输入的异步迭代器创建实例,立刻调用next()并把返回的Promise加入pending集合,这样所有迭代器会同时开始执行各自的逻辑。
  2. 竞赛式等待:每次用Promise.race找到pending中最先完成的Promise——不管是产出了新值、迭代器完成,还是抛出了错误。
  3. 产出与续跑:如果是正常产出值,先把值yield出去,然后立刻给该迭代器发起下一个next()请求,把新的Promise放回pending,让它继续参与下一轮竞赛。
  4. 错误处理:默认情况下,单个迭代器抛出的错误会终止整个合并过程;如果你希望单个错误不影响其他迭代器,可以把throw winner.error换成错误日志加continue。

这种方式完全符合你的需求:不需要缓冲所有元素,逐个产出;不关心顺序,最快就绪的值先输出;也能处理无限序列——只要不终止循环,它会一直消费所有迭代器的产出。

内容的提问来源于stack exchange,提问作者sdgfsdh

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.29 08:00:35