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

如何使用Airflow调度Cloud Dataflow的PostgreSQL转BigQuery ETL任务?

使用Airflow调度Dataflow PostgreSQL到BigQuery的ETL任务

前提准备

  • 确保你的Airflow环境已配置GCP连接(连接ID建议设为google_cloud_default),对应的服务账号需具备:
    • Dataflow作业提交与管理权限
    • BigQuery数据写入权限
    • PostgreSQL(或Cloud SQL)数据读取权限(若为Cloud SQL,需额外配置Cloud SQL代理相关权限)
  • 确认你的Dataflow任务已能手动运行成功,记录核心信息:
    • 模板化任务:模板的GCS存储路径、所需参数(如PostgreSQL连接信息、BigQuery目标表)
    • Python管道任务:管道代码的GCS路径或Airflow可访问路径、运行参数

方法1:调度Dataflow模板化任务

若你的Dataflow任务已导出为GCS上的模板,使用CloudDataflowTemplatedJobOperator是最便捷的方式。

示例DAG代码

from airflow import DAG
from airflow.providers.google.cloud.operators.dataflow import CloudDataflowTemplatedJobOperator
from datetime import datetime, timedelta

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

with DAG(
    'dataflow_postgres_to_bq',
    default_args=default_args,
    description='调度Dataflow同步PostgreSQL到BigQuery',
    schedule_interval='0 2 * * *',  # 每天凌晨2点运行
    start_date=datetime(2024, 1, 1),
    catchup=False,
    tags=['dataflow', 'postgres', 'bigquery'],
) as dag:

    run_dataflow_job = CloudDataflowTemplatedJobOperator(
        task_id='run_postgres_to_bq_dataflow',
        template_path='gs://your-bucket/dataflow-templates/postgres-to-bq-template',
        job_name='postgres-to-bq-sync-{{ ds_nodash }}',  # 用日期生成唯一作业名,避免冲突
        parameters={
            'postgres-host': 'your-postgres-host',
            'postgres-port': '5432',
            'postgres-database': 'your-db-name',
            'postgres-user': 'your-db-user',
            'postgres-password': '{{ var.value.postgres_db_password }}',  # 用Airflow变量存储敏感信息
            'bq-project': 'your-gcp-project-id',
            'bq-dataset': 'your-bq-dataset',
            'bq-table': 'your-bq-target-table',
            # 其他模板要求的参数
        },
        location='us-central1',  # 你的Dataflow作业运行区域
        gcp_conn_id='google_cloud_default',
    )

    run_dataflow_job

方法2:调度Python管道形式的Dataflow任务

若你的Dataflow任务是Python脚本(如postgres_to_bq.py),可使用DataflowPythonOperator直接提交作业。

示例DAG代码

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

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

with DAG(
    'dataflow_postgres_to_bq_python',
    default_args=default_args,
    description='调度Dataflow Python管道同步PostgreSQL到BigQuery',
    schedule_interval='0 2 * * *',
    start_date=datetime(2024, 1, 1),
    catchup=False,
    tags=['dataflow', 'postgres', 'bigquery'],
) as dag:

    run_dataflow_python = DataflowPythonOperator(
        task_id='run_postgres_to_bq_pipeline',
        py_file='gs://your-bucket/dataflow-scripts/postgres_to_bq.py',  # 脚本在GCS的存储路径
        job_name='postgres-to-bq-py-sync-{{ ds_nodash }}',
        options={
            'runner': 'DataflowRunner',
            'project': 'your-gcp-project-id',
            'region': 'us-central1',
            'temp_location': 'gs://your-bucket/dataflow-temp',
            'postgres-host': 'your-postgres-host',
            'postgres-db': 'your-db-name',
            'bq-table': 'your-gcp-project-id:your-bq-dataset.your-target-table',
            # 其他管道所需参数
        },
        gcp_conn_id='google_cloud_default',
    )

    run_dataflow_python

关键注意事项

  • 敏感信息管理:数据库密码等敏感内容禁止硬编码,优先使用Airflow的Variables或加密连接存储
  • 作业命名唯一性:作业名必须包含日期或唯一标识,防止重复提交时出现作业名称冲突
  • 区域一致性:Dataflow作业运行区域尽量与BigQuery数据集、GCS存储桶区域保持一致,降低跨区域传输成本与延迟
  • 监控告警:可在DAG中添加失败告警逻辑(如邮件通知),或集成GCP监控对Dataflow作业状态进行跟踪

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.10 16:15:30