Apache Airflow Amazon 提供包实战使用 RedshiftToS3Operator 将 Amazon Redshift 表数据卸载UNLOAD到 Amazon S3【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflowApache Airflow 的 Amazon 提供包apache-airflow[amazon]内置了RedshiftToS3Operator传输操作符用于将 Amazon Redshift 表或查询结果通过 Redshift 的UNLOAD命令导出为 S3 文件。本文以仓库内官方文档 redshift_to_s3.rst 为核心骨架结合源码、系统级示例 DAG 与单元测试系统讲解该操作符的完整用法、参数语义、底层执行原理、Redshift Data API 接入方式与 OpenLineage 血缘支持帮助读者快速搭建Redshift → S3的数据导出链路。前置任务环境准备与连接配置使用该传输操作符之前需要完成三件事原文参见 prerequisite_tasks.rst创建必要的 AWS 资源在 AWS 控制台或通过 AWS CLI 准备好目标 S3 桶、Redshift 集群以及具备相应权限的 IAM 角色/用户。通过 pip 安装 Amazon 提供包pip install apache-airflow[amazon]配置 AWS 连接Connection在 Airflow 中建立类型为aws的连接默认连接 ID 为aws_default用于 S3 访问同时建立 Redshift 数据库连接默认连接 ID 为redshift_default。从源码看RedshiftToS3Operator实际依赖三个 HookS3Hook 负责获取 S3 凭据RedshiftSQLHook 负责执行 SQL而 RedshiftDataHook 则在采用 Redshift Data API 模式时被使用。三者共同完成了从 Redshift 读、往 S3 写的职责划分。RedshiftToS3Operator 是什么RedshiftToS3Operator的核心职责是把 Amazon Redshift 中的表数据通过 SQLUNLOAD命令导出为 S3 文件。它位于 redshift_to_s3.py类的 docstring 声明其行为为 Execute an UNLOAD command to s3 as a CSV with headers默认导出为带 CSV 选项的文件include_header可控制是否输出表头。在仓库的 DAG 拓扑中它常与以下操作符搭配使用构成完整的Redshift 与 S3 双向数据流转RedshiftCreateClusterOperator/RedshiftDeleteClusterOperator测试环境下的集群生命周期管理RedshiftDataOperator执行建表、插入、删表等辅助 SQLS3CreateBucketOperator/S3DeleteBucketOperatorS3 桶的创建与清理S3KeySensor验证 UNLOAD 产物是否已落地到 S3S3ToRedshiftOperator反向传输S3 → Redshift。完整示例见系统测试 example_redshift_s3_transfers.py。完整可运行的示例 DAG官方文档通过exampleinclude指令直接内嵌了系统测试 DAG 中的核心片段example_redshift_s3_transfers.py。最基础的用法如下transfer_redshift_to_s3 RedshiftToS3Operator( task_idtransfer_redshift_to_s3, redshift_data_api_kwargs{ database: DB_NAME, cluster_identifier: redshift_cluster_identifier, db_user: DB_LOGIN, wait_for_completion: True, }, s3_bucketbucket_name, s3_keyS3_KEY, schemaPUBLIC, tableREDSHIFT_TABLE, )对照该示例 DAG 的完整上下文example_redshift_s3_transfers.py其中使用的关键常量含义如下DB_LOGIN adminuser、DB_NAME devRedshift 数据库的登录用户名与库名S3_KEY s3_output_S3 目标 key 前缀REDSHIFT_TABLE test_table源表名示例中该表通过RedshiftDataOperator以如下 SQL 创建CREATE TABLE test_table ( fruit_id INTEGER, name VARCHAR NOT NULL, color VARCHAR NOT NULL ); INSERT INTO test_table VALUES ( 1, Banana, Yellow);执行产物验证由于table_as_file_name默认为TrueUNLOAD 产物会以{s3_key}/{table}_为前缀落地。示例 DAG 中通过S3KeySensor轮询等待产物出现example_redshift_s3_transfers.pycheck_if_key_exists S3KeySensor( task_idcheck_if_key_exists, bucket_namebucket_name, bucket_keyf{S3_KEY}/{REDSHIFT_TABLE}_0000_part_00, )这里test_table_0000_part_00正是 UNLOAD 命令默认的文件命名格式前缀_0000_part_00多分片时依次递增可作为断言导出成功的关键路径。核心参数详解源码级RedshiftToS3Operator构造函数的完整签名redshift_to_s3.py如下RedshiftToS3Operator( *, s3_bucket: str, s3_key: str, schema: str | None None, table: str | None None, select_query: str | None None, redshift_conn_id: str redshift_default, aws_conn_id: str | None NOTSET, verify: bool | str | None None, unload_options: list | None None, autocommit: bool False, include_header: bool False, parameters: Iterable | Mapping | None None, table_as_file_name: bool True, redshift_data_api_kwargs: dict | None None, **kwargs, )各参数语义整理如下参数类型默认值说明s3_bucketstr必填目标 S3 桶名s3_keystr必填目标 S3 key当table_as_file_nameFalse时必须包含完整文件名schemastr | NoneNoneRedshift 库中的 schema 名仅在使用table且未提供select_query时生效卸载临时表时不要传tablestr | NoneNoneRedshift 表名与schema配合生成默认查询select_querystr | NoneNone自定义 SELECT 查询优先级高于默认的SELECT * FROM schema.tableredshift_conn_idstrredshift_defaultRedshift 连接 IDSQL 模式aws_conn_idstr | NoneNOTSET回退aws_defaultS3/AWS 连接 ID若连接extras中含role_arn则走 IAM Role 授权verifybool | str | NoneNoneS3 连接 SSL 校验False跳过校验字符串为 CA 证书 bundle 路径unload_optionslist[]附加的 UNLOAD 选项列表如[HEADER, PARALLEL OFF, CSV]autocommitboolFalse为True时 UNLOAD 语句自动提交否则在连接关闭前提交include_headerboolFalse为True时在 S3 文件中输出列头自动追加HEADER选项parametersIterable | Mapping | NoneNone渲染 SQL 查询时使用的参数table_as_file_nameboolTrue为True时以表名作为 S3 文件名前缀{s3_key}/{table}_为保持向后兼容默认开启redshift_data_api_kwargsdict | NoneNone传入 Redshift Data API 模式时供 Hookexecute_query使用的参数禁止包含sql与parameters关于aws_conn_id的特别说明源码中对aws_conn_id的处理比较特殊redshift_to_s3.py当用户显式传入时标记conn_setTrue未传入时静默回退为aws_default。执行阶段redshift_to_s3.py会先尝试读取该 AWS 连接若连接的extras中存在role_arn则凭据块直接使用aws_iam_role{role_arn}形式将授权完全交由 Redshift 的 IAM Role 完成推荐的生产实践否则通过S3Hook.get_credentials()获取临时/永久凭据并用 build_credentials_block 拼接为aws_access_key_id...;aws_secret_access_key...;token...格式——若凭据来自 STS带 token会自动把token一并写入凭据块。两种执行模式SQL 连接 vs Redshift Data API该操作符支持两种向 Redshift 提交 UNLOAD 的方式通过是否设置redshift_data_api_kwargs自动切换源码属性 use_redshift_data。模式一传统 SQL 连接默认不传redshift_data_api_kwargs时操作符实例化RedshiftSQLHook并通过hook.run(unload_query, autocommit, parameters...)执行redshift_to_s3.py。此模式要求 Airflow 能通过redshift_conn_id直连 Redshift 数据库如使用 psycopg2 驱动。模式二Redshift Data API免直连传入redshift_data_api_kwargs时操作符改用RedshiftDataHook的execute_query方法异步提交语句并轮询结果redshift_to_s3.py。可用的键值参见 RedshiftDataHook.execute_query 的签名常用项包括database目标数据库名cluster_identifierRedshift 集群标识符db_user数据库用户名secret_arn用于数据库访问的 Secrets Manager 密钥 ARNstatement_nameSQL 语句名称wait_for_completion是否等待执行完成默认Truepoll_interval轮询间隔秒数默认10workgroup_nameRedshift Serverless 工作组名与cluster_identifier互斥session_id要复用的查询会话 IDUUID4 格式session_keep_alive_seconds会话在查询结束后保持存活的秒数上限 24 小时。注意sql与parameters两个键会被操作符明确拒绝源码 redshift_to_s3.py因为 SQL 由操作符内部构建。单元测试 test_redshift_to_s3.py 专门验证了传入这两个非法键会抛出AirflowException。临时表与会话复用系统测试中还有一个值得借鉴的模式为临时表导出数据并复用已建立的 Data API 会话example_redshift_s3_transfers.pycreate_tmp_table RedshiftDataOperator( task_idcreate_tmp_table, cluster_identifierredshift_cluster_identifier, databaseDB_NAME, db_userDB_LOGIN, sql_create_table(REDSHIFT_TMP_TABLE, is_tempTrue) _insert_data(REDSHIFT_TMP_TABLE), wait_for_completionTrue, session_keep_alive_seconds600, ) transfer_redshift_to_s3_reuse_session RedshiftToS3Operator( task_idtransfer_redshift_to_s3_reuse_session, redshift_data_api_kwargs{ wait_for_completion: True, session_id: {{ task_instance.xcom_pull(task_idscreate_tmp_table, keysession_id) }}, }, s3_bucketbucket_name, s3_keyS3_KEY_3, tableREDSHIFT_TMP_TABLE, )这里有两个要点卸载临时表时不要传schema。源码default_select_query属性redshift_to_s3.py中有 schema 时生成SELECT * FROM {schema}.{table}无 schema 时直接生成SELECT * FROM {table}——临时表不在用户 schema 下传 schema 会导致查询失败。Data API 会话可通过 XCom 在任务间传递建表任务开启session_keep_alive_seconds保留会话卸载任务通过session_id复用同一会话从而让临时表在同一会话内可见。RedshiftDataHook.execute_query会校验session_id必须是合法 UUID4redshift_data.py。底层执行原理UNLOAD 查询是如何构建的操作符在execute阶段依次完成以下步骤redshift_to_s3.py确定 S3 key若设置了table且table_as_file_nameTrue将s3_key重写为f{s3_key}/{table}_redshift_to_s3.py确定查询select_query优先否则回退到default_select_querySELECT * FROM ...两者都缺失时抛出ValueError(Please specify either a table orselect_queryto fetch the data.)补充 HEADERinclude_headerTrue且unload_options中尚无HEADER时自动追加构建凭据块按上文role_arn/AKSK 分支生成拼装最终 SQL_build_unload_queryredshift_to_s3.py生成如下模板UNLOAD ($${select_query}$$) TO s3://{s3_bucket}/{s3_key} credentials {credentials_block} {unload_options};其中$$...$$是 PostgreSQL 的美元引号定界符用于安全包裹可能包含单引号的 SELECT 语句_build_unload_query还会先用正则re.sub(r(.?), r\1, select_query)把转义过的单引号还原redshift_to_s3.py。单元测试 test_redshift_to_s3.py 通过多组参数化用例验证了含单引号的自定义查询在卸载前后保持一致。执行并记录日志分别输出 Executing UNLOAD command... 与 UNLOAD command complete... 两条日志。unload_options会被逐项以制表符连接后插入 SQL 模板因此你可传入任意合法的 UNLOAD 选项例如unload_options[ CSV, HEADER, PARALLEL OFF, GZIP, MANIFEST VERBOSE, ]注意UNLOAD 产物默认按数据分片输出多个文件如需单文件可加PARALLEL OFF如需清单文件可加MANIFEST。自定义 SELECT 查询与优先级规则select_query让你不必整表导出而是按需投影、过滤甚至多表 JOIN。其优先级规则被单元测试明确锁定test_redshift_to_s3.py同时传select_query、table、schema时select_query生效table/schema被忽略只传table/schema时使用默认查询SELECT * FROM schema.table只传select_query而不传table时table_as_file_name不再改变 keyexpected_s3_key保持原样。典型用法RedshiftToS3Operator( task_idunload_filtered, s3_bucketmy-bucket, s3_keyanalytics/orders, table_as_file_nameFalse, select_query( SELECT customer_id, order_date, total_amount FROM analytics.orders WHERE order_date 2025-01-01 ), unload_options[CSV, HEADER, PARALLEL OFF], include_headerTrue, )OpenLineage 血缘支持该操作符实现了get_openlineage_facets_on_completeredshift_to_s3.py可为每次导出生成输入/输出数据集的血缘信息输出数据集s3://{s3_bucket}/{s3_key}输入数据集redshift://{authority}命名空间下的{database}.{schema}.{table}使用默认SELECT *时输出文件会继承输入表的 schema facet并生成 IDENTITY 类型的columnLineage列级一一对应使用自定义select_query时会通过 SQL 解析器提取查询中的所有输入表支持 JOIN 多表并为每张表附加 schema/documentation facet解析失败时输出ExtractionErrorRunFacet。单元测试验证了 SQL 模式与 Data API 模式产出的血缘完全一致test_redshift_to_s3.py自定义查询场景下 JOIN 的两张表schema1.customers、schema2.orders均被正确识别为输入数据集test_redshift_to_s3.py。若 Airflow 实例启用了 OpenLineage这些信息会自动上报到血缘后端。参考文档官方操作符文档本文所依据的主文档redshift_to_s3.rst操作符源码redshift_to_s3.py完整系统测试 DAGexample_redshift_s3_transfers.py单元测试含 SQL 构建与 OpenLineage 断言test_redshift_to_s3.pyRedshift Data API Hookredshift_data.py凭据块构建工具redshift.pyAWS 官方对 UNLOAD 命令的完整语义选项、权限、命名规则可查阅 Amazon Redshift 数据库开发人员指南中的 UNLOAD 章节boto3 库中redshift与s3服务的客户端文档分别对应 Data API 与 S3 交互的底层接口。将上述仓库内的示例与源码结合使用即可在生产 DAG 中稳定落地Redshift 表 → S3 文件的批量导出任务。【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
