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

如何通过AWS Kinesis Firehose将Proto/JSON转换为Parquet

刚好做过类似的流处理转Parquet存S3的场景,给你梳理一套完整的落地步骤,完美适配你每分钟10000条实体的需求:

落地分步指南

第一步:先把Proto转成JSON(推荐路径)

虽然Firehose理论上能处理Proto,但用JSON当中间格式转Parquet会省超多事——不用额外写复杂的Schema映射逻辑。你的进程里直接用对应语言的Proto工具库转就行:

  • Java:用com.google.protobuf:protobuf-java-util里的JsonFormat类
  • Python:用protobuf库的json_format.MessageToJson()方法
  • Go:用github.com/golang/protobuf/jsonpb包

举个Python的极简示例,直接能用:

from google.protobuf import json_format
import your_proto_def  # 导入你自己写的Proto定义模块

# 假设你生成的控制信号实体是control_signal_obj
control_signal_obj = your_proto_def.ControlSignal(...)  # 填充你的实体数据
json_str = json_format.MessageToJson(control_signal_obj, preserving_proto_field_name=True)

转出来的JSON字段名和Proto完全对应,后续转Parquet的时候Schema对齐超轻松。

第二步:配置Kinesis Data Firehose投递流

这是核心环节,按以下步骤在AWS控制台操作就行:

  • 选数据源:直接选Direct PUT or other sources,因为你的进程要直接发数据到Firehose。
  • 开启格式转换:打开Record format conversion,目标格式选Apache Parquet。这里要注意:
    • Firehose依赖Glue Data Catalog管理Parquet的Schema,所以你要么提前在Glue里手动建表,要么让Firehose自动推断。强烈建议手动建表,避免自动推断的字段类型出错(比如把数字当成字符串)。比如在Glue里定义表的字段:参考点(float)、观测值(float)、控制信号(float)、生成时间(timestamp)等,和你的JSON结构一一对应。
    • 配置转换时,指定你在Glue里创建的数据库和表名,Firehose就会用这个Schema把JSON转成Parquet。
  • 配置S3目标:选你的目标S3桶,还可以设置路径前缀(比如year=!{timestamp:yyyy}/month=!{timestamp:MM}/day=!{timestamp:dd}/),按日期分目录,后续用Athena查询的时候效率更高。
  • 缓冲区设置:你每分钟10000条数据,把缓冲区大小设为1MB或者缓冲区时间设为60秒都行——既能保证数据及时写入,又不会生成一堆Parquet小文件(小文件太多会严重拖累后续分析性能)。

第三步:从你的进程发数据到Firehose

用AWS SDK批量发送,别一条一条发,不然API调用成本会很高。举个Python用boto3的示例:

import boto3
import json

# 初始化Firehose客户端,替换成你的AWS区域
firehose = boto3.client('firehose', region_name='us-east-1')

def batch_send_to_firehose(json_records):
    # 把多条JSON转成Firehose要求的换行分隔格式
    firehose_records = [{'Data': json.dumps(record) + '\n'} for record in json_records]
    resp = firehose.put_record_batch(
        DeliveryStreamName='your-firehose-stream-name',  # 替换成你的流名称
        Records=firehose_records
    )
    # 处理失败的记录,避免丢数
    if resp['FailedPutCount'] > 0:
        print(f"Warning: {resp['FailedPutCount']} records failed to send")

另外,一定要给你的进程所在的角色(比如ECS任务角色、EC2实例角色)加firehose:PutRecordBatch的权限,不然会报权限错误。

第四步:验证和后续分析

  • 数据写到S3后,直接用Athena查就行,因为Glue Data Catalog已经有对应的表结构了。比如写个SQL查最近7天的数据:
SELECT * FROM your_glue_db.your_control_signals_table 
WHERE generation_time >= current_date - interval '7' day;
  • 可以下载一个Parquet文件,用parquet-tools(本地装一个)查看结构,确认字段类型和数据都没问题。

一些实用优化建议

  • Schema更新要同步:如果后续你的Proto字段改了,一定要同步更新Glue里的表结构,不然Firehose转换会失败,数据会跑到错误目录里。
  • 定期查错误目录:Firehose会把转换失败的记录写到S3的error/前缀目录,每周抽5分钟看看,排查格式问题。
  • 监控关键指标:在CloudWatch里盯Firehose的DeliveryToS3Success和ConversionSuccess指标,一旦掉下去马上排查。

内容的提问来源于stack exchange,提问作者KunalC

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.22 08:18:10