如何在Apache Beam的ReadFromText中读取多GCS路径/多桶文件?
解决Apache Beam ReadFromText读取多路径(跨GCS桶)的问题
直接给你几个可行的解决方案,都是实际项目里验证过的:
直接传入路径列表
ReadFromText本身就支持接收字符串列表作为输入路径,不管是不是跨GCS桶,只要把每个带通配符的路径单独列出来就行。示例代码:from apache_beam.io import ReadFromText with Pipeline(options=pipeline_options) as p: lines = p | ReadFromText([ 'gs://bucket1/folder/*', 'gs://bucket2/folder/*', 'gs://bucket3/another/path/*.txt' # 还能混合不同的通配符规则 ])这种方式最稳妥,Beam会自动处理所有匹配的文件,不用依赖shell的语法扩展。
避免shell风格的大括号语法
你之前用的gs://{bucket1/folder,bucket2/folder}/*是bash/zsh这类shell的扩展语法,Dataflow的命令行解析器不认识这种写法,所以才会报「unrecognized arguments」错误。别再用这种shell特有的写法,换成Beam支持的列表形式。通过命令行参数传递多路径(弹性模板场景)
如果是用弹性模板,需要把多路径作为参数传入的话,别用大括号,而是用逗号分隔的字符串,然后在代码里拆分:- 启动模板时的命令行参数:
gcloud dataflow flex-template run ... --parameters input_paths="gs://bucket1/folder/*,gs://bucket2/folder/*" - 代码里解析参数:
from apache_beam.options.pipeline_options import PipelineOptions class MyOptions(PipelineOptions): @classmethod def _add_argparse_args(cls, parser): parser.add_argument('--input_paths', type=str, required=True) pipeline_options = PipelineOptions() custom_options = pipeline_options.view_as(MyOptions) input_paths_list = custom_options.input_paths.split(',') with Pipeline(options=pipeline_options) as p: lines = p | ReadFromText(input_paths_list)
要是路径里有特殊字符(虽然GCS路径一般不会有),可以换个分隔符比如分号,只要前后统一就行。
- 启动模板时的命令行参数:
总结一下:Beam的ReadFromText对多路径的支持是直接通过列表实现的,不用搞shell那套扩展语法,按上面的方式来就能解决跨桶读取的问题。
内容的提问来源于stack exchange,提问作者salvob
相关产品推荐
相关产品推荐

