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

MSK Kafka触发的Python Lambda写入S3遇KeyError:'Records'问题求助

MSK Kafka触发Lambda写入S3时出现KeyError: 'Records'问题解决

我创建了一个由MSK Kafka触发的Python Lambda函数,用于将数据写入S3存储桶,但运行时出现KeyError错误,错误信息为"'Records'"。

原代码

import json
import boto3
import base64
import uuid

s3 = boto3.client('s3')

def lambda_handler(event, context):
    for record in event['Records']:
        # Decode the base64-encoded Kafka data
        kinesis_data = base64.b64decode(record['kinesis']['data']).decode('utf-8')
        payload = json.loads(kinesis_data)

        # Generate a unique key for each object
        unique_key = f"{payload['event_type']}/{str(uuid.uuid4())}.json"
        
        # Process the payload as needed and write it to S3
        s3.put_object(
            Bucket='dev-mos-xcorr-broadcast',
            Key=unique_key,
            Body=json.dumps(payload)
        )

    return {'statusCode': 200, 'body': 'Data written to S3'}

错误详情

{
"errorMessage": "'Records'",
"errorType": "KeyError",
"requestId": "60be31e6-ce4a-4ee7-a3ab-f8aff3f49717",
"stackTrace": [
"  File \"/var/task/lambda_function.py\", line 9, in lambda_handler\n    for record in event['Records']:\n"
]
}

解决方案

错误原因

你的代码采用了Kinesis触发Lambda的事件格式逻辑,但MSK Kafka触发的Lambda事件结构完全不同:

  • 没有大写的Records字段,取而代之的是小写的records
  • 事件按topic分组,每个topic下包含多条消息,不存在kinesis嵌套字段

修改后的代码

import json
import boto3
import base64
import uuid

s3 = boto3.client('s3')

def lambda_handler(event, context):
    # 遍历所有topic下的消息
    for topic, messages in event['records'].items():
        for msg in messages:
            try:
                # 解码Kafka消息的value内容
                payload = json.loads(base64.b64decode(msg['value']).decode('utf-8'))
                # 生成唯一S3对象键,兼容缺失event_type的情况
                unique_key = f"{payload.get('event_type', 'unknown')}/{str(uuid.uuid4())}.json"
                # 写入S3存储桶
                s3.put_object(
                    Bucket='dev-mos-xcorr-broadcast',
                    Key=unique_key,
                    Body=json.dumps(payload)
                )
            except Exception as e:
                print(f"处理消息失败: {str(e)}")
                # 可选:将失败消息写入死信队列或留存日志

    return {'statusCode': 200, 'body': '数据已写入S3'}

额外提示

  • 测试时可以在函数开头添加print(json.dumps(event)),查看完整的MSK触发事件结构,确认字段匹配
  • 加入异常处理逻辑,避免单条消息处理失败导致整个函数报错退出
  • 使用payload.get()替代直接索引字段,防止payload结构异常引发新的KeyError

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.12 05:12:33