GCS触发Cloud Function启动DataFlow时CSV输入路径配置问题
解决方案
你不需要使用TextIO.read().metadata的写法,只需要把输入文件路径调整为DataFlow的运行时参数,再在Cloud Function触发时传入触发文件的GCS路径即可,具体修改步骤如下:
1. 调整DataFlow管道代码,添加运行时参数
首先在管道代码中新增对输入路径参数的解析,替换原来写死的GCS路径:
import apache_beam as beam from apache_beam.options.pipeline_options import PipelineOptions # 自定义参数类,用于接收运行时传入的配置 class MyPipelineOptions(PipelineOptions): @classmethod def _add_argparse_args(cls, parser): parser.add_argument( "--input_file", required=True, help="输入CSV文件的GCS路径" ) # 主管道逻辑修改 p = beam.Pipeline(options=options) pipeline_options = options.view_as(MyPipelineOptions) raw_values = ( p | "ReadTable" >> beam.io.textio.ReadFromText(pipeline_options.input_file) | "custFunc" >> beam.Map(CallAPI) | "writeTable" >> beam.io.WriteToBigQuery('newtable', project='project1', dataset='test', schema=table_schema, write_disposition=beam.io.BigQueryDisposition.WRITE_TRUNCATE, create_disposition=beam.io.BigQueryDisposition.CREATE_IF_NEEDED) )
修改完成后将管道打包为DataFlow经典模板或者Flex模板,上传到GCS存储桶中即可。
2. 调整Cloud Function触发逻辑
在GCS触发的Cloud Function代码中,拿到触发事件对应的文件路径,拼接为gs://[存储桶名]/[文件路径]的格式,作为input_file参数传入DataFlow作业的启动请求中即可。
注意:如果同一时间有多个文件上传触发多次作业,要注意BigQuery的
WRITE_TRUNCATE配置会导致后续作业覆盖前一次的写入结果,如有并发写入需求可调整写入策略。
内容的提问来源于stack exchange,提问作者GK89
相关产品推荐
相关产品推荐

