Apache Beam使用DataflowRunner时无法写入BigQuery且遗留大量临时文件问题求助
Apache Beam使用DataflowRunner时无法写入BigQuery且遗留大量临时文件问题求助
我目前在开发一个Apache Beam管道,用于处理数据并写入BigQuery。使用DirectRunner时管道运行完全正常,但切换到DataflowRunner后,管道能无错误无警告地完成,却没有任何数据插入到BigQuery中。另外,我发现Cloud Storage桶的临时目录(gs://my-bucket/temp/bq_load/...)里留下了大量的临时文件,目标表中也完全没有数据。
以下是我的管道结构:
worker_options.sdk_container_image = '...' with beam.Pipeline(options=pipeline_options) as p: processed_data = ( p | "ReadFiles" >> beam.Create(FILE_LIST) | "ProcessFiles" >> beam.ParDo(ProcessAvroFileDoFn()) | "WriteToBigQuery" >> beam.io.WriteToBigQuery( table=f"{PROJECT_ID}:{DATASET_ID}.{TABLE_ID}", schema=BQ_SCHEMA, write_disposition=beam.io.BigQueryDisposition.WRITE_APPEND, create_disposition=beam.io.BigQueryDisposition.CREATE_IF_NEEDED ) )
关键观察
- 使用DirectRunner时管道成功运行,数据能正常写入BigQuery,没有任何问题。
- 使用DataflowRunner时,管道无错误无警告完成,但:
- 没有数据写入BigQuery
- 桶中遗留大量临时文件(比如
bq_load/...路径下的文件)
- 处理的数据是有效的NDJSON格式
- BigQuery的Schema和数据结构完全匹配
已尝试的操作
- 检查遗留的临时文件:下载了临时文件并验证内容是有效的NDJSON行,手动用
bq load命令上传这个文件到BigQuery完全可行。 - 测试其他数据集:尝试了多种不同的输入,但问题依然存在。
- 查看Dataflow日志:在Dataflow监控控制台里检查了日志,但没有发现任何错误或警告。
- 更换服务账号:使用权限不足的服务账号时管道会抛出错误,所以看起来不太可能是Worker的权限问题。
我之前看到过一个类似的讨论线程,但里面并没有得到有效的解决方法,希望在这里能找到答案。
备注:内容来源于stack exchange,提问作者Spine Feast
相关产品推荐
相关产品推荐

