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

JavaScript如何实现fork函数以复用Async Generator异步生成器?

异步生成器fork函数实现

核心实现满足两个硬性要求:

  • 原始异步生成器仅迭代一次,避免重复执行高开销的数据库查询
  • 仅缓冲不同副本消费速度差异部分的数据,不会全量加载到内存
function fork(asyncGenerator, copyCount = 2) {
    // 为每个副本维护独立的待消费队列和等待resolve回调列表
    const copyQueues = Array.from({length: copyCount}, () => []);
    const waitResolvers = Array.from({length: copyCount}, () => []);
    let sourceDone = false;
    let isPullingSource = false;

    // 统一从原始生成器拉取数据,分发给所有副本队列
    async function pullFromSource() {
        if (isPullingSource || sourceDone) return;
        isPullingSource = true;
        try {
            const result = await asyncGenerator.next();
            if (result.done) {
                sourceDone = true;
                // 给所有副本发送结束标记
                copyQueues.forEach(queue => queue.push({done: true}));
                // 唤醒所有等待数据的副本迭代器
                waitResolvers.forEach(resolverList => resolverList.forEach(resolve => resolve()));
                return;
            }
            // 新数据推送到所有副本队列
            copyQueues.forEach(queue => queue.push({value: result.value, done: false}));
            // 唤醒等待数据的迭代器
            waitResolvers.forEach((resolverList, idx) => {
                const currentQueue = copyQueues[idx];
                while (resolverList.length && currentQueue.length) resolverList.shift()();
            });
        } finally {
            isPullingSource = false;
        }
    }

    // 为每个副本生成独立的异步迭代器
    return copyQueues.map((queue, copyIndex) => {
        const myResolvers = waitResolvers[copyIndex];
        return async function* () {
            while (true) {
                // 队列有数据直接消费
                if (queue.length) {
                    const item = queue.shift();
                    if (item.done) return;
                    yield item.value;
                    continue;
                }
                // 原始生成器已结束直接退出
                if (sourceDone) return;
                // 无数据则等待新数据拉取
                await new Promise(resolve => myResolvers.push(resolve));
                // 触发原始数据拉取
                pullFromSource();
            }
        }();
    });
}

使用说明

  • 调用fork(events, n)可以生成n个独立副本,默认生成2个副本
  • 内存占用仅由消费最慢的副本的落后量决定,只要所有副本都在持续消费,就不会出现内存溢出
  • 原始生成器只会被遍历一次,不会重复触发数据库查询

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.04 04:45:03