十年匠心定制 · 商业建站与技术教学双线并行 咨询热线:400-886-1026 service@lmnt.cn
ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

从 Amazon S3 Glacier 到 Google Cloud Storage 的跨云迁移:Apache Airflow GlacierToGCSOperator 实战指南

从 Amazon S3 Glacier 到 Google Cloud Storage 的跨云迁移:Apache Airflow GlacierToGCSOperator 实战指南 从 Amazon S3 Glacier 到 Google Cloud Storage 的跨云迁移Apache Airflow GlacierToGCSOperator 实战指南【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflowAmazon S3 Glacier 是 AWS 提供的一种安全、持久且成本极低的云存储级别专为数据归档和长期备份设计。当企业需要把历史归档数据从 AWS 迁移到 Google Cloud 生态例如切换到 GCP 统一存储、让数据进入 BigQuery 分析链路时Apache Airflow 的GlacierToGCSOperator提供了一条开箱即用的自动化路径。本文将围绕 glacier_to_gcs.rst 展开结合算子源码、Hook 实现与仓库内测试用例完整讲解该跨云迁移任务的前置条件、参数语义、执行原理与内存注意事项读完即可在真实 DAG 中落地 Glacier → GCS 的数据搬运。为什么需要 Glacier 到 GCS 的跨云迁移Glacier 存储类的核心优势是极低的存储成本和面向长期归档的持久性设计适合冷数据、合规归档与容灾备份。但在多云混合架构中数据往往需要跨云流动业务重心迁移到 Google Cloud需要把历史归档一并迁入 GCS归档数据需要被 GCP 侧的 Dataflow、BigQuery、Vertex AI 等生态消费希望在统一的数据湖GCS中集中管理所有冷热数据。GlacierToGCSOperator正是为这一场景设计的专用 transfer 算子它位于 Amazon provider 的 transfers 目录中承担从 Amazon Glacier vault 取数据 → 落到 Google Cloud Storage 桶的完整任务。前置条件按 prerequisite_tasks.rst 的要求使用该算子前需要完成以下准备准备 AWS 资源在AWS Console或AWS CLI中创建好源 Glacier vault以及可选的待归档文件并保证运行 Airflow 的账号拥有对该 vault 执行initiate-job、get-job-output等操作的权限。安装 Amazon providerpip install apache-airflow[amazon]详细的安装说明参见仓库中airflow-core/docs/installation/目录下的安装文档。配置 AWS Connection在 Airflow 中建立aws_default连接或自定义连接名提供访问 Glacier 所需的 AWS 访问密钥与区域信息。需要特别指出的是从 算子源码 可以看到该算子同时依赖GlacierHookAmazon provider与GCSHookGoogle provider。因此除apache-airflow[amazon]外运行环境还需要安装 Google 相关依赖例如apache-airflow[google]并在 Airflow 中配置一个可用的 GCP 连接默认连接名为google_cloud_default否则任务会在执行阶段因找不到 Hook 依赖或连接而失败。GlacierToGCSOperator 核心参数详解算子定义位于 providers/amazon/src/airflow/providers/amazon/aws/transfers/glacier_to_gcs.py构造签名如下GlacierToGCSOperator( *, aws_conn_id: str | None aws_default, gcp_conn_id: str google_cloud_default, vault_name: str, bucket_name: str, object_name: str, gzip: bool, chunk_size: int 1024, google_impersonation_chain: str | Sequence[str] | None None, **kwargs, )参数类型默认值说明aws_conn_idstr \| Noneaws_default指向 AWS 连接用于创建GlacierHook访问 Glacier 服务gcp_conn_idstrgoogle_cloud_default指向 GCP 连接用于创建GCSHook上传对象vault_namestr必填执行任务的 Glacier vault 名称支持模板bucket_namestr必填目标 Google Cloud Storage 桶名称支持模板object_namestr必填上传到 GCS 桶中的对象名支持模板gzipbool必填是否在上传前对本地文件/文件数据进行 gzip 压缩chunk_sizeint1024从 Glacier vault 下载数据时的分块大小字节google_impersonation_chainstr \| Sequence[str] \| NoneNone可选的 Google 服务账号模拟链用于以短时凭证模拟目标账号执行 GCS 操作参数语义补充说明模板字段源码中声明了template_fields (vault_name, bucket_name, object_name)意味着这三个字段支持 Jinja 模板渲染可以在运行时通过{{ ti.xcom_pull(...) }}或{{ ds }}等上下文动态生成目标路径适合批量归档迁移场景。chunk_size的取值逻辑chunk_size是 Glacier 侧下载的分块字节数。从系统测试 DAG 的注释可以确认如果 chunk_size 大于实际文件大小则整个文件会被一次性下载反之数据将按指定块大小被分批读取。合理调小chunk_size可以降低单次内存峰值但会增加迭代开销。gzip与 GCS 上传该参数会原样透传给GCSHook.upload用于控制上传前是否对文件进行 gzip 压缩从而减少 GCS 侧存储占用与传输流量。google_impersonation_chain当以字符串传入时该账号必须授予发起账号Service Account Token CreatorIAM 角色以序列传入时列表中的相邻身份依次授予前一身份该角色最终以最后一个账号的身份发起请求该参数同样支持模板渲染。完整 DAG 示例最小可运行示例文档核心片段原文档通过exampleinclude从系统测试 DAG 中抽取了GlacierToGCSOperator的核心用法见 example_glacier_to_gcs.py 中howto_transfer_glacier_to_gcs标记段transfer_archive_to_gcs GlacierToGCSOperator( task_idtransfer_archive_to_gcs, vault_namevault_name, bucket_namegcs_bucket_name, object_namegcs_object_name, gzipFalse, # Override to match your needs # If chunk size is bigger than actual file size # then whole file will be downloaded chunk_size1024, )这是一个真正可复制、可运行的算子实例指定源 vault、目标桶与对象名后算子会在执行时自动完成 Glacier 清单检索、结果流式下载与 GCS 上传。生产级完整 DAG含 Job 创建与等待GlacierToGCSOperator只负责取数 上传但在真实场景中Glacier 的检索需要先创建检索任务并等待其完成。仓库中的系统测试 DAG 给出了完整的生产链路串联了 Glacier 生态的配套算子与传感器create_glacier_job GlacierCreateJobOperator(task_idcreate_glacier_job, vault_namevault_name) JOB_ID {{ task_instance.xcom_pull(create_glacier_job)[jobId] }} wait_for_operation_complete GlacierJobOperationSensor( vault_namevault_name, job_idJOB_ID, task_idwait_for_operation_complete, ) upload_archive_to_glacier GlacierUploadArchiveOperator( task_idupload_data_to_glacier, vault_namevault_name, bodybTest Data ) transfer_archive_to_gcs GlacierToGCSOperator( task_idtransfer_archive_to_gcs, vault_namevault_name, bucket_namegcs_bucket_name, object_namegcs_object_name, gzipFalse, chunk_size1024, ) chain( create_vault(vault_name), create_glacier_job, wait_for_operation_complete, upload_archive_to_glacier, transfer_archive_to_gcs, delete_vault(vault_name), )完整依赖关系为GlacierCreateJobOperator发起 inventory-retrieval 任务 →GlacierJobOperationSensor轮询任务完成其 job_id 通过 XCom 从上游拉取→GlacierUploadArchiveOperator写入测试数据 →GlacierToGCSOperator执行跨云迁移 → 清理任务删除测试 vault。这套编排展示了创建检索任务 → 等待完成 → 迁移的标准姿势值得在真实 DAG 中复用。执行原理源码级调用链解析GlacierToGCSOperator.execute()的核心实现非常精简但背后串联了两大云 SDK 的完整链路。逐步拆解如下对应 glacier_to_gcs.py第 1 步初始化双云 Hookglacier_hook GlacierHook(aws_conn_idself.aws_conn_id) gcs_hook GCSHook( gcp_conn_idself.gcp_conn_id, impersonation_chainself.impersonation_chain, )GlacierHook继承自AwsBaseHook构造时强制指定client_typeglacier见 hooks/glacier.py本质是对boto3.client(glacier)的薄封装GCSHook则负责与 Google Cloud Storage 交互。第 2 步发起清单检索任务job_id glacier_hook.retrieve_inventory(vault_nameself.vault_name)对应GlacierHook.retrieve_inventory()hooks/glacier.py其内部调用get_conn().initiate_job(vaultNamevault_name, jobParameters{Type: inventory-retrieval})向 Glacier 提交一个inventory-retrieval类型的异步任务返回的响应中包含jobId。第 3 步流式分块下载任务结果with tempfile.NamedTemporaryFile() as temp_file: glacier_data glacier_hook.retrieve_inventory_results( vault_nameself.vault_name, job_idjob_id[jobId] ) stream glacier_data[body] for chunk in stream.iter_chunks(chunk_sizeself.chunk_size): temp_file.write(chunk) temp_file.flush()retrieve_inventory_results()内部调用get_job_output(vaultName, jobId)hooks/glacier.py返回响应中的body是一个 botocore 的StreamingBody。算子不会一次性把整个结果载入内存而是调用iter_chunks(chunk_size...)按块迭代写入本地NamedTemporaryFile临时文件——这就是chunk_size参数真正发挥作用的位置。第 4 步上传到 GCS 并返回对象 URIgcs_hook.upload( bucket_nameself.bucket_name, object_nameself.object_name, filenametemp_file.name, gzipself.gzip, ) return fgs://{self.bucket_name}/{self.object_name}GCSHook.uploadproviders/google/src/airflow/providers/google/cloud/hooks/gcs.py支持filename/data两种数据来源本算子采用本地文件上传方式并将gzip参数透传。上传成功后算子返回gs://bucket/object形式的对象 URI 字符串可供下游任务通过 XCom 消费。注意虽然下载侧采用了流式分块但数据最终会先完整落到 worker 节点的临时文件中再由 GCSHook 整体上传因此本地磁盘占用与对象大小成正比内存与磁盘容量都是规划任务时需要考虑的约束。内存使用注意事项官方明确警告原文档与算子 docstring 都给出了同一条重要警告GlacierToGCSOperator的可用性依赖于 worker 的内存容量传输大文件可能导致 worker 主机内存耗尽、任务失败。实际使用建议对于超大归档文件建议评估拆分策略例如先在 AWS 侧按对象切分再逐个迁移避免单任务处理过大的 Glacier 检索结果合理设置chunk_size以控制下载侧的流式读取粒度——chunk_size越小单次读入内存的数据量越少但迭代次数增多为运行该算子的 worker 规划充足的临时磁盘空间临时文件会完整落盘与内存余量若数据量级较大优先考虑GlacierCreateJobOperatorGlacierJobOperationSensor 本算子的编排组合并在调度层面错峰执行。测试验证调用链如何被确认仓库为该算子提供了两层测试保障可作为理解其行为的权威参考单元测试test_glacier_to_gcs.py 通过 mock 断言了完整的调用链GlacierHook以aws_conn_id实例化并依次调用retrieve_inventory(vault_name...)与retrieve_inventory_results(vault_name..., job_id...)GCSHook以gcp_conn_id与impersonation_chain实例化并调用upload(bucket_name..., object_name..., gzipFalse, filename...)。这验证了算子确实以Glacier 检索 → 临时文件 → GCS 上传的顺序执行且各参数被正确透传。系统测试example_glacier_to_gcs.py 则是可直接运行的端到端 DAG真实创建 vault、创建检索任务、等待完成、上传测试数据、执行跨云迁移、最终清理资源并兼容 Airflow 2.x 与 3.x 两种运行环境。小结GlacierToGCSOperator以极小的实现成本封装了 AWS 与 GCP 两侧的底层 SDK 交互通过GlacierHook发起并拉取 inventory-retrieval 任务结果以流式分块方式落到临时文件再经GCSHook.upload上传至目标桶最终返回gs://URI。实际使用时请务必遵守官方警告结合任务规模规划 worker 的内存与临时磁盘容量并参考系统测试 DAG 将创建任务 → 等待完成 → 迁移的完整链路纳入调度。更多 Glacier 生态算子如GlacierCreateJobOperator、GlacierUploadArchiveOperator可参见 providers/amazon/docs/operators/s3/glacier.rst底层 Hook 实现位于 hooks/glacier.py。【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表