首页
/ Apache Airflow Amazon Redshift 操作指南:Redshift SQL、Redshift Data 与集群全生命周期管理

Apache Airflow Amazon Redshift 操作指南:Redshift SQL、Redshift Data 与集群全生命周期管理

2026-09-12 11:58:15作者:明树来

本文是 Apache Airflow 官方 Amazon Provider 中 Redshift 操作符的完整使用指南,围绕 操作符总览文档 展开。Amazon 为 Redshift 提供了两条查询路径——Python 连接器(Postgres 协议)Redshift Data API(HTTPS),Airflow 对两者均有官方支持;除此之外,本指南还覆盖集群创建、暂停、恢复、快照、删除与状态感知等全生命周期操作。读完本文,你将掌握 SQLExecuteQueryOperatorRedshiftDataOperator 与全套 Redshift*Cluster*Operator/RedshiftClusterSensor 的参数语义、源码级实现原理,并能直接组装出可运行的 Redshift 编排 DAG。

一、两条查询路径:先理解再选型

原文档开篇就给出关键结论:Amazon 提供两种查询 Redshift 的方式,而 Airflow 两者都支持:

  1. Python 连接器(redshift_connector,走 Postgres 协议)——在 Airflow 中对应 Amazon Redshift SQL 路线,即使用 Postgres 类型的连接(Connection)。
  2. Amazon Redshift Data API(HTTPS API,基于 boto3)——在 Airflow 中对应 Amazon Redshift Data 路线,即使用 RedshiftDataOperator

选型原则非常直接,文档给出的建议是:

  • 要走 API(HTTP)路线,选择 Amazon Redshift DataRedshiftDataOperator);
  • 要走 Postgres 连接路线,选择 Amazon Redshift SQLSQLExecuteQueryOperator + Redshift 连接)。

两者的本质差异在于:RedshiftDataOperator 通过 AWS API 提交并轮询语句状态,完全不需要 Postgres 连接,因此天然适合无法直连数据库网络(如 VPC 之外)的调度环境;而 SQLExecuteQueryOperator 则需要在 Airflow 中配置可直连的 Redshift 连接信息。这一区分贯穿本文全部示例。

二、环境准备(Prerequisite Tasks)

文档在 前提任务片段 中明确,使用这些操作符需要完成三件事:

  1. 通过 AWS ConsoleAWS CLI 创建必要的资源(如集群、安全组、子网组等);

  2. 安装 API 依赖库:

    pip install 'apache-airflow[amazon]'
    

    详细的安装方式参见 Airflow 安装文档

  3. 配置 AWS 连接(Connection),参见 连接配置

此外,Redshift Data 与 Cluster 系列操作符还共享一份通用参数(见 generic_parameters 片段):

  • aws_conn_id:AWS 连接 ID,默认 aws_default;若设为 None 则使用默认 boto3 行为(不查连接)。
  • region_name:AWS 区域名,默认 None,此时取 AWS 连接 Extra 中的 region_name。
  • verify:是否校验 SSL 证书;可传 False 关闭校验,或传 CA 证书包文件路径;默认 None 时取连接配置。
  • botocore_config:用于构造 botocore.config.Config 的字典,可配置重试(如 mode: standardmax_attempts: 10)、超时(connect_timeoutread_timeout)、tcp_keepalive 等,用于规避限流异常。注意:显式传入空字典 {} 会覆盖连接中的配置。

三、Amazon Redshift SQL:用 SQLExecuteQueryOperator 执行查询

Amazon Redshift SQL 文档 中,执行 SQL 查询的推荐方式是使用通用的 SQLExecuteQueryOperator,配合文档中所指的 redshift 连接类型(连接配置细节见本文第十节)。该操作符属于 apache-airflow-providers-common-sql,完整示例位于 example_sql_execute_query.py,核心片段如下:

execute_query = SQLExecuteQueryOperator(
    task_id="execute_query",
    sql=f"SELECT 1; SELECT * FROM {AIRFLOW_DB_METADATA_TABLE} LIMIT 1;",
    split_statements=True,
    return_last=False,
)

要点说明:

  • sql 可包含多条语句,配合 split_statements=True 拆分执行;
  • return_last=False 表示不返回最后一条语句的结果;
  • 该操作符同样支持 TaskFlow 装饰器形式 @task.sql(...),见同一示例文件;
  • 不配置 Redshift 连接就执行 SQL 的场景,文档明确指向 RedshiftDataOperator(下一节),二者形成互补。

四、Amazon Redshift Data:用 RedshiftDataOperator 执行语句

Amazon Redshift Data 文档 的核心是 RedshiftDataOperator(类定义见 redshift_data.py)。它与 RedshiftSQLOperator 的关键区别正如文档所言:通过 AWS API 查询并取回数据,无需 Postgres 连接

4.1 基础用法

来自系统测试 DAG example_redshift.py 的官方示例:

create_table_redshift_data = RedshiftDataOperator(
    task_id="create_table_redshift_data",
    cluster_identifier=redshift_cluster_identifier,
    database=DB_NAME,
    db_user=DB_LOGIN,
    sql=[
        """
        CREATE TABLE IF NOT EXISTS fruit (
        fruit_id INTEGER,
        name VARCHAR NOT NULL,
        color VARCHAR NOT NULL
        );
    """
    ],
    poll_interval=POLL_INTERVAL,
    wait_for_completion=True,
)

4.2 核心参数语义(源自源码 docstring 与构造签名)

redshift_data.py__init__ 签名可确认以下参数及默认值:

参数 默认值 说明
sql 必填 单条或多条 SQL 语句(str 或 list)
database None 数据库名
cluster_identifier None 集群唯一标识符
db_user None 数据库用户名
parameters None SQL 语句参数列表
secret_arn None 用于数据库访问的 Secret ARN
statement_name None SQL 语句名称
with_event False 是否向 EventBridge 发送事件
wait_for_completion True 是否等待语句完成
poll_interval 10 轮询语句状态的间隔(秒)
return_sql_result False True 时返回 SQL 结果,False(默认)时返回 statement ID
workgroup_name None Redshift Serverless 工作组名,与 cluster_identifier 互斥
session_id None 会话标识,用于会话复用
session_keep_alive_seconds None 查询结束后会话保持存活秒数,最长 24 小时
cancel_on_kill True 任务被杀时是否取消正在运行的 Redshift 语句(含 deferrable 等待期间)
deferrable 取配置 operators.default_deferrable 是否可延迟执行
durable None(默认开启) 持久化执行开关,见第六节

另外注意 template_fields 包含 cluster_identifierdatabasesqldb_userparametersstatement_nameworkgroup_namesession_id,且 template_ext = (".sql",),意味着 sql 参数支持 Jinja 模板与 .sql 文件引用。

4.3 多语句批量执行

示例中通过 sql 列表一次提交多条 INSERT:

insert_data = RedshiftDataOperator(
    task_id="insert_data",
    cluster_identifier=redshift_cluster_identifier,
    database=DB_NAME,
    db_user=DB_LOGIN,
    sql=[
        "INSERT INTO fruit VALUES ( 1, 'Banana', 'Yellow');",
        "INSERT INTO fruit VALUES ( 2, 'Apple', 'Red');",
        "INSERT INTO fruit VALUES ( 3, 'Lemon', 'Yellow');",
        ...
    ],
    poll_interval=POLL_INTERVAL,
    wait_for_completion=True,
)

五、复用会话:临时表的跨任务协作

Redshift Data API 的会话机制可用于跨任务共享 临时表。文档给出的模式是:

  • 上游任务中设置 session_keep_alive_seconds
  • 下游任务中通过 XCom 取回 session_id 并传入 session_id 参数。

官方示例(见 example_redshift.py):

create_tmp_table_data_api = RedshiftDataOperator(
    task_id="create_tmp_table_data_api",
    cluster_identifier=redshift_cluster_identifier,
    database=DB_NAME,
    db_user=DB_LOGIN,
    sql=[
        """
        CREATE TEMPORARY TABLE tmp_people (
        id INTEGER,
        first_name VARCHAR(100),
        age INTEGER
        );
    """
    ],
    poll_interval=POLL_INTERVAL,
    wait_for_completion=True,
    session_keep_alive_seconds=600,
)

insert_data_reuse_session = RedshiftDataOperator(
    task_id="insert_data_reuse_session",
    sql=[
        "INSERT INTO tmp_people VALUES ( 1, 'Bob', 30);",
        "INSERT INTO tmp_people VALUES ( 2, 'Alice', 35);",
        "INSERT INTO tmp_people VALUES ( 3, 'Charlie', 40);",
    ],
    poll_interval=POLL_INTERVAL,
    wait_for_completion=True,
    session_id="{{ task_instance.xcom_pull(task_ids='create_tmp_table_data_api', key='session_id') }}",
)

这里有两个值得注意的实现细节:

  1. session_keep_alive_seconds=600 保证临时表所在会话在查询结束后仍存活 10 分钟,给下游任务留出执行窗口;
  2. 下游任务通过 xcom_pull(..., key='session_id') 动态获取上游返回的会话 ID——上游任务执行后 XCom 中即包含 session_id(由 Redshift Data API 在 ExecuteStatement 时返回,作为操作符结果透传)。

六、Durable Execution:崩溃安全的语句执行

这是 RedshiftDataOperator 在 Airflow 3.3+ 上引入的重要能力,文档用较大篇幅专门讲解,源码实现位于 redshift_data.pyResumableJobMixin 机制中。

6.1 工作原理

RedshiftDataOperator 提交语句后在 worker 上轮询至完成。默认情况下它以 durable(持久化) 模式运行,使其具备崩溃安全特性:

  • 轮询开始前,Redshift 的 statement ID 会被持久化到 任务状态存储(task state store)
  • 若 worker 崩溃或被抢占、任务随后重试,操作符会重新连接到 Redshift 中已在执行的语句,而不是重新提交 SQL

重试时操作符会检查先前语句的状态:

  • 仍在运行 → 重连并继续轮询;
  • 已成功 → 立即返回,不重新提交;
  • 已失败终止,或 ID 过期无法找到 → 重新提交 SQL。

6.2 边界行为

  • wait_for_completion=False 同样受保护:即使该次尝试从不轮询,只要提交成功,语句 ID 也在提交后立即持久化,因此重试时依然重连而非重复提交。
  • 状态存储不可用:运行时若任务状态存储不可用,操作符会记录"崩溃恢复已禁用"并回退到原有行为(重试时重新提交)。
  • 清理与保留期:持久化的 statement ID 不会被自动删除,只有执行 airflow state-store clean 才会清理。若任务的 retry_delay 大于 [state_store] default_retention_days(默认 30 天)且期间执行了清理,则下次重试时 ID 已不存在,操作符将重新提交 SQL。因此清理调度周期不应短于最长的 retry_delay
  • 清除任务(Clearing)视为重试:若语句已成功,清除不会删除已存储的 ID,下次尝试会读回并直接返回、不重复提交。这涉及 Airflow 的可恢复任务语义与 <a href="https://link.gitcode.com/i/b4ba7103e5f21c1ce6392a1dbfb77d25" target="_blank">state_store] clear_on_success 配置,详见 [可恢复任务文档。
  • deferrable 优先级更高:durable 只作用于同步路径;当 deferrable=True 时,Triggerer 本身已跨等待跟踪语句,此时 durable 不生效。

6.3 版本前提与退出机制

  • durable 依赖任务状态存储,要求 Airflow 3.3 或更高版本;在更早版本上该标志无效:显式设置只发出警告(见源码中 _warn_and_disable_durable_pre_3_3 函数),操作符始终在重试时重新提交 SQL。
  • 如需退出该机制(重试时总是重新提交),设置 durable=False
statement = RedshiftDataOperator(
    task_id="redshift_data",
    database="dev",
    sql="SELECT * FROM table",
    cluster_identifier="cluster_identifier",
    durable=False,
)

七、集群全生命周期:Redshift Cluster 系列操作符

Amazon Redshift (Cluster) 文档 覆盖了集群管理的六个操作符与一个传感器,全部类定义位于 redshift_cluster.py

7.1 创建集群:RedshiftCreateClusterOperator

create_cluster = RedshiftCreateClusterOperator(
    task_id="create_cluster",
    cluster_identifier=redshift_cluster_identifier,
    vpc_security_group_ids=[security_group_id],
    cluster_subnet_group_name=cluster_subnet_group_name,
    publicly_accessible=False,
    cluster_type="single-node",
    node_type="ra3.large",
    master_username=DB_LOGIN,
    master_user_password=DB_PASS,
)

从源码构造签名可补充以下默认参数:cluster_type 默认 multi-nodedb_name 默认 devnumber_of_nodes 默认 1port 默认 5439cluster_version 默认 1.0automated_snapshot_retention_period 默认 1allow_version_upgrade 默认 Trueencrypted 默认 False。此外还支持 availability_zonepreferred_maintenance_windowcluster_parameter_group_namekms_key_idiam_rolestags 等扩展参数,以及运行行为参数:wait_for_completion 默认 Falsemax_attempt 默认 5poll_interval 默认 60deferrable 默认取全局配置、delete_cluster_on_failure 默认 True(失败时自动清理并重试,cleanup_timeout_seconds 默认 300)。

7.2 暂停与恢复:RedshiftPauseClusterOperator / RedshiftResumeClusterOperator

pause_cluster = RedshiftPauseClusterOperator(
    task_id="pause_cluster",
    cluster_identifier=redshift_cluster_identifier,
)

resume_cluster = RedshiftResumeClusterOperator(
    task_id="resume_cluster",
    cluster_identifier=redshift_cluster_identifier,
)

两者均支持 deferrable=True延迟模式运行:任务从 Airflow worker 槽位释放,轮询交由 Triggerer(触发器)完成,节省 worker 资源。默认 wait_for_completion=Falsepoll_interval=30max_attempts=30

7.3 快照:创建与删除

create_cluster_snapshot = RedshiftCreateClusterSnapshotOperator(
    task_id="create_cluster_snapshot",
    cluster_identifier=redshift_cluster_identifier,
    snapshot_identifier=redshift_cluster_snapshot_identifier,
    poll_interval=30,
    max_attempt=100,
    retention_period=1,
    wait_for_completion=True,
)

delete_cluster_snapshot = RedshiftDeleteClusterSnapshotOperator(
    task_id="delete_cluster_snapshot",
    cluster_identifier=redshift_cluster_identifier,
    snapshot_identifier=redshift_cluster_snapshot_identifier,
)

创建快照操作符的源码默认值:retention_period 默认 -1(不设保留期)、wait_for_completion 默认 Falsepoll_interval 默认 15max_attempt 默认 20;删除快照操作符默认 wait_for_completion=Truepoll_interval=10

7.4 删除集群:RedshiftDeleteClusterOperator

delete_cluster = RedshiftDeleteClusterOperator(
    task_id="delete_cluster",
    cluster_identifier=redshift_cluster_identifier,
)
delete_cluster.trigger_rule = TriggerRule.ALL_DONE
delete_cluster.max_attempts = 50

源码默认参数:skip_final_cluster_snapshot=True(默认跳过最终快照)、wait_for_completion=Truepoll_interval=30max_attempts=30。示例 DAG 中将其 trigger_rule 设为 ALL_DONE,保证无论前面任务成败都会执行清理,并将 max_attempts 提升到 50 以应对删除的长时间等待。

7.5 等待集群状态:RedshiftClusterSensor

wait_cluster_available = RedshiftClusterSensor(
    task_id="wait_cluster_available",
    cluster_identifier=redshift_cluster_identifier,
    target_status="available",
    poke_interval=15,
    timeout=60 * 30,
)

RedshiftClusterSensor 持续检查集群状态,直到达到 target_status(如 availablepaused)或其他终态。它是典型的轮询型传感器,poke_interval 控制探测间隔,timeout 控制超时。

八、端到端示例 DAG:完整生命周期编排

将上述组件串起来的官方系统测试 DAG 位于 example_redshift.py,其任务依赖链完整呈现了"创建 → 等待可用 → 快照 → 暂停 → 恢复 → 数据操作 → 清理"的编排思路:

chain(
    test_context,              # 测试上下文(获取安全组、子网组等变量)
    create_cluster,            # 创建集群
    wait_cluster_available,    # 等待 available
    create_cluster_snapshot,   # 创建快照
    wait_cluster_available_before_pause,
    pause_cluster,             # 暂停集群
    wait_cluster_paused,       # 等待 paused
    resume_cluster,            # 恢复集群
    wait_cluster_available_after_resume,
    create_table_redshift_data,  # Redshift Data API 建表
    insert_data,                 # 批量插入
    delete_cluster_snapshot,     # 删除快照
    delete_cluster,              # 删除集群(ALL_DONE)
)

会话复用路径则单独并行执行:

chain(
    wait_cluster_available_after_resume,
    create_tmp_table_data_api,    # 建临时表,session_keep_alive_seconds=600
    insert_data_reuse_session,    # 通过 XCom 复用 session_id
    delete_cluster_snapshot,
)

该 DAG 也展示了两个实用技巧:通过 SystemTestContextBuilder 从 Airflow 变量获取外部资源(安全组、子网组);通过 watcher 任务配合 ALL_DONE 触发器正确标记成功/失败。生产环境中可参照此结构,把集群创建/销毁放在低频的初始化与清理分支,把日常 SQL 任务放在中间环节。

九、选择指引与组合建议

综合原文档与源码,两条执行路径的适用场景可归纳如下:

维度 Amazon Redshift SQL(SQLExecuteQueryOperator Amazon Redshift Data(RedshiftDataOperator
连接方式 Postgres 类型 Redshift 连接(直连) AWS API(boto3,无需数据库连接)
适用网络 可从 worker 直连数据库 任意网络可达 AWS API 即可
返回值 数据库查询结果 statement ID 或 SQL 结果(return_sql_result=True
特殊能力 通用 SQL 操作符生态、TaskFlow 装饰器 会话复用、durable 崩溃安全、Serverless 支持
依赖 common-sql provider amazon provider + boto3

典型组合:集群级操作(创建/暂停/恢复/删除/快照)用 Cluster 系列操作符 + RedshiftClusterSensor 把控状态;日常 ETL 查询按网络环境选择 SQLExecuteQueryOperatorRedshiftDataOperator;需要跨任务共享临时表时使用 session_id/session_keep_alive_seconds 会话复用;对重试安全要求高的生产任务,在 Airflow 3.3+ 上保持 durable 默认开启。

十、Redshift 连接配置与认证方式

Redshift 连接文档SQLExecuteQueryOperator 路线的基础。连接类型支持 Redshift 集成,默认连接 ID 为 redshift_default,配置字段如下:

  • Host(可选):Redshift 集群端点。
  • Schema(可选):Redshift 数据库名。
  • Login(可选):认证用户名。
  • Password(可选):认证密码。
  • Port(可选):交互端口(默认 5439)。
  • Extra(可选):JSON 字典形式的扩展参数,支持 redshift_connector 的全部连接参数。

认证方式有三种:

  1. 数据库认证:填写 Schema、Host、Login、Password、Port(如 Host 为 redshift-cluster-1.123456789.us-west-1.redshift.amazonaws.com,Port 为 5439)。

  2. 凭据认证(Credentials):使用连接中的凭据,必须提供 Port(默认 5439),且假设其余字段(如 Login)为空;此方式下用 cluster_identifier 取代 Host 和 Port 唯一标识集群。

  3. IAM 认证:通过 IAM 获取临时密码连接。要求提供 Port、Login、Schema,Extra 示例:

    {
      "iam": true,
      "cluster_identifier": "redshift-cluster-1",
      "port": 5439,
      "region": "us-east-1",
      "db_user": "awsuser",
      "database": "dev",
      "profile": "default"
    }
    

    若 Extra 中未设置 cluster_identifier,会自动从 Host 字段推断。IAM 认证同样支持 Redshift Serverless:将 is_serverless 设为 true 并提供 serverless_work_group,可用 serverless_token_duration_seconds 控制临时密码有效期(最小 900 秒、最大 3600 秒、默认 3600 秒):

    {
      "iam": true,
      "is_serverless": true,
      "serverless_work_group": "default",
      "serverless_token_duration_seconds": 3600,
      "port": 5439,
      "region": "us-east-1",
      "database": "dev",
      "profile": "default"
    }
    

注意:通过 URI 方式配置连接时,所有组件必须进行 URL 编码。其余 AWS 连接(供 RedshiftDataOperator/Cluster 操作符使用)的通用配置可参考 AWS 连接文档

十一、参考文档索引

登录后查看全文
热门项目推荐
相关项目推荐

项目优选

收起
kernelkernel
deepin linux kernel
C
34
18
ops-transformerops-transformer
本项目是CANN提供的transformer类大模型算子库,实现网络在NPU上加速计算。
C++
1.16 K
2.78 K
docsdocs
暂无描述
Markdown
904
5.83 K
ops-nnops-nn
本项目是CANN提供的神经网络类计算算子库,实现网络在NPU上加速计算。
C++
934
1.86 K
pytorchpytorch
作为 Ascend for PyTorch 社区的核心组件,TorchNPU 是昇腾专为 PyTorch 打造的深度学习适配插件,使 PyTorch 框架能够直接调用昇腾 NPU,为开发者提供昇腾 AI 处理器的超强算力。
Python
862
1.36 K
kernelkernel
openEuler内核是openEuler操作系统的核心,既是系统性能与稳定性的基石,也是连接处理器、设备与服务的桥梁。
C
535
606
ops-mathops-math
本项目是CANN提供的数学类基础计算算子库,实现网络在NPU上加速计算。
C++
1.38 K
1.47 K
jiuwenswarmjiuwenswarm
JiuwenSwarm 是一款基于openJiuwen开发的智能AI Agent,它能够将大语言模型的强大能力,通过你日常使用的各类通讯应用,直接延伸至你的指尖。
Python
4.02 K
1.03 K
AscendNPU-IRAscendNPU-IR
AscendNPU-IR是基于MLIR(Multi-Level Intermediate Representation)构建的,面向昇腾亲和算子编译时使用的中间表示,提供昇腾完备表达能力,通过编译优化提升昇腾AI处理器计算效率,支持通过生态框架使能昇腾AI处理器与深度调优
C++
549
399
cann-learning-hubcann-learning-hub
CANN 学习中心仓,支持在线互动运行、边学边练,提供教程、示例与优化方案,一站式助力昇腾开发者快速上手。
Jupyter Notebook
1.07 K
537