news 2026/9/13 19:05:34

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

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
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 中的GlacierHookAwsBaseHook的子类,构造时固定传入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 GlacierJobOperationSensor

2.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。要点包括:

  • 默认连接 IDaws_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_nameprofile_namerole_arnassume_role_methodconfig_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_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.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_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.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(...),把accountIdvaultNamearchiveDescriptionbodychecksum原样透传给 boto3(见 operators/glacier.py)。单元测试 test_glacier.py 精确断言了upload_archive的调用参数,包括accountId=Nonechecksum=None的默认行为。

注意:该 Operator 的template_fields同样包含vault_name(见 operators/glacier.py),其余参数如body不支持模板渲染。


五、Sensor:等待 Glacier 任务进入终态

5.1GlacierJobOperationSensor:轮询任务状态

功能:等待某个 Glacier 任务的状态进入终态。由于库存检索任务是异步的,发起任务后必须用该 Sensor 阻塞等待,直到describe_job返回的状态码为Succeeded

核心参数

参数类型默认值说明
vault_namestr任务所在的 Glacier Vault 名称(必填)
job_idstrretrieve_inventory()返回的任务 ID(必填)
poke_intervalint60 * 20(1200 秒)两次轮询之间的等待秒数
modestrreschedule传感器运行模式,pokereschedule
aws_conn_idstr | Noneaws_defaultAWS 连接 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):

  1. 调用self.hook.describe_job(vault_name=..., job_id=...)查询任务状态;
  2. StatusCode == "Succeeded":记录成功日志,返回True,传感器结束;
  3. StatusCode == "InProgress":记录处理中日志,返回False,继续轮询;
  4. 其他任何状态码:抛出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 与源码中还包含一个完整的跨云传输组件GlacierToGCSOperator(transfers/glacier_to_gcs.py),它常与上文组件串联成端到端数据链路,值得一并掌握。

执行流程execute()方法,见 transfers/glacier_to_gcs.py):

  1. 构造GlacierHookGCSHook(分别由aws_conn_idgcp_conn_id指定连接);
  2. 调用retrieve_inventory发起库存检索任务并取得jobId
  3. 调用retrieve_inventory_results获取任务输出(get_job_output),拿到 StreamingBody;
  4. 使用tempfile.NamedTemporaryFile()创建临时文件,按chunk_size分块读取 Glacier 数据并写入;
  5. 通过gcs_hook.upload(...)将临时文件上传到 GCS 的指定 bucket 与 object;
  6. 返回gs://{bucket}/{object}格式的 GCS 对象地址。

核心参数

参数类型默认值说明
aws_conn_idstr | Noneaws_defaultAWS 连接 ID
gcp_conn_idstrgoogle_cloud_defaultGCP 连接 ID
vault_namestr源 Glacier Vault 名称
bucket_namestr目标 GCS bucket
object_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.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_inventoryretrieve_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_inventoryupload_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),仅供参考

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/13 19:04:50

iii Worker Registry 完全指南:浏览、安装与管理可插拔 Worker

iii Worker Registry 完全指南:浏览、安装与管理可插拔 Worker 【免费下载链接】iii Effortlessly compose, extend, and observe every service in real-time for the first time ever. 项目地址: https://gitcode.com/GitHub_Trending/mo/iii 导读 本指南…

作者头像 李华
网站建设 2026/9/13 19:03:45

PowerPC Linux PCI 总线 EEH 错误恢复机制深度解析

PowerPC Linux PCI 总线 EEH 错误恢复机制深度解析 【免费下载链接】linux Linux kernel source tree 项目地址: https://gitcode.com/GitHub_Trending/li/linux 导读 本文基于 Linux 内核源码树中的 Documentation/arch/powerpc/eeh-pci-error-recovery.rst&#xff0…

作者头像 李华