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

如何用Python在Airflow中调度多个Redshift(PL/pgSQL)存储过程?

Airflow调度Redshift存储过程实现方案

一、Redshift数据库连接配置

UI配置方式(常用)

  • 登录Airflow控制台,进入Admin > Connections页面
  • 点击**+ Add a new record**,填写以下关键字段:
    • Conn Id:自定义标识(比如redshift_default,后续代码会用到)
    • Conn Type:选择Amazon Redshift
    • Host:Redshift集群的端点地址(格式类似your-cluster.xxxxxx.us-west-2.redshift.amazonaws.com)
    • Schema:要执行存储过程的目标数据库名
    • Login:Redshift访问用户名
    • Password:Redshift用户密码
    • Port:默认5439,按实际集群配置填写
  • 保存后,代码即可通过这个Conn Id调用数据库连接

代码配置方式(适合自动化部署)

如果需要批量或脚本化配置连接,可使用以下代码(生产环境建议结合Airflow Secrets或环境变量存储敏感信息):

from airflow.models import Connection
from airflow.utils.db import create_session

with create_session() as session:
    conn = Connection(
        conn_id='redshift_default',
        conn_type='postgres',
        host='your-cluster-endpoint',
        schema='target_db',
        login='db_user',
        password='db_pass',
        port=5439,
    )
    session.add(conn)
    session.commit()

二、DAG核心实现:并行任务+每周调度

完整代码示例

from airflow import DAG
from airflow.operators.dummy import DummyOperator
from airflow.providers.postgres.operators.postgres import PostgresOperator
from datetime import datetime, timedelta

# DAG默认参数
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(
    'redshift_stored_procs_weekly',
    default_args=default_args,
    description='每周并行调度Redshift存储过程,stored_p_1略早启动',
    schedule_interval='0 0 * * 0',  # 每周日凌晨0点执行,等价于@weekly
    catchup=False,
    tags=['redshift', 'stored_procedure'],
) as dag:

    # 起始空任务,用于统一触发所有存储过程任务
    start = DummyOperator(task_id='start')

    # stored_p_1 配置30秒延迟启动,确保比其他任务早执行
    run_sp1 = PostgresOperator(
        task_id='run_stored_p_1',
        postgres_conn_id='redshift_default',
        sql="CALL stored_p_1();",
        execution_delay=timedelta(seconds=30),
    )

    # stored_p_2、stored_p_3 与start直接关联,和sp1并行执行
    run_sp2 = PostgresOperator(
        task_id='run_stored_p_2',
        postgres_conn_id='redshift_default',
        sql="CALL stored_p_2();",
    )

    run_sp3 = PostgresOperator(
        task_id='run_stored_p_3',
        postgres_conn_id='redshift_default',
        sql="CALL stored_p_3();",
    )

    # 任务依赖配置:start触发所有任务,sp1通过延迟实现早启动
    start >> [run_sp1, run_sp2, run_sp3]

三、特殊需求解决方案:stored_p_1早启动但不阻塞其他任务

因为直接设置任务依赖会导致其他任务等待sp1完成,违反你的需求,这里用两个关键点解决:

  1. 统一起始触发:用DummyOperator作为所有任务的父节点,确保三个存储过程任务都能并行启动
  2. 延迟启动sp1:给run_sp1添加execution_delay参数,让它比其他任务早30秒启动,保证初期数据写入完成后,其他任务能访问到数据

如果需要更精确的数据依赖(比如确保sp1写入的特定数据存在),可以在stored_p_2和stored_p_3的PL/pgSQL代码中添加轮询逻辑:

-- 在stored_p_2开头添加数据检查
WHILE NOT EXISTS (SELECT 1 FROM target_table WHERE data_status = 'initialized') LOOP
    PERFORM pg_sleep(5); -- 每5秒检查一次,直到数据就绪
END LOOP;

-- 执行存储过程核心逻辑

四、关键细节说明

  • 并行执行:通过start >> [run_sp1, run_sp2, run_sp3]的依赖关系,三个任务会同时被触发,实现并行调度
  • 调度周期:schedule_interval='0 0 * * 0'表示每周日凌晨0点执行,也可以用@weekly简化写法
  • 生产环境注意:禁止硬编码密码,建议使用Airflow Secrets Manager或环境变量存储敏感凭证

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.02 22:20:41