使用Golang Apache Beam BigQueryIO写入BigQuery分区表的问题求助
问题背景
我是Apache Beam新手,尝试构建简单Pipeline,通过bigqueryio包将PCollection写入BigQuery分区表。由于Golang&Beam SDK的文档和示例多为Java/Python版本,入门难度较大。我使用github.com/apache/beam/sdks/v2/go/pkg/beam/io/xlang/bigqueryio包编写了代码:
col := createCol(s, line) outTable := "project:dataset:table$20220808" bigqueryio.Write(s, outTable, col, bigqueryio.CreateDisposition(bigqueryio.CreateIfNeeded), bigqueryio.WriteExpansionAddr(""))
运行时出现以下错误:
tried cross-language for beam:transform:org.apache.beam:schemaio_bigquery_write:v1 against autojava::sdks:java:io:google-cloud-platform:expansion-service:runExpansionService and failed expanding external transform error in starting expansion service, StartService(): context deadline exceeded
疑问与解答
本地运行时扩展服务问题
本地自动启动Java扩展服务时常因环境依赖(如Java版本不兼容、依赖包下载超时)触发超时错误,必须手动启动扩展服务:- 启动命令示例:
java -jar beam-sdks-java-io-google-cloud-platform-expansion-service-2.x.x.jar --port 8097 - 代码中修改
WriteExpansionAddr参数为手动启动的服务地址,比如localhost:8097
- 启动命令示例:
两个bigqueryio包的选择逻辑
github.com/apache/beam/sdks/v2/go/pkg/beam/io/xlang/bigqueryio:跨语言实现,依赖Java扩展服务,支持更丰富的写入配置,但需要处理跨语言服务的稳定性问题,适合复杂写入场景github.com/apache/beam/sdks/v2/go/pkg/beam/io/bigqueryio:Go原生实现,无需依赖Java服务,稳定性更高,但配置选项有限。如果仅需基础写入功能,优先选原生包;若需要自定义写入策略(如分区覆盖),则需使用跨语言包并解决扩展服务问题
WRITE_TRUNCATE及WriteDisposition缺失的解决办法
跨语言包当前暴露的配置不全,但可通过WithAdditionalParameters传递底层Java组件支持的参数:bigqueryio.Write(s, outTable, col, bigqueryio.CreateDisposition(bigqueryio.CreateIfNeeded), bigqueryio.WriteExpansionAddr("localhost:8097"), bigqueryio.WithAdditionalParameters(map[string]string{ "writeDisposition": "WRITE_TRUNCATE", }), )
该参数会直接传递给Java端的BigQuery写入组件,实现覆盖现有分区的效果
内容的提问来源于stack exchange,提问作者Avishay28

