如何通过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
相关产品推荐
相关产品推荐

