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 交互:
| 组件 | 模块路径 | 职责 |
|---|---|---|
| Hook | hooks/glacier.py | 封装boto3.client("glacier"),提供发起库存检索任务、获取任务结果、查询任务状态三个底层方法 |
| Operator | operators/glacier.py | 提供GlacierCreateJobOperator(发起任务)与GlacierUploadArchiveOperator(上传归档) |
| Sensor | sensors/glacier.py | 提供GlacierJobOperationSensor,轮询任务直到进入终态 |
| Transfer | transfers/glacier_to_gcs.py | 将 Glacier 任务结果搬运到 Google Cloud Storage(属于跨云传输的延伸用法) |
从源码结构看,hooks/glacier.py 中的GlacierHook是AwsBaseHook的子类,构造时固定传入client_type="glacier",也就是说所有上层的 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(vaultName=vault_name) # 清理 Vault(触发规则为 ALL_DONE,保证测试结束一定执行) boto3.client("glacier").delete_vault(vaultName=vault_name)2.3 配置 AWS 连接
所有 Glacier 组件默认使用连接 ID 为aws_default的 Amazon Web Services Connection,详细配置方式见 providers/amazon/docs/connections/aws.rst。要点包括:
- 默认连接 ID:
aws_default。若运行环境中${HOME}/.aws/存在凭据文件且连接的用户名/密码字段为空,会自动读取其中的凭据; - 凭据来源:支持 boto3 官方凭据链(环境变量、IAM Profile 等),也可在连接中直接指定 Access Key / Secret Key,或通过 Extra 中的
role_arn进行角色扮演(STS assume role); - Region:新版 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_id=None,而不是留空,以避免日志告警。
在 connections/aws.rst 中还可以看到用 Python 代码构造连接并生成AIRFLOW_CONN_*环境变量的示例:
from airflow.models.connection import Connection conn = Connection( conn_id="sample_aws_connection", conn_type="aws", login="YOUR_AWS_ACCESS_KEY_ID", password="YOUR_AWS_SECRET_ACCESS_KEY", extra={"region_name": "eu-central-1"}, ) env_key = f"AIRFLOW_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_name
- AWS 区域名;
- 设为
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_name | str | 执行任务的 Glacier Vault 名称(必填) |
aws_conn_id | str | None | AWS 连接 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.py):
create_glacier_job = GlacierCreateJobOperator(task_id="create_glacier_job", vault_name=vault_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_name=VAULT_NAME调用一次retrieve_inventory。
4.2GlacierUploadArchiveOperator:向 Vault 上传归档
功能:向指定的 Glacier Vault 添加(上传)一个归档(archive)。
核心参数:
| 参数 | 类型 | 说明 |
|---|---|---|
vault_name | str | Vault 名称(必填) |
body | bytes 或可 seek 的文件对象 | 要上传的数据(必填) |
checksum | str | None | 数据的 SHA256 树哈希;不提供时由 AWS 自动计算填充 |
archive_description | str | None | 归档的描述信息 |
account_id | str | None | 拥有 Vault 的 AWS 账户 ID;默认使用签名请求所用凭据对应的账户 |
aws_conn_id | str | None | AWS 连接 ID,默认aws_default |
典型用法(来自系统测试 DAG 的[START howto_operator_glacier_upload_archive]片段,见 example_glacier_to_gcs.py):
upload_archive_to_glacier = GlacierUploadArchiveOperator( task_id="upload_data_to_glacier", vault_name=vault_name, body=b"Test Data" )底层调用链:execute()直接调用self.hook.conn.upload_archive(...),把accountId、vaultName、archiveDescription、body、checksum原样透传给 boto3(见 operators/glacier.py)。单元测试 test_glacier.py 精确断言了upload_archive的调用参数,包括accountId=None、checksum=None的默认行为。
注意:该 Operator 的
template_fields同样包含vault_name(见 operators/glacier.py),其余参数如body不支持模板渲染。
五、Sensor:等待 Glacier 任务进入终态
5.1GlacierJobOperationSensor:轮询任务状态
功能:等待某个 Glacier 任务的状态进入终态。由于库存检索任务是异步的,发起任务后必须用该 Sensor 阻塞等待,直到describe_job返回的状态码为Succeeded。
核心参数:
| 参数 | 类型 | 默认值 | 说明 |
|---|---|---|---|
vault_name | str | — | 任务所在的 Glacier Vault 名称(必填) |
job_id | str | — | retrieve_inventory()返回的任务 ID(必填) |
poke_interval | int | 60 * 20(1200 秒) | 两次轮询之间的等待秒数 |
mode | str | reschedule | 传感器运行模式,poke或reschedule |
aws_conn_id | str | None | aws_default | AWS 连接 ID |
模式说明(来自源码 docstring,见 sensors/glacier.py):
poke:传感器在整体执行期间占用一个 worker 槽位,在两次轮询之间睡眠。适用于预期运行时间短或需要短轮询间隔的场景;reschedule:条件未满足时释放 worker 槽位,稍后重新调度。适用于等待时间较长的场景,此时轮询间隔应大于一分钟,以免给调度器造成过大负载。
该 Sensor 默认采用reschedule模式且poke_interval长达 20 分钟,正是因为 Glacier 的库存检索任务通常需要数小时才能完成,长时间占用 worker 槽位并不经济。
典型用法(来自系统测试 DAG 的[START howto_sensor_glacier_job_operation]片段,见 example_glacier_to_gcs.py):
wait_for_operation_complete = GlacierJobOperationSensor( vault_name=vault_name, job_id=JOB_ID, task_id="wait_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}),使任务失败。
状态枚举定义在源码顶部的JobStatus(IN_PROGRESS = "InProgress"、SUCCEEDED = "Succeeded",见 sensors/glacier.py)。单元测试 test_glacier.py 覆盖了成功、处理中、空状态码、Failed状态四种分支,其中异常分支用pytest.raises(AirflowException, match=...)验证了失败信息。
六、延伸:Glacier 到 GCS 的跨云传输
虽然原文档正文聚焦于任务创建与状态等待,但仓库的系统测试 DAG 与源码中还包含一个完整的跨云传输组件GlacierToGCSOperator(transfers/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_id | str | None | aws_default | AWS 连接 ID |
gcp_conn_id | str | google_cloud_default | GCP 连接 ID |
vault_name | str | — | 源 Glacier Vault 名称 |
bucket_name | str | — | 目标 GCS bucket |
object_name | str | — | 目标 GCS object 名称 |
gzip | bool | —(必填) | 是否在上传时压缩文件数据 |
chunk_size | int | 1024 | 从 Glacier 下载的块大小(字节);若大于实际文件大小,整个文件将一次性下载 |
google_impersonation_chain | str | Sequence[str] | None | None | 可选的 GCP 服务账号模拟链 |
注意事项:源码 docstring 明确警告该 Operator 依赖内存使用,传输大文件可能表现不佳(见 transfers/glacier_to_gcs.py),在大文件场景下应谨慎评估。
用法示例(来自系统测试 DAG 的[START howto_transfer_glacier_to_gcs]片段,见 example_glacier_to_gcs.py):
transfer_archive_to_gcs = GlacierToGCSOperator( task_id="transfer_archive_to_gcs", vault_name=vault_name, bucket_name=gcs_bucket_name, object_name=gcs_object_name, gzip=False, # 如果 chunk size 大于实际文件大小,则整个文件会被一次性下载 chunk_size=1024, )单元测试 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), # 创建 Vault(TEST SETUP) create_glacier_job, # 发起库存检索任务 wait_for_operation_complete, # 等待任务进入 Succeeded upload_archive_to_glacier, # 上传归档 transfer_archive_to_gcs, # 搬运到 GCS delete_vault(vault_name), # 清理 Vault(TEARDOWN) )该 DAG 的关键设计点:
- Vault 命名隔离:Vault 名、GCS bucket 名、object 名都以测试上下文生成的
env_id为前缀,避免多环境互相干扰; - 异步任务编排:
GlacierCreateJobOperator只负责发起任务,jobId通过 XCom 模板传递给GlacierJobOperationSensor等待完成;由于retrieve_inventory是异步任务,实际耗时可达数小时,因此 Sensor 默认mode="reschedule"、poke_interval=1200秒; - 触发规则:清理任务
delete_vault使用TriggerRule.ALL_DONE,保证无论测试主体成功还是失败都会执行清理(见 example_glacier_to_gcs.py); - Airflow 版本兼容:DAG 同时兼容 Airflow 2 与 3(通过
AIRFLOW_V_3_0_PLUS条件导入,见 example_glacier_to_gcs.py)。
八、测试与验证:如何确认组件行为
仓库为 Glacier 相关组件提供了系统测试与单元测试两层验证,可作为自行编写测试的参照:
- 系统测试 DAG:providers/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),仅供参考