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

使用Lambda函数向DynamoDB上传实时数据:测试正常生产失败求助

Kinesis触发Lambda写入DynamoDB生产环境失败问题排查与解决

问题背景

通过Kinesis数据生成器结合Lambda函数将实时数据上传至AWS DynamoDB,测试用例运行正常,但生产环境无法完成数据上传。

现有代码与数据示例

Lambda函数代码

import boto3
import json
import base64

def lambda_handler(event, context):
   dynamo_db = boto3.resource('dynamodb')
   table = dynamo_db.Table('dyno-table')
   try:
      data = [record.get('kinesis').get('data') for record in event['Records']]
    
      for record in data:
         recordeddata = base64.b64decode(record).decode("utf-8")
         decoded_data_dic = json.loads(recordeddata)
         with table.batch_writer() as batch_writer:
            batch_writer.put_item(Item=decoded_data_dic)
      return data
    except Exception as e:
      return str(e)

Kinesis发送的JSON数据示例

{
"Records": [
{
  "kinesis": {
    "kinesisSchemaVersion": "1.0",
    "partitionKey": "4159727726",
    "sequenceNumber": "49639282004144167688658152268735854270994120426068639794",
    "data": "eyJmaXJzdE5hbWUiOiJBbnRvbmV0dGUiLCJsYXN0TmFtZSI6IlRoaWVsIiwiaWQiOjYsImlwIjoiMTAwLjIyNy4zOC4yMTMiLH0K",
    "approximateArrivalTimestamp": 1680243715.662
  },
  "eventSource": "aws:kinesis",
  "eventVersion": "1.0",
  "eventID": "shardId-000000000003:49639282004144167688658152268735854270994120426068639794",
  "eventName": "aws:kinesis:record",
  "invokeIdentityArn": "arn:aws:iam::777050133147:role/AvaniK-LambdaRole",
  "awsRegion": "ap-northeast-1",
  "eventSourceARN": "arn:aws:kinesis:ap-northeast-1:777050133147:stream/DeepakR-datastream"
}
]
}

数据生成Schema

{
  "firstName":"{{name.firstName}}",
  "lastName":"{{name.lastName}}",
  "id":{{random.number(70)}},
  "ip":"{{internet.ip}}",
}

问题根源分析

  1. 无效JSON格式:数据Schema中最后一个字段"ip":"{{internet.ip}}",末尾多了一个逗号,导致生成的JSON存在尾逗号(JSON规范不允许尾逗号),json.loads解析时会抛出异常,中断数据写入流程。测试用例可能使用了手动构造的合法JSON,因此未触发该问题。
  2. BatchWriter使用不当:原代码在循环内重复创建batch_writer,每条数据单独开启批量写入会话,效率低下且可能导致资源未正确释放。
  3. 错误处理不规范:原代码仅返回错误字符串,未抛出异常导致Kinesis无法触发重试,同时缺少日志记录,难以排查生产环境问题。
  4. 权限隐患:生产环境Lambda角色可能缺少DynamoDB批量写入权限(测试环境可能使用高权限角色)。

解决方案

1. 修正数据Schema

移除最后一个字段末尾的逗号,确保生成合法JSON:

{
  "firstName":"{{name.firstName}}",
  "lastName":"{{name.lastName}}",
  "id":{{random.number(70)}},
  "ip":"{{internet.ip}}"
}

2. 优化Lambda函数代码

调整BatchWriter位置,增加日志记录与异常抛出,方便调试与重试:

import boto3
import json
import base64
import logging

# 配置日志输出
logger = logging.getLogger()
logger.setLevel(logging.INFO)

def lambda_handler(event, context):
    dynamo_db = boto3.resource('dynamodb')
    table = dynamo_db.Table('dyno-table')
    try:
        # 单次批量写入会话处理所有记录
        with table.batch_writer() as batch_writer:
            for record in event['Records']:
                kinesis_data = record.get('kinesis', {}).get('data')
                if not kinesis_data:
                    logger.warning("检测到空Kinesis数据,跳过处理")
                    continue
                # 分步处理解码与解析,单独捕获异常
                try:
                    decoded_str = base64.b64decode(kinesis_data).decode("utf-8")
                    data_item = json.loads(decoded_str)
                    batch_writer.put_item(Item=data_item)
                    logger.info(f"成功写入数据: {data_item}")
                except base64.binascii.Error as decode_err:
                    logger.error(f"Base64解码失败: {str(decode_err)},原始数据: {kinesis_data}")
                except json.JSONDecodeError as json_err:
                    logger.error(f"JSON解析失败: {str(json_err)},解码后字符串: {decoded_str}")
        return {"statusCode": 200, "message": "所有数据处理完成"}
    except Exception as e:
        logger.error(f"全局处理异常: {str(e)}")
        raise e  # 抛出异常触发Kinesis重试

3. 验证IAM权限

确保Lambda角色拥有DynamoDB批量写入权限,添加如下IAM策略:

{
    "Version": "2012-10-17",
    "Statement": [
        {
            "Effect": "Allow",
            "Action": "dynamodb:BatchWriteItem",
            "Resource": "arn:aws:dynamodb:ap-northeast-1:777050133147:table/dyno-table"
        }
    ]
}

4. 查看CloudWatch日志

登录AWS控制台,进入CloudWatch查看Lambda的执行日志,确认是否存在解码错误、权限错误等信息,这是排查生产环境问题的核心手段。

内容的提问来源于stack exchange,提问作者Deepak Singh Rajput

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.26 12:32:18