У меня есть файл protobuf с именем HistogramBins.proto, вот так:
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 : it's not found as HistogramBins_pb2.HistogramBins
При этом следующий пример работает нормально. На этот раз я использую protobuf из пакета Google:
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__)))
В чем разница между моим protobuf и их?
Это файл protobuf Google:
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;
}