登录
首页 >  Golang >  Go问答

Apache Beam Go SDK 实现数据流传输

来源:stackoverflow

时间:2024-03-10 18:15:27 104浏览 收藏

对于一个Golang开发者来说,牢固扎实的基础是十分重要的,golang学习网就来带大家一点点的掌握基础知识点。今天本篇文章带大家了解《Apache Beam Go SDK 实现数据流传输》,主要介绍了,希望对大家的知识积累有所帮助,快点收藏起来吧,否则需要时就找不到了!

问题内容

我一直在使用 go beam sdk (v2.13.0),但无法获取在 gcp dataflow 上运行的字数统计示例。它进入崩溃循环,尝试启动 org.apache.beam.runners.dataflow.worker.dataflowrunnerharness。使用 direct 运行程序在本地运行时,该示例可以正确执行。

该示例与上面给出的原始示例完全没有修改。

堆栈跟踪是:

org.apache.beam.vendor.grpc.v1p13p1.com.google.protobuf.invalidprotocolbufferexception: protocol message had invalid utf-8. 
at org.apache.beam.vendor.grpc.v1p13p1.com.google.protobuf.invalidprotocolbufferexception.invalidutf8(invalidprotocolbufferexception.java:148) 
at org.apache.beam.vendor.grpc.v1p13p1.com.google.protobuf.codedinputstream$streamdecoder.readstringrequireutf8(codedinputstream.java:2353) 
at org.apache.beam.model.pipeline.v1.runnerapi$functionspec.(runnerapi.java:59611) 
at org.apache.beam.model.pipeline.v1.runnerapi$functionspec.(runnerapi.java:59572) 
at org.apache.beam.model.pipeline.v1.runnerapi$functionspec$1.parsepartialfrom(runnerapi.java:60241) 
at org.apache.beam.model.pipeline.v1.runnerapi$functionspec$1.parsepartialfrom(runnerapi.java:60235) 
at org.apache.beam.vendor.grpc.v1p13p1.com.google.protobuf.codedinputstream$streamdecoder.readmessage(codedinputstream.java:2424) 
at org.apache.beam.model.pipeline.v1.runnerapi$coder.(runnerapi.java:27531) 
at org.apache.beam.model.pipeline.v1.runnerapi$coder.(runnerapi.java:27489) 
at org.apache.beam.model.pipeline.v1.runnerapi$coder$1.parsepartialfrom(runnerapi.java:28410) 
at org.apache.beam.model.pipeline.v1.runnerapi$coder$1.parsepartialfrom(runnerapi.java:28404) 
at org.apache.beam.model.pipeline.v1.runnerapi$coder$builder.mergefrom(runnerapi.java:28028) 
at org.apache.beam.model.pipeline.v1.runnerapi$coder$builder.mergefrom(runnerapi.java:27868) 
at org.apache.beam.vendor.grpc.v1p13p1.com.google.protobuf.codedinputstream$streamdecoder.readmessage(codedinputstream.java:2408) 
at org.apache.beam.vendor.grpc.v1p13p1.com.google.protobuf.mapentrylite.parsefield(mapentrylite.java:128) 
at org.apache.beam.vendor.grpc.v1p13p1.com.google.protobuf.mapentrylite.parseentry(mapentrylite.java:184) 
at org.apache.beam.vendor.grpc.v1p13p1.com.google.protobuf.mapentry.(mapentry.java:106) 
at org.apache.beam.vendor.grpc.v1p13p1.com.google.protobuf.mapentry.(mapentry.java:50) 
at org.apache.beam.vendor.grpc.v1p13p1.com.google.protobuf.mapentry$metadata$1.parsepartialfrom(mapentry.java:70) 
at org.apache.beam.vendor.grpc.v1p13p1.com.google.protobuf.mapentry$metadata$1.parsepartialfrom(mapentry.java:64) 
at org.apache.beam.vendor.grpc.v1p13p1.com.google.protobuf.codedinputstream$streamdecoder.readmessage(codedinputstream.java:2424) 
at org.apache.beam.model.pipeline.v1.runnerapi$components.(runnerapi.java:930) 
at org.apache.beam.model.pipeline.v1.runnerapi$components.(runnerapi.java:848) 
at org.apache.beam.model.pipeline.v1.runnerapi$components$1.parsepartialfrom(runnerapi.java:2714) 
at org.apache.beam.model.pipeline.v1.runnerapi$components$1.parsepartialfrom(runnerapi.java:2708) 
at org.apache.beam.vendor.grpc.v1p13p1.com.google.protobuf.codedinputstream$streamdecoder.readmessage(codedinputstream.java:2424) 
at org.apache.beam.model.pipeline.v1.runnerapi$pipeline.(runnerapi.java:2892) 
at org.apache.beam.model.pipeline.v1.runnerapi$pipeline.(runnerapi.java:2850) 
at org.apache.beam.model.pipeline.v1.runnerapi$pipeline$1.parsepartialfrom(runnerapi.java:3981) 
at org.apache.beam.model.pipeline.v1.runnerapi$pipeline$1.parsepartialfrom(runnerapi.java:3975) 
at org.apache.beam.vendor.grpc.v1p13p1.com.google.protobuf.abstractparser.parsepartialfrom(abstractparser.java:221) 
at org.apache.beam.vendor.grpc.v1p13p1.com.google.protobuf.abstractparser.parsefrom(abstractparser.java:239) 
at org.apache.beam.vendor.grpc.v1p13p1.com.google.protobuf.abstractparser.parsefrom(abstractparser.java:244) 
at org.apache.beam.vendor.grpc.v1p13p1.com.google.protobuf.abstractparser.parsefrom(abstractparser.java:49) 
at org.apache.beam.vendor.grpc.v1p13p1.com.google.protobuf.generatedmessagev3.parsewithioexception(generatedmessagev3.java:311) 
at org.apache.beam.model.pipeline.v1.runnerapi$pipeline.parsefrom(runnerapi.java:3222) 
at org.apache.beam.runners.dataflow.worker.dataflowworkerharnesshelper.getpipelinefromenv(dataflowworkerharnesshelper.java:131) 
at org.apache.beam.runners.dataflow.worker.dataflowrunnerharness.main(dataflowrunnerharness.java:59)

我使用了示例中指定的 docker 映像,并且还使用相同的标签 (v2.13.0) 从我自己的 docker 进行了尝试,但仍然遇到相同的错误。我意识到它还没有准备好投入生产,但我希望示例能够正常工作。

按照开始时的说明,我像这样运行了这项工作:

wordcount --input gs://dataflow-samples/shakespeare/kinglear.txt \
--output gs://example-bucket/counts \
--runner dataflow \
--project example-project \
--temp_location gs://example-bucket/tmp/ \
--staging_location gs://example-bucket/binaries/ \
--worker_harness_container_image=apache-docker-beam-snapshots-docker.bintray.io/beam/go:20180515

我再次尝试了入门中提供的 docker,以及使用 v2.13.0 构建的 docker。

我的示例文件 go.mod 是:

module example.org/wordcount

go 1.12

require (
    cloud.google.com/go v0.41.0 // indirect
    github.com/apache/beam v2.13.0+incompatible
    github.com/pkg/errors v0.8.1 // indirect
    golang.org/x/net v0.0.0-20190628185345-da137c7871d7 // indirect
    google.golang.org/grpc v1.22.0 // indirect
)

这可能是什么原因造成的?


解决方案


Dataflow 并未正式支持 Apache Beam Go SDK。不过,一些用户已经能够使用它。我怀疑这个版本可能有问题。您也许可以尝试不同的版本。

您可以在 Beam mailing list 上与其他用户讨论哪些版本适合他们(但不受支持)。

今天关于《Apache Beam Go SDK 实现数据流传输》的内容就介绍到这里了,是不是学起来一目了然!想要了解更多关于的内容请关注golang学习网公众号!

声明:本文转载于:stackoverflow 如有侵犯,请联系study_golang@163.com删除
相关阅读
更多>
最新阅读
更多>
课程推荐
更多>