Kinesis Firehose经Lambda转换后报错,是否需重新编码数据?
问题原因与解决方案
错误原因
Kinesis Firehose对Lambda数据转换的输出格式有严格要求:每个返回的output_record中的data字段必须是Base64编码的字符串,而你的代码中直接将解码后的Python字典赋值给了data,这不符合格式规范,导致Firehose报错。
正确处理方式
Firehose无法直接传递结构化的字典数据,必须将处理后的数据重新编码为Base64格式返回。你需要在Lambda中完成解码、解析操作后,再把数据重新序列化为JSON字符串并做Base64编码,之后返回给Firehose,最后由Fluentd接收后解码解析。
修正后的Lambda代码
import base64 import json def lambda_handler(event, context): output = [] for record in event['records']: # 解码原始Base64数据并转为字符串 payload = base64.b64decode(record['data']).decode('utf-8') # 解析为字典(此处可添加自定义解析逻辑) data_dict = json.loads(payload) # 将处理后的字典重新转为JSON字符串,再做Base64编码 processed_data = json.dumps(data_dict).encode('utf-8') encoded_data = base64.b64encode(processed_data).decode('utf-8') # 构造符合Firehose要求的输出记录 output_record = { 'recordId': record['recordId'], 'result': 'Ok', 'data': encoded_data } output.append(output_record) print(f'Successfully processed {len(event["records"])} records.') return {'records': output}
补充说明
- 若需要对解析后的
data_dict做自定义处理(比如提取字段、过滤数据等),直接在解析步骤之后添加对应逻辑即可,不影响后续的重新编码流程。 - Fluentd端收到的仍是Base64编码数据,需在Fluentd配置中添加解码逻辑(如使用
base64_decode过滤器)还原为JSON格式,再进行后续解析。
内容的提问来源于stack exchange,提问作者fledgling
相关产品推荐
相关产品推荐

