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

求助:修复AWS Kinesis转S3的Lambda代码,实现单条JSON独立换行

解决Kinesis Stream转S3时JSON记录合并到同一行的问题

嘿,我看你遇到了Lambda处理Kinesis数据后,所有JSON记录都挤在S3同一行的问题,这在Firehose配合Lambda做数据转换的场景里挺常见的,我来帮你修复代码。

问题出在哪?

你的代码里用了json.dumps(..., indent=0),这个参数会让生成的JSON内部产生换行(虽然indent=0不会缩进,但每个键值对都会换行),而且更关键的是,你没有给每条处理后的记录添加换行分隔符,导致Firehose把所有记录的内容直接拼接在一起,最终写到S3里就变成了一大行。

我们需要两个小调整就能解决:

  1. 生成紧凑的单行JSON(别用indent参数,让JSON挤成一行)
  2. 给每条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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.28 03:27:45