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

如何在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.16 14:26:58