如何修改Dataflow流式代码,仅监听处理新增文件并写入BigQuery
修改Dataflow流式代码仅处理GCS新增文件
要实现仅监听并处理后续新增的GCS文件,核心是调整文件读取逻辑以跳过已有文件,同时修正流式作业的运行方式。以下是修改后的完整代码及说明:
关键修改点
- 替换
ReadFromText为FileIO.match+FileIO.read组合,通过watch模式监听新增文件 - 设置
watch的start_time为管道启动时的当前时间,确保只处理启动后新增的文件 - 移除不必要的
while True循环,流式Dataflow作业会持续运行,无需手动重启
修改后的代码
import logging import json from datetime import datetime import apache_beam as beam from apache_beam.options.pipeline_options import PipelineOptions, StandardOptions from apache_beam.io.filesystem import CompressionTypes # 配置日志 logging.basicConfig(level=logging.INFO) # 定义管道选项 pipeline_options = PipelineOptions( project='moneyfellows-data', runner='DataflowRunner', job_name='streaming-bundle-test', temp_location='gs://mf-staging-area/temp', region='us-central1' ) standard_options = pipeline_options.view_as(StandardOptions) standard_options.streaming = True # 定义管道 with beam.Pipeline(options=pipeline_options) as pipeline: # 监听GCS新增文件(仅处理管道启动后新增的文件) lines = ( pipeline | "MatchNewFiles" >> beam.io.FileIO.match( file_pattern=input_pattern, # 设置监听起始时间为当前时间,忽略已有文件 watch=beam.io.FileIO.Watch(start_time=datetime.now()) ) | "ReadMatchedFiles" >> beam.io.FileIO.read(compression_type=CompressionTypes.GZIP) | "ExtractLines" >> beam.FlatMap(lambda file: file.readlines()) | "ParseJSON" >> beam.Map(json.loads) ) # 写入BigQuery lines | "WriteToBigQuery" >> beam.io.WriteToBigQuery( output_table, create_disposition=beam.io.BigQueryDisposition.CREATE_IF_NEEDED, write_disposition=beam.io.BigQueryDisposition.WRITE_APPEND, ) # 启动流式管道,无需循环 logging.info("启动流式Dataflow作业,开始监听GCS新增文件") pipeline_result = pipeline.run() pipeline_result.wait_until_finish()
说明
FileIO.match的watch参数:启用后会持续监听GCS路径下的新增文件,start_time设为当前时间,确保管道启动前已存在的文件不会被匹配到。- 移除循环:流式Dataflow作业提交后会在Google Cloud上持续运行,自动监听新增文件,原代码的
while True会导致重复提交作业,属于错误用法。 - 文件读取流程:
FileIO.match匹配到新增文件后,通过FileIO.read读取文件内容,再用FlatMap提取每行数据,后续逻辑和原代码保持一致。
内容的提问来源于stack exchange,提问作者Dina Elhusseiny
相关产品推荐
相关产品推荐

