AWS Kinesis Firehose写入S3的JSON格式与字段特殊字符处理问题
解决方案
核心修改逻辑
- 新增键名清理工具函数,支持递归处理嵌套JSON结构,可移除键名中所有非字母数字字符,完全匹配你给出的
RelatedAWSResources:0/type转换为RelatedAWSResources0type的需求 - 修复原代码全局
output变量的隐患:Lambda容器复用场景下全局变量会累积历史处理记录,造成数据重复,改为在函数内部初始化结果列表 - 新增JSON解析异常捕获,避免非法格式的 payload 导致整个Lambda执行中断
- 原有加换行的逻辑保留,同时兼容键名转换需求
完整修改后的Lambda代码
import json import boto3 import base64 import re def clean_key_name(key): # 移除非字母数字的所有特殊字符,可按需调整正则匹配规则 return re.sub(r'[^a-zA-Z0-9]', '', key) def clean_json_keys(obj): # 递归处理JSON中所有层级的键名 if isinstance(obj, dict): new_dict = {} for k, v in obj.items(): new_key = clean_key_name(k) new_dict[new_key] = clean_json_keys(v) return new_dict elif isinstance(obj, list): return [clean_json_keys(item) for item in obj] else: return obj def lambda_handler(event, context): output = [] for record in event['records']: payload = base64.b64decode(record['data']).decode('utf-8') try: # 解析JSON并清理键名 raw_data = json.loads(payload) cleaned_data = clean_json_keys(raw_data) cleaned_payload = json.dumps(cleaned_data) # 加换行符 row_w_newline = cleaned_payload + "\n" result = "Ok" except json.JSONDecodeError: # 非法JSON可标记为处理失败,Firehose会将其投递到配置的错误前缀目录 row_w_newline = payload + "\n" result = "ProcessingFailed" encoded_data = base64.b64encode(row_w_newline.encode('utf-8')) output_record = { 'recordId': record['recordId'], 'result': result, 'data': encoded_data } output.append(output_record) print('Processed {} records.'.format(len(event['records']))) return {'records': output}
注意事项
- 键名清理规则可灵活调整:如果需要保留下划线等合法字符,只需修改
clean_key_name函数中的正则表达式为r'[^a-zA-Z0-9_]'即可 - 异常处理逻辑可根据业务需求调整:如果不需要拦截非法JSON,可将
ProcessingFailed改为Ok,非法记录会保留原始内容加换行后投递到S3 - 修改完成后可先在Lambda控制台用测试事件验证转换效果,确认键名转换、换行添加符合预期后再绑定到Firehose上线
内容的提问来源于stack exchange,提问作者Bokambo
相关产品推荐
相关产品推荐

