Apache Airflow Amazon Redshift Connection 完整配置指南数据库认证、IAM 与 Redshift Serverless 详解【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow本篇技术指南以 Apache Airflow Amazon Providerapache-airflow-providers-amazon中的 Redshift Connection 类型为核心系统讲解在 Airflow 中配置 Amazon Redshift 连接的字段语义、三种认证方式数据库密码认证、凭据认证、IAM 认证以及 Redshift Serverless 场景下的扩展参数。读完本文你将掌握redshift_conn_id的完整配置方法能够根据实际集群形态Provisioned 集群或 Serverless 工作组正确组装 Connection 的 Host/Schema/Login/Port 与 Extra JSON并能理解RedshiftSQLHook底层如何将连接参数转换为redshift_connector的调用为编写 Redshift 相关 DAG 与排查连接问题打下基础。Redshift Connection 在 Airflow 中的定位Redshift Connection 是 Airflow 用于与 Amazon Redshift 交互的连接类型其唯一用途是支撑 Airflow 与 Redshift 之间的集成。在仓库中真正消费该连接的是 Amazon Provider 的 Hook 层——尤其是 RedshiftSQLHook它继承自airflow.providers.common.sql.hooks.sql.DbApiHook负责执行 SQL 语句这一核心场景conn_name_attr redshift_conn_id即 DAG 中通过redshift_conn_id参数引用该连接default_conn_name redshift_default即默认连接 IDconn_type redshift即连接类型注册名supports_autocommit True支持自动提交。除了RedshiftSQLHook同一连接还被 RedshiftHookboto3.client(redshift)的薄封装用于集群生命周期管理如create_cluster、cluster_status、delete_cluster以及数据传输 Operatorredshift_to_s3.py、s3_to_redshift.py复用后者同样通过RedshiftSQLHook(redshift_conn_id...)建立会话。认证方式总览连接 Redshift 的认证可以复用redshift_connectorAWS 官方 Python 驱动即 amazon-redshift-python-driver所支持的任何认证方式包括数据库密码认证Database Authentication直接在 Connection 字段中填写 Login/Password凭据认证Credentials Authentication使用 Connection 中保存的登录凭据但以cluster_identifier替代 Host/Port 来唯一定位集群IAM 认证IAM Authentication通过 AWS IAM 获取临时数据库口令需要配合aws_conn_id对应的 AWS 凭据IdP 插件认证使用 Identity Provider 插件属于驱动级扩展能力。从源码看RedshiftSQLHook._get_conn_params()是认证分流的实现核心if conn.extra_dejson.get(iam, False): login, password, port self.get_iam_token(conn) else: login conn.login password conn.password port conn.port即 Extra 中iam: true是切换 IAM 认证的总开关redshift_sql.py。默认连接 IDRedshift Connection 的默认连接 ID 为redshift_default。在 UI 中新建连接时Airflow 会自动预填该 ID在 DAG 中也可以显式传入例如from airflow.providers.amazon.aws.hooks.redshift_sql import RedshiftSQLHook hook RedshiftSQLHook(redshift_conn_idredshift_default)连接字段逐一说明在 Airflow UI 的 Admin → Connections 中新建 Redshift 类型连接或在配置管理中以键值对方式录入时涉及以下字段官方文档标注均为可选实际按认证方式有硬性要求详见下文字段是否必填语义Host可选Amazon Redshift 集群端点Cluster EndpointSchema可选Amazon Redshift 数据库名Login可选连接 Redshift 时使用的用户名Password可选连接 Redshift 时使用的密码Port可选与 Redshift 交互的端口Extra可选JSON 字典形式的额外参数透传给redshift_connector字段到驱动参数的映射同样由_get_conn_params()完成if login: conn_params[user] login if password: conn_params[password] password if port: conn_params[port] port if conn.host: conn_params[host] conn.host if conn.schema: conn_params[database] conn.schema可以看到Schema 字段在底层被映射为驱动参数database这是理解 Redshift Connection 的关键——UI 上的 DatabaseSchema 字段即redshift_connector的database参数。Hook 的get_ui_field_behaviour()也把 UI 显示重命名为{login: User, schema: Database}便于用户在界面中理解。Extra 字段Extra 以 JSON 字典形式填写会被透传给redshift_connector.connect(**conn_kwargs)。get_conn()的组装逻辑为conn_params self._get_conn_params() conn_kwargs_dejson self.conn.extra_dejson conn_kwargs: dict {**conn_params, **conn_kwargs_dejson} return redshift_connector.connect(**conn_kwargs)即先以标准字段生成基础参数再用 Extra 字典覆盖/追加redshift_sql.py。这意味着驱动支持的连接参数如ssl、iam、cluster_identifier、region、db_user、database、profile等都可以放进 Extra并会覆盖同名标准字段。官方驱动文档对支持的参数有完整清单本文后续的 IAM 与 Serverless 示例是 Airflow 侧最常用的组合。若通过 URI 方式配置连接如redshift://user:passhost:5439/dbname请确保 URI 中的每个组件都经过 URL 编码URL-encoded尤其是密码中若含有、:、/等特殊字符。数据库认证Database Authentication这是最直观的方式把集群端点、数据库、用户名、密码、端口直接填到标准字段中。官方文档给出的示例值SchemaDevHostredshift-cluster-1.123456789.us-west-1.redshift.amazonaws.comLoginawsuserPassword********Port5439此时无需 Extra 参数Hook 直接使用conn.login/conn.password/conn.port构造连接。适合开发、测试环境或使用静态数据库密码的生产场景。凭据认证Credentials Authentication凭据认证使用 Connection 中保存的登录凭据连接 Redshift其特殊之处在于不通过 Host/Port 定位集群而是通过cluster_identifier唯一标识集群Port 为必填若未提供使用默认值5439其余 Connection 标准字段如Login假设为空在该方式中cluster_identifier取代Host与Port。从源码结构看cluster_identifier通过 Extra 传入Hook 在 IAM 流程中会优先读取conn.extra_dejson.get(cluster_identifier)仅在未设置时才回退到从 Host 推断见下文 IAM 认证小节。凭据认证适合希望把集群标识与网络端点解耦的自动化场景。IAM 认证IAM AuthenticationIAM 认证通过 AWS IAM 为用户生成临时数据库口令避免在 Airflow 中保存长期有效的数据库密码。其流程与要求如下在 Hook 初始化时使用给定的 AWS IAM 配置aws_conn_id默认aws_default获取临时口令Port 为必填若未提供使用默认值5439Login 与 Schema 均为必填除上述字段外Connection 的其他标准字段假设为空若 Extra 中未设置cluster_identifier则会自动从 Connection 的Host字段推断——具体逻辑见源码get_iam_token()取 Host 的第一个.之前的片段作为集群标识例如redshift-cluster-1.ccdre4hpd39h.us-east-1.redshift.amazonaws.com会被解析为redshift-cluster-1redshift_sql.py。官方文档给出的 Extra 配置示例{ iam: true, cluster_identifier: redshift-cluster-1, port: 5439, region: us-east-1, db_user: awsuser, database: dev, profile: default }各参数作用iam: true开启 IAM 认证总开关cluster_identifier目标集群标识用于get_cluster_credentialsAPI 调用未设置时从 Host 推断port连接端口必填项region集群所在区域db_user需要为其生成临时口令的数据库用户database目标数据库名profile使用的 AWS 配置文件profile。底层实现中get_iam_token()通过AwsBaseHook(aws_conn_idself.aws_conn_id, client_typeredshift).conn调用 boto3 的get_cluster_credentials(DbUser..., DbName..., ClusterIdentifier..., AutoCreateFalse)获取DbPassword与DbUser再用临时口令建立驱动连接redshift_sql.py。注意AutoCreateFalse即不会自动创建用户。提示RedshiftSQLHook的 docstring 特别强调——使用 IAM 认证时把iam放进 Extra 并置为truePassword 字段留空默认使用aws_default连接获取临时令牌也可以在初始化 Hook 时用aws_conn_id覆盖。从 Host 自动推断 cluster_identifier当 Extra 未显式给出cluster_identifier时Hook 的推断规则_get_identifier_from_hostname为以amazonaws.com结尾且恰好 6 段的标准端点my-cluster.id.region.redshift.amazonaws.com→ 取第 1 段与第 3 段拼接为my-cluster.region以amazonaws.com.cn结尾AWS 中国区端点7 段同样取第 1 段与第 3 段其他格式如直接用 IP 连接无法解析时回退为整个 Host 字符串并输出 debug 日志提示预期格式。见 redshift_sql.pyIAM 认证 Redshift Serverless从 Redshift Serverless 诞生起Airflow 就支持通过同一连接类型访问 Serverless 工作组。官方文档给出的 Extra 配置{ iam: true, is_serverless: true, serverless_work_group: default, serverless_token_duration_seconds: 3600, port: 5439, region: us-east-1, database: dev, profile: default }新增/扩展参数语义参数含义约束/默认值is_serverless置为true表示目标是 Redshift Serverless必须配合serverless_work_groupserverless_work_groupServerless 工作组名称必填未设置会抛出AirflowExceptionserverless_token_duration_seconds返回的临时口令有效时长秒最小 900最大 3600默认 3600在底层get_iam_token()检测到is_serverlessTrue后会改用AwsBaseHook(aws_conn_idself.aws_conn_id, client_typeredshift-serverless)调用get_credentials(dbNameconn.schema, workgroupNameserverless_work_group, durationSecondsserverless_token_duration_seconds)换取临时口令redshift_sql.py。因此serverless_work_group缺失时Hook 会直接报错提示Please set serverless_work_group in redshift connection to use IAM with Redshift Serverless.临时口令最长 1 小时3600 秒最短 15 分钟900 秒可按任务耗时调整此时Schema 字段database依然必填因为它作为dbName传入get_credentials。三种认证方式的 Extra 汇总对照场景标准字段Extra 关键参数必填要求数据库认证Host、Schema、Login、Password、Port无Host/Schema/Login/Password/Port凭据认证其余字段留空cluster_identifierPort默认 5439IAM 认证集群Host用于推断、Login、Schema、Portiam: true、cluster_identifier可选Login、Schema、PortIAM 认证ServerlessSchema作为dbName、Portiam: true、is_serverless: true、serverless_work_groupSchema、Port、serverless_work_group连接消费方Hook 与 Operator 使用示例配置好连接后在 DAG 中的典型用法如下。直接执行 SQL基于RedshiftSQLHookfrom airflow import DAG from airflow.providers.amazon.aws.hooks.redshift_sql import RedshiftSQLHook from airflow.utils.dates import days_ago with DAG(dag_idredshift_conn_demo, scheduleNone, start_datedays_ago(1)) as dag: def run_query(): hook RedshiftSQLHook(redshift_conn_idredshift_default) hook.run(SELECT 1;, autocommitTrue) run_query()在数据迁移场景中同一连接被 S3 ⇄ Redshift 传输类任务复用。例如 RedshiftToS3Operator / 对应传输逻辑 内部即通过RedshiftSQLHook(redshift_conn_id...)执行 UNLOAD 语句S3 → Redshift 侧则用其执行 COPY 语句并可调用get_table_primary_key(table, schema)默认 schema 为public查询主键以支撑 upsert 逻辑。这意味着连接配置正确与否直接影响这些高阶任务。此外该连接还参与 OpenLineage 血缘提取RedshiftSQLHook.get_openlineage_database_info()会基于连接构造 Redshift 命名空间scheme 为redshiftauthority 形如cluster.region:port并通过SVV_REDSHIFT_COLUMNS系统视图做跨库列信息查询redshift_sql.py。如果集群启用了数据血缘追踪连接字段Host/Port/Schema的准确性会直接影响血缘命名空间的正确性。常见问题与排查建议连接超时 / 被拒绝get_conn()内置了基于tenacity的重试最多 5 次指数退避最大 20 秒OperationalError连接超时与InterfaceError连接被拒绝都会触发重试redshift_sql.py。若仍失败优先核对 Host 端点、Port默认 5439与安全组/网络可达性。IAM 认证报错提示设置cluster_identifier或host说明 Extra 未提供cluster_identifier且 Host 为空二者至少其一集群场景建议显式设置避免依赖推断。Serverless IAM 认证报错检查is_serverless是否置为true且serverless_work_group是否正确填写。URI 配置乱码检查所有组件是否 URL 编码尤其是密码中的特殊字符。字段空值导致参数缺失_get_conn_params()对空值字段直接跳过if login:等判断驱动将使用自身默认值需要显式指定port的认证方式IAM请务必填写 Port。本文涉及的连接文档原文位于 providers/amazon/docs/connections/redshift.rst核心实现可继续阅读 RedshiftSQLHook 与对应单元测试 test_redshift_sql.py其中测试用例覆盖了iam、cluster_identifier从 Extra 读取与推断等多种组合可作为配置行为的权威验证参考。【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考