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

Node.js用async迭代器解决AWS CloudWatch日志上传sequenceToken冲突问题

问题根因

原代码通过on('data')绑定异步回调,事件触发时会直接执行回调逻辑,不会等待前一次回调的异步操作完成,导致多条日志上传请求并发发起,后发起的请求用的还是旧的sequenceToken,触发Cloudwatch接口报错。

修复方案(基于async迭代器实现)

利用Node.js可读流原生支持的async迭代器特性,用for await...of顺序消费日志流,天然保证前一条日志上传完成、更新完token后,才会处理下一条日志:

const build = require('pino-abstract-stream');

const stream = async (options) => {
    // 创建AWS连接
    const client = await createClient();
 
    // 获取初始sequenceToken
    let sequenceToken = await getInitSequenceToken(client);

    return build(function (source) {
        // 用async迭代器顺序消费日志流
        (async () => {
            for await (const obj of source) {
                const command = new PutLogEventsCommand({
                    logGroupName: "api",
                    logStreamName: `executive-${Config.env}-${Config.location}`,
                    logEvents: [
                        {
                            message: obj.msg,
                            timestamp: obj.time,
                        },
                    ],
                    sequenceToken,
                });

                try {
                    // 等待当前上传请求完成
                    const response = await client.send(command);
                    // 更新token,下一次循环直接使用新值
                    sequenceToken = response.nextSequenceToken;
                } catch (err) {
                    // 可根据业务需要添加上传失败的降级逻辑,比如重试、落本地缓存等
                    console.error('日志上传Cloudwatch失败', err);
                }
            }
        })();
    });
};
优化建议(可选)

如果日志量较大,单条上传性能太低,可以在迭代器内添加攒批逻辑:

  • 维护一个日志数组缓冲区,每次循环先把日志 push 进缓冲区
  • 当缓冲区长度达到阈值(Cloudwatch单接口最多支持1万条/5MB大小)或者超过指定超时时间(比如500ms),再批量调用接口上传
  • 上传完成后清空缓冲区,继续消费后续日志

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.04 19:24:01