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

如何修改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()

说明

  1. FileIO.match的watch参数:启用后会持续监听GCS路径下的新增文件,start_time设为当前时间,确保管道启动前已存在的文件不会被匹配到。
  2. 移除循环:流式Dataflow作业提交后会在Google Cloud上持续运行,自动监听新增文件,原代码的while True会导致重复提交作业,属于错误用法。
  3. 文件读取流程:FileIO.match匹配到新增文件后,通过FileIO.read读取文件内容,再用FlatMap提取每行数据,后续逻辑和原代码保持一致。

内容的提问来源于stack exchange,提问作者Dina Elhusseiny

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.26 21:32:23