You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在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. 方案1不可行:ReadAllFromText仅接收具体文件路径,不会解析通配符,因此传入*.json这类通配符会导致找不到文件,直接报错。
  2. 方案2逻辑错误:循环中不断给p赋值,第一次循环后p从Pipeline对象变为PCollection,第二次循环尝试在PCollection上执行ReadFromText(源操作),会直接抛出异常。
  3. 方案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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.05 13:15:55