Apache Airflow Amazon Redshift 操作指南:Redshift SQL、Redshift Data 与集群全生命周期管理
本文是 Apache Airflow 官方 Amazon Provider 中 Redshift 操作符的完整使用指南,围绕 操作符总览文档 展开。Amazon 为 Redshift 提供了两条查询路径——Python 连接器(Postgres 协议) 与 Redshift Data API(HTTPS),Airflow 对两者均有官方支持;除此之外,本指南还覆盖集群创建、暂停、恢复、快照、删除与状态感知等全生命周期操作。读完本文,你将掌握 SQLExecuteQueryOperator、RedshiftDataOperator 与全套 Redshift*Cluster*Operator/RedshiftClusterSensor 的参数语义、源码级实现原理,并能直接组装出可运行的 Redshift 编排 DAG。
一、两条查询路径:先理解再选型
原文档开篇就给出关键结论:Amazon 提供两种查询 Redshift 的方式,而 Airflow 两者都支持:
- Python 连接器(redshift_connector,走 Postgres 协议)——在 Airflow 中对应 Amazon Redshift SQL 路线,即使用 Postgres 类型的连接(Connection)。
- Amazon Redshift Data API(HTTPS API,基于 boto3)——在 Airflow 中对应 Amazon Redshift Data 路线,即使用
RedshiftDataOperator。
选型原则非常直接,文档给出的建议是:
- 要走 API(HTTP)路线,选择 Amazon Redshift Data(
RedshiftDataOperator); - 要走 Postgres 连接路线,选择 Amazon Redshift SQL(
SQLExecuteQueryOperator+ Redshift 连接)。
两者的本质差异在于:RedshiftDataOperator 通过 AWS API 提交并轮询语句状态,完全不需要 Postgres 连接,因此天然适合无法直连数据库网络(如 VPC 之外)的调度环境;而 SQLExecuteQueryOperator 则需要在 Airflow 中配置可直连的 Redshift 连接信息。这一区分贯穿本文全部示例。
二、环境准备(Prerequisite Tasks)
文档在 前提任务片段 中明确,使用这些操作符需要完成三件事:
-
通过 AWS Console 或 AWS CLI 创建必要的资源(如集群、安全组、子网组等);
-
安装 API 依赖库:
pip install 'apache-airflow[amazon]'详细的安装方式参见 Airflow 安装文档。
-
配置 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: standard、max_attempts: 10)、超时(connect_timeout、read_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_identifier、database、sql、db_user、parameters、statement_name、workgroup_name、session_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') }}",
)
这里有两个值得注意的实现细节:
session_keep_alive_seconds=600保证临时表所在会话在查询结束后仍存活 10 分钟,给下游任务留出执行窗口;- 下游任务通过
xcom_pull(..., key='session_id')动态获取上游返回的会话 ID——上游任务执行后 XCom 中即包含session_id(由 Redshift Data API 在ExecuteStatement时返回,作为操作符结果透传)。
六、Durable Execution:崩溃安全的语句执行
这是 RedshiftDataOperator 在 Airflow 3.3+ 上引入的重要能力,文档用较大篇幅专门讲解,源码实现位于 redshift_data.py 的 ResumableJobMixin 机制中。
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-node、db_name 默认 dev、number_of_nodes 默认 1、port 默认 5439、cluster_version 默认 1.0、automated_snapshot_retention_period 默认 1、allow_version_upgrade 默认 True、encrypted 默认 False。此外还支持 availability_zone、preferred_maintenance_window、cluster_parameter_group_name、kms_key_id、iam_roles、tags 等扩展参数,以及运行行为参数:wait_for_completion 默认 False、max_attempt 默认 5、poll_interval 默认 60、deferrable 默认取全局配置、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=False、poll_interval=30、max_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 默认 False、poll_interval 默认 15、max_attempt 默认 20;删除快照操作符默认 wait_for_completion=True、poll_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=True、poll_interval=30、max_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(如 available、paused)或其他终态。它是典型的轮询型传感器,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 查询按网络环境选择 SQLExecuteQueryOperator 或 RedshiftDataOperator;需要跨任务共享临时表时使用 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 的全部连接参数。
认证方式有三种:
-
数据库认证:填写 Schema、Host、Login、Password、Port(如 Host 为
redshift-cluster-1.123456789.us-west-1.redshift.amazonaws.com,Port 为5439)。 -
凭据认证(Credentials):使用连接中的凭据,必须提供 Port(默认
5439),且假设其余字段(如 Login)为空;此方式下用 cluster_identifier 取代 Host 和 Port 唯一标识集群。 -
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 连接文档。
十一、参考文档索引
- 操作符总览:providers/amazon/docs/operators/redshift/index.rst
- Redshift Data 详解:redshift_data.rst
- Redshift SQL 详解:redshift_sql.rst
- Redshift Cluster 详解:redshift_cluster.rst
- 连接配置:providers/amazon/docs/connections/redshift.rst
- 操作符源码:redshift_data.py、redshift_cluster.py
- 端到端系统测试示例:example_redshift.py
- 通用 SQL 执行示例:example_sql_execute_query.py
- 相关核心概念:任务状态存储、可恢复任务
atomcodeClaude Code 的开源替代方案。连接任意大模型,编辑代码,运行命令,自动验证 — 全自动执行。用 Rust 构建,极致性能。 | An open-source alternative to Claude Code. Connect any LLM, edit code, run commands, and verify changes — autonomously. Built in Rust for speed. Get StartedRust4.21 K638- DDeepSeek-V4.1-FlashDeepSeek-V4.1-Flash 是一个多模态混合专家(MoE)模型,拥有 5520 亿骨干参数,并支持最多一百万 token 的上下文长度。该模型原生支持图像和文本输入,并以自回归方式生成文本Python380
cherry-studio🍒 Cherry Studio 是一款支持多个 LLM 提供商的桌面客户端TypeScript2 K146
hello-agents📚 《从零开始构建智能体》——从零开始的智能体原理与实践教程Python48167
new-apiAI模型聚合管理中转分发系统,一个应用管理您的所有AI模型,支持将多种大模型转为统一格式调用,支持OpenAI、Claude、Gemini等格式,可供个人或者企业内部管理与分发渠道使用。🍥 A Unified AI Model Management & Distribution System. Aggregate all your LLMs into one app and access them via an OpenAI-compatible API, with native support for Claude (Messages) and Gemini formats.Go20743
JeecgBoot🔥企业级低代码平台集成了AI应用平台,帮助企业快速实现低代码开发和构建AI应用!前后端分离架构 SpringBoot,SpringCloud、Mybatis,Ant Design4、 Vue3.0、TS+vite!强大的代码生成器让前后端代码一键生成,无需写任何代码! 引领AI低代码开发模式: AI生成->OnlineCoding-> 代码生成-> 手工MERGE,显著的提高效率,又不失灵活~Java34251