如何在Lambda中存储最新时间值,实现API增量获取避免S3重复
Lambda增量同步API数据到S3,避免重复记录
问题背景
每5分钟运行Lambda从客户API拉取数据时,API返回过去7天的全量数据,导致S3重复存储。需要实现:读取上一次同步的最新CreatedWhen时间,作为API请求参数获取增量数据,同时保存最新时间供下次使用。
实现方案
Lambda本地存储是临时的,必须用S3存储元数据来持久化上次同步的时间。具体步骤:
1. 读取上一次同步的时间
在S3中创建一个专门的元数据文件(比如V1/Charge Code/last_sync_time.json),每次Lambda启动时先读取这个文件。如果文件不存在,使用初始时间1970-01-01T00:00:00.0000000。
2. 带时间参数请求API
将读取到的时间作为过滤参数传入API请求,只获取该时间之后创建的数据(需确认API支持的参数格式,示例假设用CreatedWhenAfter参数)。
3. 提取并保存最新的CreatedWhen时间
遍历API返回的ChargeCodes数组,找出最大的CreatedWhen值,写入S3的元数据文件,覆盖旧值。
4. 仅上传增量数据到S3
将本次获取的增量数据保存到S3,避免重复。
修改后的完整代码
var AWS = require('aws-sdk'); AWS.config.update({ region: 'us-east-1' }); var s3 = new AWS.S3(); const https = require('https'); // 元数据文件的S3路径 const LAST_SYNC_FILE_KEY = 'V1/Charge Code/last_sync_time.json'; // 初始同步时间(首次运行用) const DEFAULT_SYNC_TIME = '1970-01-01T00:00:00.0000000'; // 读取上一次同步时间 const getLastSyncTime = async () => { try { const response = await s3.getObject({ Bucket: 'commtrac', Key: LAST_SYNC_FILE_KEY }).promise(); const data = JSON.parse(response.Body.toString()); return data.lastSyncTime || DEFAULT_SYNC_TIME; } catch (err) { // 文件不存在时返回初始时间 if (err.code === 'NoSuchKey') { return DEFAULT_SYNC_TIME; } throw err; } }; // 保存最新同步时间到S3 const saveLastSyncTime = async (lastSyncTime) => { await s3.putObject({ Bucket: 'commtrac', Key: LAST_SYNC_FILE_KEY, Body: JSON.stringify({ lastSyncTime }), ContentType: 'application/json' }).promise(); }; // 提取响应中最新的CreatedWhen时间 const getLatestCreatedWhen = (apiResponse) => { const chargeCodes = JSON.parse(apiResponse).ChargeCodes || []; if (chargeCodes.length === 0) { return DEFAULT_SYNC_TIME; } // 按CreatedWhen降序排序,取第一个 chargeCodes.sort((a, b) => new Date(b.CreatedWhen) - new Date(a.CreatedWhen)); return chargeCodes[0].CreatedWhen; }; const doGetRequest = (lastSyncTime) => { return new Promise((resolve, reject) => { const username = "xxx"; const password = "xxx"; const auth = "Basic " + Buffer.from(username + ":" + password).toString("base64"); // 拼接带时间过滤的API路径(需确认API实际参数名,这里假设是CreatedWhenAfter) const path = `/rest/v1/chargeCode?CreatedWhenAfter=${encodeURIComponent(lastSyncTime)}`; const options = { host: 'www.xxxxxxxxxxx.com', path: path, headers: { Authorization: auth}, method: 'GET' }; var body=''; const req = https.request(options, (res) => { res.on('data', function (chunk) { body += chunk; }); res.on('end', function () { console.log("增量数据结果", body.toString()); resolve(body); }); }); req.on('error', (e) => { reject(e.message); }); req.end(); }); }; exports.handler = async (event, context, callback) => { try { // 1. 获取上次同步时间 const lastSyncTime = await getLastSyncTime(); // 2. 请求增量数据 const result = await doGetRequest(lastSyncTime); // 3. 提取本次最新的CreatedWhen时间 const latestCreatedWhen = getLatestCreatedWhen(result); // 4. 保存最新同步时间到S3 await saveLastSyncTime(latestCreatedWhen); // 5. 生成文件名并上传增量数据到S3 let date_ob = new Date(); let year = date_ob.getFullYear(); let month = ("0" + (date_ob.getMonth() + 1)).slice(-2); let day = date_ob.getDate(); let hours = date_ob.getHours(); let minutes = date_ob.getMinutes(); let seconds = date_ob.getSeconds(); const bucket_name = 'Charge Code - ' + year + '-' + month + '-' + day; const objectName = 'Charge Code - ' + year + '-' + month + '-' + day +'-' + hours + '-' + minutes + '-' + seconds +'.json'; const uploadParams = { Bucket: 'commtrac', Key: `V1/Charge Code/${bucket_name}/${objectName}`, Body: result, ContentType: 'application/json', ACL: 'public-read' }; await s3.upload(uploadParams).promise(); console.log('增量数据上传成功', bucket_name + ' || ' + objectName); // 返回响应 const response = { "statusCode": 200, "headers": {"my_header": "my_value"}, "body": JSON.stringify(result), "isBase64Encoded": false }; callback(null, response); } catch (error) { console.error('同步失败', error); callback(error); } };
关键说明
- 元数据文件:用
last_sync_time.json存储上次同步的最新时间,确保Lambda重启后能读取到历史值。 - API参数适配:需要确认客户API支持的时间过滤参数名(比如可能是
since、modifiedAfter等),修改path中的参数即可。 - 时间排序:通过排序ChargeCodes数组,确保提取到的是最新的
CreatedWhen值。 - 错误处理:处理元数据文件不存在的情况,同时用
async/await简化异步逻辑,避免回调嵌套。
内容的提问来源于stack exchange,提问作者rusiru guruge
相关产品推荐
相关产品推荐

