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

Cloud Composer新手求助:基于Airflow传感器实现GCS-CSV-Dataflow-BigQuery流程

Cloud Composer + Dataflow + BigQuery 自动处理GCS新增CSV文件实现方案

一、前置环境准备

  • 启用Google Cloud相关服务:Cloud Composer、Dataflow、BigQuery、Cloud Storage API
  • 创建两个GCS存储桶:一个用于存放源文件(如your-input-bucket),一个用于Dataflow临时文件存储(如your-temp-bucket)
  • 在BigQuery中创建目标数据集(如csv_processing)和对应表,表结构需匹配转换后的字段,示例SQL:
    CREATE TABLE csv_processing.target_table (
        original_id STRING,
        user_name STRING,
        order_amount INT64
    )
    
  • 确保Cloud Composer的服务账号拥有以下权限:GCS存储桶读写权限、Dataflow作业提交权限、BigQuery表写入权限

二、Dataflow 转换代码实现

编写Python脚本csv_to_bq_transform.py,完成CSV读取、字段重命名、格式转换及BigQuery写入:

import argparse
import apache_beam as beam
from apache_beam.options.pipeline_options import PipelineOptions, GoogleCloudOptions
from apache_beam.io.gcp.bigquery import WriteToBigQuery
import csv

class TransformCsvRow(beam.DoFn):
    def process(self, element):
        # 解析源CSV行(假设源字段为:id, name, amount)
        reader = csv.DictReader([element], fieldnames=['id', 'name', 'amount'])
        row = next(reader)
        
        # 字段重命名+格式转换
        transformed_row = {
            'original_id': row['id'],
            'user_name': row['name'].strip(),
            'order_amount': int(row['amount']) if row['amount'].isdigit() else None
        }
        yield transformed_row

def run():
    parser = argparse.ArgumentParser()
    parser.add_argument('--input_gcs_path', required=True, help='GCS源文件路径,如gs://your-input-bucket/test.csv')
    parser.add_argument('--output_bq_table', required=True, help='BigQuery目标表,如your-project-id:csv_processing.target_table')
    parser.add_argument('--temp_gcs_location', required=True, help='Dataflow临时文件路径,如gs://your-temp-bucket/temp/')
    known_args, pipeline_args = parser.parse_known_args()

    pipeline_options = PipelineOptions(pipeline_args)
    google_cloud_options = pipeline_options.view_as(GoogleCloudOptions)
    google_cloud_options.temp_location = known_args.temp_gcs_location

    with beam.Pipeline(options=pipeline_options) as p:
        (
            p
            | '读取CSV文件' >> beam.io.ReadFromText(known_args.input_gcs_path, skip_header_lines=1)
            | '字段转换与重命名' >> beam.ParDo(TransformCsvRow())
            | '写入BigQuery' >> WriteToBigQuery(
                known_args.output_bq_table,
                write_disposition=beam.io.BigQueryDisposition.WRITE_APPEND,
                create_disposition=beam.io.BigQueryDisposition.CREATE_IF_NEEDED
            )
        )

if __name__ == '__main__':
    run()

将该脚本上传到GCS的固定路径,比如gs://your-temp-bucket/dataflow-scripts/csv_to_bq_transform.py

三、Airflow DAG 编写(核心触发逻辑)

编写DAG文件gcs_csv_trigger_dag.py,通过GCSObjectSensor监听新增CSV文件,自动触发Dataflow作业:

from airflow import DAG
from airflow.providers.google.cloud.sensors.gcs import GCSObjectSensor
from airflow.providers.google.cloud.operators.dataflow import DataflowPythonOperator
from datetime import datetime, timedelta

default_args = {
    'owner': 'airflow',
    'depends_on_past': False,
    'start_date': datetime(2024, 1, 1),
    'email_on_failure': False,
    'email_on_retry': False,
    'retries': 1,
    'retry_delay': timedelta(minutes=5),
}

with DAG(
    'gcs_csv_auto_trigger',
    default_args=default_args,
    description='监听GCS新增CSV文件,自动触发Dataflow转换至BigQuery',
    schedule_interval=None,  # 依赖传感器触发,不设置定时调度
    catchup=False,
    tags=['gcs', 'dataflow', 'bigquery'],
) as dag:

    # 监听GCS中新增的CSV文件
    gcs_csv_sensor = GCSObjectSensor(
        task_id='wait_for_new_csv',
        bucket='your-input-bucket',
        object_pattern='*.csv',  # 匹配所有CSV后缀文件
        mode='reschedule',  # 文件未出现时释放worker资源,定时重试
        poke_interval=30,  # 每30秒检查一次
        timeout=3600,  # 超时时间1小时,可按需调整
    )

    # 触发Dataflow转换作业
    execute_dataflow = DataflowPythonOperator(
        task_id='run_csv_transform',
        py_file='gs://your-temp-bucket/dataflow-scripts/csv_to_bq_transform.py',
        options={
            'input_gcs_path': 'gs://your-input-bucket/{{ task_instance.xcom_pull(task_ids="wait_for_new_csv") }}',
            'output_bq_table': 'your-project-id:csv_processing.target_table',
            'temp_gcs_location': 'gs://your-temp-bucket/temp/',
            'project': 'your-project-id',
            'region': 'us-central1',  # 替换为你的GCP区域
        },
        dag=dag,
    )

    # 设置任务依赖:传感器检测到文件后触发Dataflow作业
    gcs_csv_sensor >> execute_dataflow

四、部署与测试

  1. 部署DAG:
    • 通过GCS浏览器将DAG文件上传到Cloud Composer的DAG目录:gs://your-composer-bucket/dags/
    • 或使用gcloud命令:
      gcloud composer environments storage dags import --environment your-composer-env --location your-region --source gcs_csv_trigger_dag.py
      
  2. 验证DAG加载:登录Airflow UI,确认gcs_csv_auto_trigger DAG已启用
  3. 测试流程:
    • 上传一个符合格式的CSV文件到your-input-bucket
    • 在Airflow UI中查看wait_for_new_csv任务状态,成功后会自动触发run_csv_transform
    • 前往Dataflow控制台确认作业运行,完成后检查BigQuery目标表是否有数据写入

关键注意事项

  • 重复文件处理:若需避免重复处理,可在Dataflow作业完成后将源文件移动至归档目录,或在传感器中通过文件名前缀/后缀过滤已处理文件
  • CSV格式兼容:如果CSV包含引号、自定义分隔符,建议改用beam.io.ReadFromCsv或csv.DictReader解析,避免手动拆分出错
  • 资源优化:传感器使用reschedule模式比poke模式更节省worker资源,适合长期监听场景

内容的提问来源于stack exchange,提问作者Nhu Dao

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.18 16:35:22