如何在DataFlow中读取日期范围内的多通配符路径到PCollection
问题:DataFlow读取指定日期范围内的事件文件
场景描述
存储服务中按日期文件夹存放大量JSON格式的事件小文件,结构如下:
2022-01-01/file1.json 2022-01-01/file2.json 2022-01-01/file3.json 2022-01-01/file4.json 2022-01-01/file5.json 2022-01-02/file6.json 2022-01-02/file7.json 2022-01-02/file8.json 2022-01-03/file9.json 2022-01-03/file10.json
需求:DataFlow作业接收起止日期作为输入,读取该日期范围内的所有文件。
尝试的方案及疑问
方案1:结合beam.Create与ReadAllFromText(带通配符)
参考资料后,尝试将带通配符的日期路径传入beam.Create,再用ReadAllFromText读取,代码如下:
def run(argv=None): # argument parser # pipeline options, google_cloud_options file_list = ['gs://bucket_1/2022-01-02/*.json', 'gs://2022-01-03/*.json'] p = beam.Pipeline(options=pipeline_options) p1 = p | "create PCol from list" >> beam.Create(file_list) \ | "read files" >> ReadAllFromText() \ | "transform" >> beam.Map(lambda x: x) \ | "write to GCS" >> WriteToText('gs://bucket_3/output') result = p.run() result.wait_until_finish()
疑问:不确定beam.Create的列表中是否支持通配符,该方向是否正确。
方案2:循环读取每个日期路径
修改代码为循环读取每个带通配符的路径:
with beam.Pipeline() as p: file_list = ['gs://ext-pub-testjeff.appspot.com/2022-01-02/*.json', 'gs://ext-pub-testjeff.appspot.com/2022-01-03/*.json'] for i, file in enumerate(file_list): p = (p | f"Read Text {i}" >> beam.io.textio.ReadFromText(file, skip_header_lines = 0)) p = (p | "write to GCS" >> WriteToText('gs://ext-pub-testjeff.appspot.com/output'))
方案3:循环读取并分别写入
再次修改为读取一个日期路径就写入一次:
with beam.Pipeline() as p: file_list = ['gs://ext-pub-testjeff.appspot.com/2022-01-02/*.json', 'gs://ext-pub-testjeff.appspot.com/2022-01-03/*.json'] for i, file in enumerate(file_list): p = (p | f"Read Text {i}" >> beam.io.textio.ReadFromText(file, skip_header_lines = 0) | f"write to GCS {i}" >> WriteToText('gs://ext-pub-testjeff.appspot.com/output'))
方案验证与最优实现
各方案问题分析
- 方案1不可行:
ReadAllFromText仅接收具体文件路径,不会解析通配符,因此传入*.json这类通配符会导致找不到文件,直接报错。 - 方案2逻辑错误:循环中不断给
p赋值,第一次循环后p从Pipeline对象变为PCollection,第二次循环尝试在PCollection上执行ReadFromText(源操作),会直接抛出异常。 - 方案3不符合需求:每个日期路径单独读取写入,会生成多份分散的输出文件,无法实现统一处理,且存在资源浪费。
最优实现方式
直接利用ReadFromText支持多通配符路径列表的特性,结合日期范围生成路径模板,代码如下:
from datetime import datetime, timedelta import apache_beam as beam from apache_beam.options.pipeline_options import PipelineOptions def generate_date_range_paths(start_date_str, end_date_str, bucket): """根据起止日期生成所有日期对应的文件路径模板""" start_date = datetime.strptime(start_date_str, "%Y-%m-%d") end_date = datetime.strptime(end_date_str, "%Y-%m-%d") date_paths = [] current_date = start_date while current_date <= end_date: date_str = current_date.strftime("%Y-%m-%d") date_paths.append(f"gs://{bucket}/{date_str}/*.json") current_date += timedelta(days=1) return date_paths def run(argv=None): pipeline_options = PipelineOptions(argv) # 实际场景建议通过命令行参数传入起止日期 start_date = "2022-01-02" end_date = "2022-01-03" bucket = "ext-pub-testjeff.appspot.com" # 生成日期范围内的所有路径模板 date_paths = generate_date_range_paths(start_date, end_date, bucket) with beam.Pipeline(options=pipeline_options) as p: (p | "读取日期范围内所有文件" >> beam.io.ReadFromText(date_paths) | "数据转换处理" >> beam.Map(lambda x: x) # 替换为实际业务转换逻辑 | "写入输出文件" >> beam.io.WriteToText( "gs://ext-pub-testjeff.appspot.com/output", file_name_suffix=".json" )) if __name__ == "__main__": run()
方案说明
ReadFromText原生支持接收多个带通配符的路径列表,会自动解析每个模板对应的所有文件,无需手动展开。- 通过日期范围生成路径模板,避免手动维护路径列表,适配任意起止日期的输入需求。
- 单一流线的Pipeline设计,实现数据统一读取、处理、写入,性能更优且符合批量处理场景。
进阶场景补充(需文件过滤/验证)
如果需要先验证文件存在性或过滤特定文件,可先通过FileSystems.match列出所有符合条件的具体文件路径,再用beam.Create+ReadAllFromText读取:
def run(argv=None): pipeline_options = PipelineOptions(argv) start_date = "2022-01-02" end_date = "2022-01-03" bucket = "ext-pub-testjeff.appspot.com" date_paths = generate_date_range_paths(start_date, end_date, bucket) # 列出所有匹配的具体文件路径 matched_files = [] for path in date_paths: matches = beam.io.FileSystems.match([path])[0].metadata_list matched_files.extend([match.path for match in matches]) with beam.Pipeline(options=pipeline_options) as p: (p | "创建文件路径列表" >> beam.Create(matched_files) | "读取所有文件" >> beam.io.ReadAllFromText() | "数据转换处理" >> beam.Map(lambda x: x) | "写入输出文件" >> beam.io.WriteToText( "gs://ext-pub-testjeff.appspot.com/output", file_name_suffix=".json" ))
该方式适合需要精细化文件控制的场景,但性能略低于直接使用ReadFromText。
内容的提问来源于stack exchange,提问作者Jeff Huang
相关产品推荐
相关产品推荐

