Node.js Lambda使用async/await时URL重复写入DynamoDB问题
重复请求、重复写入的核心原因
三个直接诱因,按出现概率排序:
- Lambda异步调用默认重试机制触发:定时事件触发Lambda属于异步调用,AWS默认对执行失败的异步调用做最多2次自动重试。只要函数执行超时、抛出未捕获错误,整个脚本就会被完整重跑,所有URL都会重新请求、重新写入,直接产生重复数据。
- 代码逻辑缺陷放大重试概率:
- 全局声明的
var id属于执行环境共享变量,Lambda热启动(执行环境复用,非冷启动)时会残留旧值,极端情况会出现数据覆盖、重复写入 getLighthouse方法中,PSI接口请求报错时仅打印日志,无任何返回值,默认返回undefined;后续执行insertRecords时读取undefined.lighthouseResult会抛出类型错误,虽然你加了catch不会中断循环,但如果请求长时间挂起导致Lambda整体超时,还是会触发重试- 未给axios配置超时时间,PSI接口偶尔响应慢时,请求会一直挂起直到触发Lambda超时
- 全局声明的
- 无幂等校验:DynamoDB写入时没有唯一键约束,不管是重试产生的重复请求还是代码逻辑问题导致的重复写入,都会直接落库,没有自动去重能力。
修复方案
按优先级操作即可解决问题:
- 先关闭不必要的自动重试:进入Lambda控制台,找到对应函数的「配置-异步调用」设置项,将重试次数修改为0,避免单次执行失败就全量重跑所有URL。
- 修复代码逻辑漏洞:
- 删除全局的
id变量,将id声明移到循环内部,避免热启动的变量污染 - 给axios请求配置30秒左右的合理超时,避免请求无限挂起
getLighthouse请求失败时明确返回null,拿到PSI结果后先做有效性校验,结果异常直接跳过当前URL,不执行后续写入逻辑- 读取URL列表的逻辑放在handler内部,避免热启动时缓存旧的URL列表
- 删除全局的
- 增加幂等校验:DynamoDB写入时增加
URL+小时级时间戳的唯一键,配置写入条件,同一个小时内同一个URL只允许写入一条数据,从根本上杜绝重复落库。
修复后的参考代码
const { v4: uuidv4 } = require('uuid'); const axios = require('axios'); const AWS = require('aws-sdk'); const docClient = new AWS.DynamoDB.DocumentClient(); // 固定配置放在全局即可 const tableName = '替换为你的DynamoDB表名'; const API_KEY = '替换为你的PSI接口密钥'; const endpoint = 'https://www.googleapis.com/pagespeedonline/v5/runPagespeed'; // 给axios统一配置30秒超时,避免请求挂起 const psiRequest = axios.create({ timeout: 30000 }); const insertRecords = async (_id, _url, _lighthouseResults) => { const metrics = _lighthouseResults.lighthouseResult.audits.metrics.details.items[0]; const params = { TableName: tableName, // 写入条件:唯一键不存在时才写入,从数据库层拦截重复数据 ConditionExpression: 'attribute_not_exists(hour_mark)', Item: { id: _id, hour_mark: `${_url}_${new Date().toISOString().slice(0,13)}`, // 小时级唯一键 created_at: new Date().toISOString(), URL: _url, metrics }, }; return docClient.put(params).promise(); } exports.putItemHandler = async (event) => { // URL列表读取逻辑放在handler内,避免热启动缓存旧数据 const urls = require('./urls.json'); // 替换为你实际的URL读取逻辑 for (const url of urls) { // id声明在循环内部,不使用全局变量 const recordId = uuidv4(); console.log(`${url} - ${recordId}`); try { const lighthouseResults = await getLighthouse(url); // 结果无效直接跳过,不执行写入 if (!lighthouseResults || !lighthouseResults.lighthouseResult) { console.log(`获取${url}的PSI结果失败,跳过当前URL`); continue; } await insertRecords(recordId, url, lighthouseResults); } catch (error) { console.log(`处理${url}出错:`, error.message); // 捕获错误后继续下一个URL,不抛出错误终止整个函数 } } console.log("全部URL处理完成"); }; const getLighthouse = async (url) => { console.log("开始请求PSI接口,目标URL:", url) try { const resp = await psiRequest.get(endpoint, { params: { key: API_KEY, url: url, category: 'performance', strategy: 'mobile' } }); return resp.data } catch (err) { console.error(`请求${url}的PSI数据失败:`, err.message); // 明确返回null,避免返回undefined导致后续取值报错 return null; } }
内容的提问来源于stack exchange,提问作者Trey Copeland
相关产品推荐
相关产品推荐

