求助:修复AWS Kinesis转S3的Lambda代码,实现单条JSON独立换行
解决Kinesis Stream转S3时JSON记录合并到同一行的问题
嘿,我看你遇到了Lambda处理Kinesis数据后,所有JSON记录都挤在S3同一行的问题,这在Firehose配合Lambda做数据转换的场景里挺常见的,我来帮你修复代码。
问题出在哪?
你的代码里用了json.dumps(..., indent=0),这个参数会让生成的JSON内部产生换行(虽然indent=0不会缩进,但每个键值对都会换行),而且更关键的是,你没有给每条处理后的记录添加换行分隔符,导致Firehose把所有记录的内容直接拼接在一起,最终写到S3里就变成了一大行。
我们需要两个小调整就能解决:
- 生成紧凑的单行JSON(别用
indent参数,让JSON挤成一行) - 给每条JSON记录末尾加个换行符
\n,确保Firehose写入时每条记录单独成行
修改后的完整代码
import base64 import json print('Loading function') def lambda_handler(event, context): print('input format json:', event['records']) output = [] for record in event['records']: payload = base64.b64decode(record['data']) payload_2 = base64.b64decode(payload) sector = json.loads(payload_2) # 生成紧凑的单行JSON,去掉indent参数 sector_dumps = json.dumps({ 'resourceType': sector['entry'][0]['resource']['resourceType'], 'status': sector['entry'][0]['resource']['status'], 'id': sector['entry'][0]['resource']['id'], 'code': sector['entry'][0]['resource']['code']['coding'], 'value': sector['entry'][0]['resource']['valueQuantity']['value'], 'unit': sector['entry'][0]['resource']['valueQuantity']['unit'], 'effectiveDateTime': sector['entry'][0]['resource']['effectiveDateTime'], 'subject': sector['entry'][0]['resource']['subject'], 'performer': sector['entry'][0]['resource']['performer'], 'category': sector['entry'][0]['resource']['category'], }) # 给每条JSON添加换行符,确保写入S3时单独成行 sector_dumps_with_newline = sector_dumps + '\n' sector_byte = sector_dumps_with_newline.encode("utf-8") sector_byte_encode = base64.b64encode(sector_byte) output_record = { 'recordId': record['recordId'], 'result': 'Ok', 'data': sector_byte_encode } output.append(output_record) print('Successfully processed {} records.'.format(len(event['records']))) return {'records': output}
关键修改说明
- 去掉
indent=0:json.dumps默认会生成没有多余空格和换行的紧凑JSON,这样单条记录的JSON本身就是完整的一行,不会内部拆分行。 - 添加
\n换行符:每条处理后的JSON末尾加个换行,Firehose在把多条记录拼接写入S3时,就会自动把每条记录分隔到单独的行,最终得到标准的JSON Lines格式文件,不管是用Athena查询还是其他工具读取都很方便。
内容的提问来源于stack exchange,提问作者Anand
相关产品推荐
相关产品推荐

