如何将带Protobuf Schema的Pub/Sub消息导出至GCS存储桶?
问题描述
我通过采用Protobuf Schema的Pub/Sub主题发布消息,此前消费者可借助该Schema正常读取并解码消息。现在我希望创建导出订阅,将这些消息读取后写入GCS存储桶的文件中,且文件需保留原Protobuf消息,以便后续通过Dataflow作业解码处理。
我已通过两种方式成功配置导出订阅并向GCS存储桶写入文件:
- 通过Cloud Console创建
Write to Cloud Storage订阅(文件格式为文本) - 通过
Pub/Sub to Text Files on Cloud Storage模板创建Dataflow作业
但两种方式生成的文件内容均无法正常解析,下载后使用protoc结合生成记录所用的Schema手动解码时,会收到Failed to parse input.错误提示。
我认为问题出在将Protobuf消息以文本文件形式写入的环节,请问能否在保持发布端Protobuf消息格式的前提下实现该需求?
问题核心是你之前的两种方式都将原始二进制Protobuf消息转成了文本格式(比如Base64编码或JSON序列化),导致protoc无法直接解析。要保留原始Protobuf格式,需要直接写入二进制数据,具体可通过以下两种方式实现:
1. 使用Pub/Sub内置导出订阅选择二进制格式
在Cloud Console创建导出订阅时,文件格式选择「Binary」而非「Text」:
- 进入Pub/Sub订阅创建页面,选择「Write to Cloud Storage」类型
- 在「File format」下拉菜单中选择「Binary」
- 完成其他配置后创建订阅
这种方式下,Pub/Sub会将消息的原始二进制数据直接写入GCS文件,无额外编码或转换。下载文件后可直接用protoc解码:
protoc --decode=YourMessageType your_schema.proto < gs://your-bucket/path/to/file.bin
2. 自定义Dataflow作业写入原始二进制数据
如果需要更灵活的控制(比如自定义文件命名、分区规则),可编写自定义Dataflow管道,直接读取Pub/Sub的原始字节消息并写入GCS二进制文件:
Python示例代码
import apache_beam as beam from apache_beam.options.pipeline_options import PipelineOptions, StandardOptions def run(): # 替换为你的项目和订阅信息 options = PipelineOptions() options.view_as(StandardOptions).streaming = True pubsub_subscription = "projects/your-project-id/subscriptions/your-subscription" gcs_output_path = "gs://your-bucket-name/output/prefix" with beam.Pipeline(options=options) as p: ( p # 读取Pub/Sub订阅的原始二进制消息 | "Read from Pub/Sub" >> beam.io.ReadFromPubSub(subscription=pubsub_subscription) # 写入GCS,禁用换行符避免破坏二进制数据 | "Write raw bytes to GCS" >> beam.io.WriteToText( file_path_prefix=gcs_output_path, file_name_suffix=".bin", append_trailing_newlines=False, shard_name_template="-SSS" # 可选:自定义分片命名规则 ) ) if __name__ == "__main__": run()
关键注意点
- 必须设置
append_trailing_newlines=False,否则每条消息末尾会被添加换行符,破坏Protobuf二进制结构 - 输出文件后缀建议用
.bin,明确标识为二进制文件 - 该管道为流式作业,会持续从Pub/Sub订阅拉取消息并写入GCS
验证方法
导出完成后,下载GCS中的文件,执行protoc解码命令,若能正常解析出消息内容,则说明格式保留正确。
内容的提问来源于stack exchange,提问作者Rohan Aletty

