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
相关产品推荐
相关产品推荐

