Golang DataFlow从PubSub写入BigQuery时遇"no root units"错误求助
问题:Direct Runner运行Beam Go管道时出现"no root units"错误
我尝试通过DataFlow从PubSub读取消息并写入BigQuery表,但使用Direct Runner运行时遇到了**"no root units"**错误。
我的代码
package main import ( "context" "encoding/json" "flag" "fmt" "github.com/apache/beam/sdks/v2/go/pkg/beam/io/bigqueryio" "github.com/apache/beam/sdks/v2/go/pkg/beam/x/debug" "github.com/apache/beam/sdks/v2/go/pkg/beam" "github.com/apache/beam/sdks/v2/go/pkg/beam/io/pubsubio" "github.com/apache/beam/sdks/v2/go/pkg/beam/log" "github.com/apache/beam/sdks/v2/go/pkg/beam/x/beamx" ) type DummyBody struct { TaskId string `json:"id" bigquery:"id"` } func buildPipeline(s beam.Scope) { rawDummyBodies := pubsubio.Read(s, "project", "topic", &pubsubio.ReadOptions{Subscription: "sub.ID"}) dummyBodies := beam.ParDo(s, func(ctx context.Context, data []byte) (DummyBody, error) { var body DummyBody if err := json.Unmarshal(data, &body); err != nil { log.Error(ctx, err) fmt.Println("Error") return body, err } fmt.Println("No Error") return body, nil }, rawDummyBodies) debug.Printf(s, "Task : %#v", dummyBodies) bigqueryio.Write(s, "project", "table", dummyBodies) } func main() { flag.Parse() beam.Init() p, s := beam.NewPipelineWithRoot() buildPipeline(s) ctx := context.Background() if err := beamx.Run(ctx, p); err != nil { log.Exitf(ctx, "Failed to execute pipeline: %v", err) } }
报错信息
2022/11/01 14:29:55 Failed to execute pipeline: translation failed
caused by:
no root units
exit status 1
解决方案
这个错误主要是Direct Runner与PubSub订阅读取模式不兼容导致的,可按以下步骤调整:
修改PubSub读取配置,启用拉模式
给pubsubio.ReadOptions添加UseSubscriptionInPullMode: true,适配Direct Runner的本地运行逻辑:rawDummyBodies := pubsubio.Read(s, "project", "topic", &pubsubio.ReadOptions{ Subscription: "sub.ID", UseSubscriptionInPullMode: true, })移除debug.Printf临时测试
部分版本中debug.Printf可能干扰管道拓扑构建,先注释掉该行再测试:// debug.Printf(s, "Task : %#v", dummyBodies)校验BigQuery表结构匹配
确保BigQuery目标表的字段名、类型与DummyBody的bigquery标签完全一致,避免写入阶段的隐性错误触发管道翻译失败。升级Beam SDK到稳定版
若使用的SDK版本存在已知兼容性问题,升级到最新的v2稳定版可解决部分底层bug。
内容的提问来源于stack exchange,提问作者Ege Soyarar
相关产品推荐
相关产品推荐

