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

如何将带Protobuf Schema的Pub/Sub消息导出至GCS存储桶?

如何将Pub/Sub的Protobuf消息以原始格式导出到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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.27 13:10:04