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

