Lambda通过Firehose向S3写入空记录问题排查
问题:Kinesis Firehose结合Lambda写入S3后文件仅为空{}的原因?
问题背景
测试事件内容:
{ "TICKER_SYMBOL": "QXZ", "SECTOR": "HEALTHCARE", "CHANGE": -0.05, "PRICE": 84.51 }
Lambda代码:
import json import base64 def lambda_handler(event, context): print(event) for record in event['records']: #Kinesis data is base64 encoded so decode here payload=base64.b64decode(record["data"]) print("Decoded payload: " + str(payload)) json_object = {} output = [] output_record = { 'recordId': record['recordId'], 'result': 'Ok', 'data': base64.b64encode(json.dumps(json_object).encode('utf-8')).decode('utf-8')} output.append(output_record) print(output) return {'records': output}
代码能打印预期输出,Firehose显示写入成功,但S3下载的文件仅为空{}。
原因分析
核心问题出在以下两点:
- 未使用解码后的payload数据:你初始化了空字典
json_object = {},完全没有把解码后的payload内容解析并赋值给它,最终输出的是空对象的base64编码,所以S3文件里只有{}。 - output列表初始化位置错误:
output = []放在for循环内部,每次循环都会重置列表,虽然单条记录时不影响返回结果,但多记录场景下会导致只保留最后一条记录的输出,逻辑存在隐患。
修复后的代码
import json import base64 def lambda_handler(event, context): print(event) output = [] # 移到循环外部,避免重复初始化 for record in event['records']: # 解码Kinesis传递的base64数据 payload = base64.b64decode(record["data"]) print("Decoded payload: " + str(payload)) # 将解码后的字节数据转为JSON对象 json_object = json.loads(payload) # 构造输出记录,使用解析后的json_object output_record = { 'recordId': record['recordId'], 'result': 'Ok', 'data': base64.b64encode(json.dumps(json_object).encode('utf-8')).decode('utf-8') } output.append(output_record) print(output) return {'records': output}
关键修改说明
- 将
output = []移至for循环外,确保所有处理后的记录都能被收集。 - 通过
json.loads(payload)将解码后的字节数据解析为JSON对象,替换原有的空字典json_object = {}。 - 输出时使用解析后的
json_object进行编码,确保数据被正确传递到S3。
内容的提问来源于stack exchange,提问作者south153
相关产品推荐
相关产品推荐

