多彩编程 多彩编程MZPH · CODE BLOG
ARTICLE DETAIL

文章详情

深耕前端与后端开发技术的一线实战笔记与踩坑复盘。

Apache Airflow 对接 Amazon S3 Glacier:归档任务编排与跨云传输实战指南

Apache Airflow 对接 Amazon S3 Glacier:归档任务编排与跨云传输实战指南 Apache Airflow 对接 Amazon S3 Glacier归档任务编排与跨云传输实战指南【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow导读Amazon S3 Glacier 是 AWS 面向数据归档与长期备份场景提供的高持久、低成本存储层级而 Apache Airflow 的 Amazon Provider 提供了与之对应的 Operator 与 Sensor帮助你在 DAG 中程序化地发起库存检索inventory-retrieval任务、向 Vault 上传归档archive并在任务完成前持续等待。本文以 providers/amazon/docs/operators/s3/glacier.rst 为核心骨架结合仓库中的 Hook、Operator、Sensor 源码与系统/单元测试完整讲解前置准备、通用参数、三个核心组件的用法以及 Glacier 到 GCS 的跨云数据搬运链路读完即可在真实 DAG 中落地一套“发起任务 → 等待完成 → 上传归档 → 搬运结果”的完整流程。一、Glacier 在 Airflow 生态中的定位Amazon Glacier 是一种安全、持久且成本极低的 S3 存储类别适合存放不常访问但必须长期保留的归档数据。在 Airflow 中Amazon Provider 通过以下模块与 Glacier 交互组件模块路径职责Hookhooks/glacier.py封装boto3.client(glacier)提供发起库存检索任务、获取任务结果、查询任务状态三个底层方法Operatoroperators/glacier.py提供GlacierCreateJobOperator发起任务与GlacierUploadArchiveOperator上传归档Sensorsensors/glacier.py提供GlacierJobOperationSensor轮询任务直到进入终态Transfertransfers/glacier_to_gcs.py将 Glacier 任务结果搬运到 Google Cloud Storage属于跨云传输的延伸用法从源码结构看hooks/glacier.py 中的GlacierHook是AwsBaseHook的子类构造时固定传入client_typeglacier也就是说所有上层的 Operator 与 Sensor 最终都经由这一个 Hook 拿到 boto3 Glacier 客户端凭证管理与区域解析逻辑统一由 Amazon Provider 的 base hook 承担。二、前置准备安装与连接配置2.1 安装 Amazon Provider使用上述组件前需要安装 Amazon Provider 及其依赖pip install apache-airflow[amazon]安装完成后即可在代码中导入from airflow.providers.amazon.aws.operators.glacier import ( GlacierCreateJobOperator, GlacierUploadArchiveOperator, ) from airflow.providers.amazon.aws.sensors.glacier import GlacierJobOperationSensor2.2 创建必要的 AWS 资源在运行 DAG 之前需要通过 AWS Console 或 AWS CLI 预先创建 Glacier Vault归档库等资源。系统测试 DAG example_glacier_to_gcs.py 展示了用 boto3 直接创建与清理 Vault 的做法可供参考import boto3 # 创建 Vault boto3.client(glacier).create_vault(vaultNamevault_name) # 清理 Vault触发规则为 ALL_DONE保证测试结束一定执行 boto3.client(glacier).delete_vault(vaultNamevault_name)2.3 配置 AWS 连接所有 Glacier 组件默认使用连接 ID 为aws_default的 Amazon Web Services Connection详细配置方式见 providers/amazon/docs/connections/aws.rst。要点包括默认连接 IDaws_default。若运行环境中${HOME}/.aws/存在凭据文件且连接的用户名/密码字段为空会自动读取其中的凭据凭据来源支持 boto3 官方凭据链环境变量、IAM Profile 等也可在连接中直接指定 Access Key / Secret Key或通过 Extra 中的role_arn进行角色扮演STS assume roleRegion新版 Provider 安装后不再默认写入{region_name: us-east-1}需要手动在连接界面配置或通过AWS_DEFAULT_REGION环境变量指定Extra 常用字段region_name、profile_name、role_arn、assume_role_method、config_kwargs用于构造 botocore Config、verify是否校验 SSL 证书、endpoint_url等跳过连接查找如果希望完全走默认的 boto3 凭据策略应在 Operator 中显式传入aws_conn_idNone而不是留空以避免日志告警。在 connections/aws.rst 中还可以看到用 Python 代码构造连接并生成AIRFLOW_CONN_*环境变量的示例from airflow.models.connection import Connection conn Connection( conn_idsample_aws_connection, conn_typeaws, loginYOUR_AWS_ACCESS_KEY_ID, passwordYOUR_AWS_SECRET_ACCESS_KEY, extra{region_name: eu-central-1}, ) env_key fAIRFLOW_CONN_{conn.conn_id.upper()} os.environ[env_key] conn.get_uri()三、通用参数所有 Glacier 组件共享的配置项来自 providers/amazon/docs/_partials/generic_parameters.rst 的通用参数适用于本节提到的全部 Operator 与 Sensor经由AwsBaseOperator/AwsBaseSensor继承。在单元测试 test_glacier.py 中可以看到这些参数如何透传到底层 Hook。3.1aws_conn_id引用 Amazon Web Services Connection 的连接 ID设为None时不进行连接查找直接使用默认 boto3 行为环境变量凭据、IAM Profile 等默认值aws_default。3.2region_nameAWS 区域名设为None或省略时使用 AWS 连接 Extra 参数中的region_name显式指定则覆盖连接中的区域值默认值None。3.3verify是否校验 SSL 证书False不校验 SSL 证书证书 bundle 路径使用指定的 CA 证书包当不想用 botocore 默认证书包时设为None或省略时使用连接 Extra 中的verify默认值None。3.4botocore_config传入字典用于构造botocore.config.Config可配置重试策略、超时、限流规避等设为None或省略时使用连接 Extra 中的config_kwargs注意传入空字典{}会覆盖连接中的 botocore 配置官方文档示例{ signature_version: unsigned, s3: { us_east_1_regional_endpoint: True, }, retries: { mode: standard, max_attempts: 10, }, connect_timeout: 300, read_timeout: 300, tcp_keepalive: True, }单元测试中的用例botocore_config{read_timeout: 42}验证了read_timeout会被正确写入op.hook._config见 test_glacier.py说明该字典最终用于构造真实的 botocore Config 对象。四、Operator创建 Glacier 任务4.1GlacierCreateJobOperator发起库存检索任务功能向指定的 Glacier Vault 发起一个 inventory-retrieval库存检索任务。Glacier 的库存检索是异步的任务在后台执行本 Operator 只是发起任务并立即返回。返回结果返回与已发起任务相关的信息字典其中包含下游任务必需的jobId。例如GlacierHook.retrieve_inventory直接返回initiate_job的响应其中response[jobId]即任务 ID见 hooks/glacier.py。核心参数参数类型说明vault_namestr执行任务的 Glacier Vault 名称必填aws_conn_idstr | NoneAWS 连接 ID默认aws_default模板字段vault_name支持 Jinja 模板渲染源码中template_fields通过aws_template_fields(vault_name)声明见 operators/glacier.py这意味着你可以用{{ ... }}表达式动态指定 Vault 名称。典型用法来自系统测试 DAG 的[START howto_operator_glacier_create_job]片段见 example_glacier_to_gcs.pycreate_glacier_job GlacierCreateJobOperator(task_idcreate_glacier_job, vault_namevault_name) JOB_ID {{ task_instance.xcom_pull(create_glacier_job)[jobId] }}第二行通过 XCom 从上游任务的结果字典中取出jobId并以模板方式传给下游的 Sensor —— 这是 Glacier 异步任务编排的惯用法。底层调用链execute()调用self.hook.retrieve_inventory(vault_name...)该方法内部构造jobParameters {Type: inventory-retrieval}并调用 boto3 的initiate_job见 operators/glacier.py。单元测试 test_glacier.py 验证了execute会以vault_nameVAULT_NAME调用一次retrieve_inventory。4.2GlacierUploadArchiveOperator向 Vault 上传归档功能向指定的 Glacier Vault 添加上传一个归档archive。核心参数参数类型说明vault_namestrVault 名称必填bodybytes 或可 seek 的文件对象要上传的数据必填checksumstr | None数据的 SHA256 树哈希不提供时由 AWS 自动计算填充archive_descriptionstr | None归档的描述信息account_idstr | None拥有 Vault 的 AWS 账户 ID默认使用签名请求所用凭据对应的账户aws_conn_idstr | NoneAWS 连接 ID默认aws_default典型用法来自系统测试 DAG 的[START howto_operator_glacier_upload_archive]片段见 example_glacier_to_gcs.pyupload_archive_to_glacier GlacierUploadArchiveOperator( task_idupload_data_to_glacier, vault_namevault_name, bodybTest Data )底层调用链execute()直接调用self.hook.conn.upload_archive(...)把accountId、vaultName、archiveDescription、body、checksum原样透传给 boto3见 operators/glacier.py。单元测试 test_glacier.py 精确断言了upload_archive的调用参数包括accountIdNone、checksumNone的默认行为。注意该 Operator 的template_fields同样包含vault_name见 operators/glacier.py其余参数如body不支持模板渲染。五、Sensor等待 Glacier 任务进入终态5.1GlacierJobOperationSensor轮询任务状态功能等待某个 Glacier 任务的状态进入终态。由于库存检索任务是异步的发起任务后必须用该 Sensor 阻塞等待直到describe_job返回的状态码为Succeeded。核心参数参数类型默认值说明vault_namestr—任务所在的 Glacier Vault 名称必填job_idstr—retrieve_inventory()返回的任务 ID必填poke_intervalint60 * 201200 秒两次轮询之间的等待秒数modestrreschedule传感器运行模式poke或rescheduleaws_conn_idstr | Noneaws_defaultAWS 连接 ID模式说明来自源码 docstring见 sensors/glacier.pypoke传感器在整体执行期间占用一个 worker 槽位在两次轮询之间睡眠。适用于预期运行时间短或需要短轮询间隔的场景reschedule条件未满足时释放 worker 槽位稍后重新调度。适用于等待时间较长的场景此时轮询间隔应大于一分钟以免给调度器造成过大负载。该 Sensor 默认采用reschedule模式且poke_interval长达 20 分钟正是因为 Glacier 的库存检索任务通常需要数小时才能完成长时间占用 worker 槽位并不经济。典型用法来自系统测试 DAG 的[START howto_sensor_glacier_job_operation]片段见 example_glacier_to_gcs.pywait_for_operation_complete GlacierJobOperationSensor( vault_namevault_name, job_idJOB_ID, task_idwait_for_operation_complete, )其中JOB_ID即上文通过 XCom 模板取得的jobId。轮询逻辑poke()方法见 sensors/glacier.py调用self.hook.describe_job(vault_name..., job_id...)查询任务状态若StatusCode Succeeded记录成功日志返回True传感器结束若StatusCode InProgress记录处理中日志返回False继续轮询其他任何状态码抛出AirflowException消息格式为Sensor failed. Job status: {Action}, code status: {StatusCode}使任务失败。状态枚举定义在源码顶部的JobStatusIN_PROGRESS InProgress、SUCCEEDED Succeeded见 sensors/glacier.py。单元测试 test_glacier.py 覆盖了成功、处理中、空状态码、Failed状态四种分支其中异常分支用pytest.raises(AirflowException, match...)验证了失败信息。六、延伸Glacier 到 GCS 的跨云传输虽然原文档正文聚焦于任务创建与状态等待但仓库的系统测试 DAG 与源码中还包含一个完整的跨云传输组件GlacierToGCSOperatortransfers/glacier_to_gcs.py它常与上文组件串联成端到端数据链路值得一并掌握。执行流程execute()方法见 transfers/glacier_to_gcs.py构造GlacierHook与GCSHook分别由aws_conn_id与gcp_conn_id指定连接调用retrieve_inventory发起库存检索任务并取得jobId调用retrieve_inventory_results获取任务输出get_job_output拿到 StreamingBody使用tempfile.NamedTemporaryFile()创建临时文件按chunk_size分块读取 Glacier 数据并写入通过gcs_hook.upload(...)将临时文件上传到 GCS 的指定 bucket 与 object返回gs://{bucket}/{object}格式的 GCS 对象地址。核心参数参数类型默认值说明aws_conn_idstr | Noneaws_defaultAWS 连接 IDgcp_conn_idstrgoogle_cloud_defaultGCP 连接 IDvault_namestr—源 Glacier Vault 名称bucket_namestr—目标 GCS bucketobject_namestr—目标 GCS object 名称gzipbool—必填是否在上传时压缩文件数据chunk_sizeint1024从 Glacier 下载的块大小字节若大于实际文件大小整个文件将一次性下载google_impersonation_chainstr | Sequence[str] | NoneNone可选的 GCP 服务账号模拟链注意事项源码 docstring 明确警告该 Operator 依赖内存使用传输大文件可能表现不佳见 transfers/glacier_to_gcs.py在大文件场景下应谨慎评估。用法示例来自系统测试 DAG 的[START howto_transfer_glacier_to_gcs]片段见 example_glacier_to_gcs.pytransfer_archive_to_gcs GlacierToGCSOperator( task_idtransfer_archive_to_gcs, vault_namevault_name, bucket_namegcs_bucket_name, object_namegcs_object_name, gzipFalse, # 如果 chunk size 大于实际文件大小则整个文件会被一次性下载 chunk_size1024, )单元测试 test_glacier_to_gcs.py 验证了执行过程中两个 Hook 的构造参数与调用序列GlacierHook(aws_conn_id...)、retrieve_inventory→retrieve_inventory_results(job_id...)→GCSHook(gcp_conn_id..., impersonation_chain...)→upload(...)与源码逻辑一一对应。七、端到端 DAG 示例完整的归档编排流程仓库中的系统测试 DAG example_glacier_to_gcs.py 把上述全部组件串成了一条完整的端到端链路可直接作为编写生产 DAG 的模板。其任务依赖关系为chain( test_context, # 测试上下文环境 ID create_vault(vault_name), # 创建 VaultTEST SETUP create_glacier_job, # 发起库存检索任务 wait_for_operation_complete, # 等待任务进入 Succeeded upload_archive_to_glacier, # 上传归档 transfer_archive_to_gcs, # 搬运到 GCS delete_vault(vault_name), # 清理 VaultTEARDOWN )该 DAG 的关键设计点Vault 命名隔离Vault 名、GCS bucket 名、object 名都以测试上下文生成的env_id为前缀避免多环境互相干扰异步任务编排GlacierCreateJobOperator只负责发起任务jobId通过 XCom 模板传递给GlacierJobOperationSensor等待完成由于retrieve_inventory是异步任务实际耗时可达数小时因此 Sensor 默认modereschedule、poke_interval1200秒触发规则清理任务delete_vault使用TriggerRule.ALL_DONE保证无论测试主体成功还是失败都会执行清理见 example_glacier_to_gcs.pyAirflow 版本兼容DAG 同时兼容 Airflow 2 与 3通过AIRFLOW_V_3_0_PLUS条件导入见 example_glacier_to_gcs.py。八、测试与验证如何确认组件行为仓库为 Glacier 相关组件提供了系统测试与单元测试两层验证可作为自行编写测试的参照系统测试 DAGproviders/amazon/tests/system/amazon/aws/example_glacier_to_gcs.py —— 真实调用 AWS 与 GCP覆盖创建 Vault、发起任务、等待、上传、传输、清理全流程文件末尾的get_test_run(dag)使其可直接通过 pytest 运行Operator 单元测试providers/amazon/tests/unit/amazon/aws/operators/test_glacier.py —— 通过 mock 验证retrieve_inventory与upload_archive的调用参数、通用参数透传aws_conn_id/region_name/verify/botocore_config以及模板字段合法性Sensor 单元测试providers/amazon/tests/unit/amazon/aws/sensors/test_glacier.py —— 覆盖poke()的成功、进行中、未知状态、失败四种分支Transfer 单元测试providers/amazon/tests/unit/amazon/aws/transfers/test_glacier_to_gcs.py —— 验证两个 Hook 的构造与调用序列。在本地验证时可以先用 mock 的 boto3 客户端跑单元测试确认编排逻辑再在具备 AWS/GCP 凭据的环境中运行系统测试 DAG 做端到端验证连接可用性可通过 Amazon Provider 的test_connection检查注意连接测试依赖 STSGetCallerIdentity仅能验证凭据有效性无法验证对 Glacier 的具体访问权限详见 connections/aws.rst。结语本文以 glacier.rst 为骨架完整覆盖了 Amazon S3 Glacier 集成所需的前置准备、通用参数、任务创建 Operator、上传 Operator、状态等待 Sensor并延伸讲解了 Glacier 到 GCS 的跨云传输与端到端 DAG 编排。核心要点可以归纳为用GlacierCreateJobOperator发起异步库存检索任务用 XCom 传递jobId用默认 reschedule 模式的GlacierJobOperationSensor等待终态再用GlacierUploadArchiveOperator写入归档、GlacierToGCSOperator完成跨云搬运。理解这些组件的底层 boto3 调用链与默认参数尤其是 Sensor 的 20 分钟轮询间隔与 reschedule 模式有助于为真实的归档与长期备份场景设计出高效、可观测、可运维的 Airflow 工作流。【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表