如何让自定义Protobuf对象支持Pickle序列化?
问题描述
我有一个名为HistogramBins.proto的Protobuf文件,内容如下:
syntax = "proto3"; message HistogramBins { string field_name = 1; repeated float bins = 2; }
通过以下命令编译:
protoc --python_out=. ./HistogramBins.proto
编译后生成HistogramBins_pb2.py文件。尝试在Apache Beam流处理中使用该类序列化时,代码如下:
import apache_beam as beam from components.HistogramBins_pb2 import HistogramBins from apache_beam.options.pipeline_options import PipelineOptions class ToProtoFn(beam.DoFn): def process(self, t): hBin = HistogramBins() hBin.field_name = t[0] hBin.bins.extend(t[1]) print(hBin) yield hBin with beam.Pipeline(options=PipelineOptions()) as p: input_collection = ( p | 'Read input data' >> beam.Create([("f1", [1.0, 2.0]), ("f2", [3.0, 4.0])]) | 'record to HistogramBins' >> beam.ParDo(ToProtoFn()) | beam.io.WriteToText('data/test2.pbtxt', coder=beam.coders.ProtoCoder(HistogramBins().__class__)) )
运行时出现错误:
_pickle.PicklingError: Can't pickle <class 'HistogramBins_pb2.HistogramBins'>: it's not found as HistogramBins_pb2.HistogramBins
而使用Google官方包中的Timestamp Protobuf类却能正常运行,代码如下:
import apache_beam as beam from google.protobuf.timestamp_pb2 import Timestamp from apache_beam.options.pipeline_options import PipelineOptions class ToProtoFn(beam.DoFn): def process(self, element): timestamp = Timestamp() timestamp.seconds, timestamp.nanos = [int(x) for x in element.strip().split(',')] print(timestamp) yield timestamp with beam.Pipeline(options=PipelineOptions()) as p: lines = (p | beam.Create(["1586753000,222333000", "1586754000,222333000"]) | beam.ParDo(ToProtoFn()) | beam.io.WriteToText('time-pb', coder=beam.coders.ProtoCoder(Timestamp().__class__)))
Google官方的Timestamp Protobuf文件内容如下:
syntax = "proto3"; package google.protobuf; option csharp_namespace = "Google.Protobuf.WellKnownTypes"; option cc_enable_arenas = true; option go_package = "github.com/golang/protobuf/ptypes/timestamp"; option java_package = "com.google.protobuf"; option java_outer_classname = "TimestampProto"; option java_multiple_files = true; option objc_class_prefix = "GPB"; message Timestamp { int64 seconds = 1; int32 nanos = 2; }
请问:
- 自定义Protobuf与官方Protobuf之间存在什么差异?
- 如何解决自定义Protobuf对象无法被Pickle序列化的问题?
解答
一、自定义与官方Protobuf的核心差异
- 模块可见性与路径标识:官方Timestamp类属于全局可识别的
google.protobuf标准包,其完全限定类路径能被Python解释器在任何进程中正确解析;而自定义的HistogramBins类属于本地模块,当Beam进行分布式调度时,Worker节点可能无法识别相对导入路径,导致pickle找不到类定义。 - 代码生成优化:官方Protobuf生成的Python代码经过针对性优化,确保类的元数据、模块路径等符合pickle序列化要求;而
protoc默认生成的自定义Protobuf代码,在模块路径标识上存在缺陷,尤其是模块不在Pythonsys.path根目录时。 - 命名空间声明:官方proto文件包含明确的
package声明与多语言命名空间配置,这些配置会转化为Python代码中的规范模块结构;自定义proto无package声明,生成的模块缺乏清晰命名空间,易在分布式环境中出现路径问题。
二、解决自定义Protobuf的Pickle序列化问题
1. 修正模块导入与分布式可见性
确保components目录是合法Python包(添加__init__.py),并将其所在路径加入Pythonsys.path。通过Beam的setup_file参数确保Worker节点能加载自定义模块:
- 创建
setup.py文件:
from setuptools import setup, find_packages setup( name='components', packages=find_packages(), )
- 运行Pipeline时指定参数:
--setup_file=setup.py
2. 直接传入Proto类而非实例的__class__
使用ProtoCoder时,直接传入类本身,避免通过实例获取类路径带来的问题:
# 替换原coder参数 coder=beam.coders.ProtoCoder(HistogramBins)
3. 给自定义proto添加package声明
修改HistogramBins.proto添加命名空间:
syntax = "proto3"; package myproject.components; message HistogramBins { string field_name = 1; repeated float bins = 2; }
重新编译后,生成的Python类模块路径更清晰,有助于pickle识别类位置。
4. 手动序列化/反序列化Protobuf对象
将Protobuf对象序列化为字节串后写入文件,读取时再反序列化:
# 修改DoFn输出字节串 class ToProtoFn(beam.DoFn): def process(self, t): hBin = HistogramBins() hBin.field_name = t[0] hBin.bins.extend(t[1]) yield hBin.SerializeToString() # 写入时无需指定ProtoCoder | beam.io.WriteToText('data/test2.pbtxt')
读取时通过HistogramBins().FromString(字节串)完成反序列化。
内容的提问来源于stack exchange,提问作者Mehran
相关产品推荐
相关产品推荐

