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

Node.js Lambda中DynamoDB batchWrite批量写入部分数据无报错失败

解决DynamoDB batchWrite递归上传时剩余数据丢失且无日志的问题

我一眼就看出了问题的核心——你的递归调用没正确等待异步操作完成,再加上错误处理和状态码设置的疏漏,导致最后10条数据的写入被Lambda提前终止,连执行或日志记录的机会都没有。

问题根源

  1. 递归调用未加await:处理超过25条数据的分支里,你调用uploadFileByBatch(payload)时没加await,Lambda会直接跳过剩余数据的写入流程,提前返回响应。函数一旦返回,未完成的异步任务会被强制终止,所以最后10条数据根本没写完,自然也不会有日志。
  2. 错误未向上层传递:uploadFileByBatch里的catch块只打印错误,却没把错误抛出去,上层handler完全感知不到出错,导致状态码一直是0,也没有错误响应返回。
  3. 未处理UnprocessedItems:DynamoDB的batchWrite可能因为吞吐量限制返回未处理的项目,你没处理这种情况,可能导致部分数据悄悄丢失。
  4. 状态码初始化错误:statusCode初始设为0,成功时也没改成200,响应状态完全不符合规范。

修复方案

  • 给递归调用uploadFileByBatch加上await,确保所有批次的写入都完成后再返回响应。
  • 让uploadFileByBatch把错误向上抛出,让上层handler能捕获并返回错误响应。
  • 新增UnprocessedItems的重试逻辑,确保所有数据都能被成功写入。
  • 正确初始化和设置statusCode:成功时设为200,出错时设为500。

修改后的完整代码

"use strict"
const AWS = require('aws-sdk');
const sha1 = require('sha1');
const documentClient = new AWS.DynamoDB.DocumentClient();

exports.handler = async function (event, context, callback) {
  let responseBody = "";
  let statusCode = 200; // 默认成功状态码
  try {
    const { excelObject } = JSON.parse(event.body);
    if(excelObject){
      await uploadFileByBatch(excelObject);
      responseBody = JSON.stringify({ message: "所有数据上传成功" });
    } else {
      responseBody = JSON.stringify({ message: "未提供有效数据" });
    }
  } catch (err) {
    statusCode = 500;
    responseBody = JSON.stringify({ error: err.message });
    console.error("上传失败:", err);
  }
  const response = {
    statusCode: statusCode,
    headers:{
      "Content-Type": "application/json",
      "access-control-allow-origin": "*"
    },
    body: responseBody
  }
  console.log(response)
  return response
}

let uploadFileByBatch = async function (payload) {
  // 处理未处理的项目(如果存在)
  if (payload.UnprocessedItems) {
    payload = payload.UnprocessedItems.Community?.map(item => item.PutRequest.Item) || [];
    if (payload.length === 0) return;
  }

  const items = [];
  const batchSize = 25;
  const currentBatch = payload.length > batchSize ? payload.slice(0, batchSize) : payload;
  
  currentBatch.forEach(obj =>{
    const hash = sha1(Buffer.from(new Date().toString()+ Math.random()));
    items.push(
      {
        PutRequest:{
          Item: {
            id: obj.id?obj.id:hash,
            organization_EN: obj.organization_EN,
            email: obj.email,
            isActive: obj.isActive
          }
        }
      }
    )
  })

  const params = {
    RequestItems:{
      "Community": items
    }
  }

  console.log(`正在上传批次,数量:${currentBatch.length}`);
  const data = await documentClient.batchWrite(params).promise();
  
  // 如果有未处理项目,递归重试
  if (Object.keys(data.UnprocessedItems).length > 0) {
    console.log("存在未处理项目,开始重试");
    await uploadFileByBatch(data);
  } 
  // 还有剩余数据,继续处理下一批
  else if (payload.length > batchSize) {
    const remainingPayload = payload.slice(batchSize);
    await uploadFileByBatch(remainingPayload);
  }
}

关键改进点

  • 递归调用添加await,确保所有批次和重试操作都完成后再返回。
  • 新增UnprocessedItems重试逻辑,解决DynamoDB吞吐量限制导致的部分数据未写入问题。
  • 上层handler添加全局try/catch,统一捕获错误并返回规范的错误响应。
  • 修正状态码设置,成功返回200、失败返回500,同时返回清晰的响应信息。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.11 07:57:25