news 2026/9/13 7:58:18

Apache Airflow Amazon S3 操作指南:12 个 Operator 与 2 个 Sensor 的完整实战详解

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
Apache Airflow Amazon S3 操作指南:12 个 Operator 与 2 个 Sensor 的完整实战详解

Apache Airflow Amazon S3 操作指南:12 个 Operator 与 2 个 Sensor 的完整实战详解

【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow

Amazon Simple Storage Service(Amazon S3)是面向互联网的对象存储服务,可用于随时随地从 Web 存储和检索任意规模的数据。在 Apache Airflow 生态中,Amazon Provider(apache-airflow[amazon])通过airflow.providers.amazon.aws.operators.s3airflow.providers.amazon.aws.sensors.s3模块,为 S3 的创建、标记、读写、复制、变换、删除与等待等场景提供了开箱即用的任务组件。本文以 providers/amazon/docs/operators/s3/s3.rst 为核心骨架,结合 源码实现 与 系统测试 DAG,完整讲解每一个 Operator 与 Sensor 的用法、关键参数及底层实现,帮助你直接在 Airflow DAG 中编排 S3 数据管道。

前置准备:使用 S3 组件前的必备条件

要使用本文介绍的所有组件,需要完成以下准备工作(参见 prerequisite_tasks.rst):

  1. 创建必要的 AWS 资源:可以通过 AWS Console 或 AWS CLI 预先创建 IAM 用户/角色、S3 存储桶等资源,并确保所用凭证具备对应操作的权限。

  2. 安装 Amazon Provider:通过 pip 安装apache-airflow[amazon]

    pip install 'apache-airflow[amazon]'

    更详细的安装说明可参考 apache-airflow 安装文档。

  3. 配置 AWS Connection:在 Airflow 中建立名为aws_default的 AWS 连接,用于提供访问凭证。

从源码看,所有 S3 Operator 均继承自AwsBaseOperator[S3Hook](见 operators/s3.py),底层通过 S3Hook 封装 boto3 客户端完成实际操作。S3Hook 对 boto3 进行了薄封装,支持模板化字段(template_fields),因此bucket_names3_key等参数可以在 DAG 中通过 Jinja 模板(如{{ ds }})动态渲染。

Operator 实战:从建桶到删除的完整生命周期

以下所有代码示例均取自 example_s3.py 中对应[START ...]/[END ...]标注区段,可在系统测试环境中直接运行验证。

创建 S3 存储桶:S3CreateBucketOperator

使用S3CreateBucketOperator创建存储桶:

create_bucket = S3CreateBucketOperator( task_id="create_bucket", bucket_name=bucket_name, )

关键参数:

  • bucket_name:要创建的存储桶名称(必填)。
  • bucket_namespace:桶的命名空间。设为account-regional可在账号区域级命名空间中创建桶,默认在全局命名空间中创建。
  • region_name:AWS 区域,不指定时使用 boto3 默认行为。

从源码(S3CreateBucketOperator.execute)可以看到其幂等处理逻辑:execute先调用hook.check_for_bucket()判断桶是否已存在,不存在才调用hook.create_bucket()创建,已存在则仅记录日志跳过,避免重复建桶报错。

删除 S3 存储桶:S3DeleteBucketOperator

使用S3DeleteBucketOperator删除存储桶:

delete_bucket = S3DeleteBucketOperator( task_id="delete_bucket", bucket_name=bucket_name, force_delete=True, )

关键参数:

  • bucket_name:要删除的存储桶名称(必填)。
  • force_delete:设为True时,会先强制删除桶内所有对象再删除桶本身,适用于清理测试环境。

在示例 DAG 中,删除任务被设置为TriggerRule.ALL_DONE(见 example_s3.py),保证无论前面任务成败都会执行清理,这是编写系统测试 DAG 时的常用手法。

设置/获取/删除桶标签:S3PutBucketTaggingOperator、S3GetBucketTaggingOperator、S3DeleteBucketTaggingOperator

设置桶标签:

put_tagging = S3PutBucketTaggingOperator( task_id="put_tagging", bucket_name=bucket_name, key=TAG_KEY, value=TAG_VALUE, )

获取桶标签:

get_tagging = S3GetBucketTaggingOperator( task_id="get_tagging", bucket_name=bucket_name, )

删除桶标签:

delete_tagging = S3DeleteBucketTaggingOperator( task_id="delete_tagging", bucket_name=bucket_name, )

S3PutBucketTaggingOperator通过keyvalue参数写入单个标签键值对;S3GetBucketTaggingOperator读取并返回桶的完整标签集合(get_bucket_taggingAPI 结果);S3DeleteBucketTaggingOperator删除桶上的全部标签。

创建/替换对象:S3CreateObjectOperator

使用S3CreateObjectOperator向桶中写入(或替换)对象:

create_object = S3CreateObjectOperator( task_id="create_object", s3_bucket=bucket_name, s3_key=key, data=DATA, replace=True, )

关键参数:

  • s3_bucket/s3_key:目标桶与对象键。
  • data:要写入的对象内容(示例中为一段 CSV 格式文本)。
  • replace:设为True时允许覆盖已存在的同名对象,否则目标已存在时会跳过或报错(取决于具体实现)。

复制对象:S3CopyObjectOperator

将一个桶中的对象复制到另一个桶:

copy_object = S3CopyObjectOperator( task_id="copy_object", source_bucket_name=bucket_name, dest_bucket_name=bucket_name_2, source_bucket_key=key, dest_bucket_key=key_2, )

使用注意(原文明确强调):

  • 所使用的 S3 连接必须同时具备源桶/源键与目标桶/目标键的访问权限。
  • 如果不希望使用目标桶的默认加密密钥,可通过 AWS KMS 指定服务端加密;使用 KMS 时必须同时提供kms_key_idkms_encryption_type两个参数,并确保角色或用户拥有使用该密钥的权限。

按前缀批量复制:S3CopyPrefixOperator

将某前缀下的所有对象复制到另一桶:

copy_prefix = S3CopyPrefixOperator( task_id="copy_prefix", source_bucket_name=bucket_name, source_bucket_prefix=f"{env_id}-", dest_bucket_name=bucket_name_2, dest_bucket_prefix=f"{env_id}-copied-", )

该 Operator 会把source_bucket_prefix下所有键批量复制到dest_bucket_prefix,适用于整目录/整前缀迁移场景。同样支持 KMS 服务端加密,使用kms_key_idkms_encryption_type时必须成对提供,且连接需同时具备源、目标两侧的访问权限。

删除一个或多个对象:S3DeleteObjectsOperator

delete_objects = S3DeleteObjectsOperator( task_id="delete_objects", bucket=bucket_name_2, keys=key_2, )

keys参数支持传入单个键或键列表,底层通过S3Hook.delete_objects调用 boto3 的delete_objectsAPI 实现批量删除。示例中该任务同样设置了TriggerRule.ALL_DONE作为收尾清理步骤。

变换对象:S3FileTransformOperator

读取源对象的数据,经变换脚本处理后写入目标对象:

file_transform = S3FileTransformOperator( task_id="file_transform", source_s3_key=f"s3://{bucket_name}/{key}", dest_s3_key=f"s3://{bucket_name_2}/{key_2}", # 以 cp 命令作为变换脚本示例 transform_script="cp", replace=True, )

关键参数:

  • source_s3_key/dest_s3_key:支持s3://bucket/key形式的完整 URI。
  • transform_script:本地可执行命令或脚本路径,接收源文件与目标文件路径作为参数(示例中直接使用系统cp命令做复制)。
  • select_expression:可选参数,可传入 Amazon S3 Select 的 SQL 表达式,先从source_s3_key中筛选出需要的数据再交给脚本处理,适合只需处理部分列/行的大文件场景。

列出对象与前缀:S3ListOperator、S3ListPrefixesOperator

列出桶内对象(可按前缀过滤):

list_keys = S3ListOperator( task_id="list_keys", bucket=bucket_name, prefix=PREFIX, )

prefix用于过滤键名以该前缀开头的对象。注意示例中PREFIX = "",空字符串前缀代表桶根目录,关于 S3 前缀的更多说明可参考 AWS 官方文档《Using prefixes》。

列出桶内前缀(按分隔符分组):

list_prefixes = S3ListPrefixesOperator( task_id="list_prefixes", bucket=bucket_name, prefix=PREFIX, delimiter=DELIMITER, )

delimiter参数(示例中为/)用于把对象按键中的分隔符折叠成"前缀",与list_objects_v2Delimiter语义一致——这在模拟目录结构、做分页或分组统计时非常实用。

读取对象内容:S3ReadObjectOperator

将对象内容以字符串形式读回:

read_object = S3ReadObjectOperator( task_id="read_object", s3_bucket=bucket_name, s3_key=key, )

该 Operator 返回对象内容字符串,可通过 XCom 传递给下游任务使用,常用于在 DAG 内部对 S3 中小文件的实时读取与加工。

Sensor 实战:等待与感知 S3 状态变化

Sensors 与 Operators 不同,它们在满足条件之前会持续等待,适合做"上游数据就绪再触发下游"的编排。

等待对象出现:S3KeySensor

S3KeySensor用于等待一个或多个键出现在指定桶中。对每个键,它调用 boto3 的head_objectAPI 检查对象是否存在;当wildcard_matchTrue时改用list_objects_v2API 做通配匹配。需要注意:每检查一个键就产生一次 API 调用,当检查大量键时会产生大量请求,需留意成本与限流。

检查单个文件:

# 检查文件是否存在 sensor_one_key = S3KeySensor( task_id="sensor_one_key", bucket_name=bucket_name, bucket_key=key, )

检查多个文件(同时存在才通过):

# 检查两个文件是否都存在 sensor_two_keys = S3KeySensor( task_id="sensor_two_keys", bucket_name=bucket_name, bucket_key=[key, key_2], )

使用正则匹配:

# 检查是否存在匹配正则表达式的文件 sensor_key_with_regex = S3KeySensor( task_id="sensor_key_with_regex", bucket_name=bucket_name, bucket_key=key_regex_pattern, use_regex=True )

示例中key_regex_pattern = ".*-key",当use_regex=Truebucket_key按正则模式匹配。

自定义校验函数check_fn:可以定义一个接收"匹配到的 S3 对象属性列表"、返回布尔值的函数——返回True表示条件满足,False表示不满足。该函数会对bucket_key中传入的每个键分别调用。之所以入参是对象列表,是因为当wildcard_matchTrue时,一个键模式可能匹配多个文件。列表中的对象属性目前只包含大小,格式为:

[{"Size": int}]

示例:检查所有匹配文件是否都大于 20 字节:

def check_fn(files: list, **kwargs) -> bool: """ 自定义检查示例:检查所有文件是否都大于 20 字节 :param files: S3 对象属性列表。 :return: 条件满足返回 true """ return all(f.get("Size", 0) > 20 for f in files)
# 检查文件是否存在且满足 check_fn 定义的模式 sensor_key_with_function = S3KeySensor( task_id="sensor_key_with_function", bucket_name=bucket_name, bucket_key=key, check_fn=check_fn, )

可延迟(deferrable)模式:将deferrable参数设为True即可让 Sensor 以可延迟模式运行——轮询工作从占用 worker 改为由 triggerer 异步执行,从而高效利用 Airflow worker 资源。注意:使用该模式要求你的 Airflow 部署中已配置并运行 triggerer 组件。

可延迟模式同样支持以上三种用法:

# 检查单个文件 sensor_one_key_deferrable = S3KeySensor( task_id="sensor_one_key_deferrable", bucket_name=bucket_name, bucket_key=key, deferrable=True, ) # 检查多个文件 sensor_two_keys_deferrable = S3KeySensor( task_id="sensor_two_keys_deferrable", bucket_name=bucket_name, bucket_key=[key, key_2], deferrable=True, ) # 正则匹配 + 可延迟 sensor_key_with_regex_deferrable = S3KeySensor( task_id="sensor_key_with_regex_deferrable", bucket_name=bucket_name, bucket_key=key_regex_pattern, use_regex=True, deferrable=True, )

等待前缀对象数稳定:S3KeysUnchangedSensor

S3KeysUnchangedSensor用于监听指定前缀下的对象数量变化,并持续等待直到超过inactivity_period(不活动期,单位秒)内对象数量不再增加才继续执行。这在等待上游持续写入文件、直到写入结束的场景(如等待批处理作业完全落盘)中非常有用:

sensor_keys_unchanged = S3KeysUnchangedSensor( task_id="sensor_keys_unchanged", bucket_name=bucket_name_2, prefix=PREFIX, inactivity_period=10, # inactivity_period 单位为秒 )

原文特别警告:该 Sensor 在 reschedule(重调度)模式下行为不正确,因为重调度调用之间会丢失桶内已列出对象的状态(对象数量统计无法在多次调用间持续累积),因此应避免将其与mode="reschedule"组合使用。与S3KeySensor一样,它也可通过deferrable=True切换为可延迟模式,由 triggerer 异步轮询。

端到端示例:完整 S3 工作流编排

将上述组件串联起来,就是一个覆盖"建桶 → 打标签 → 写对象 → 读对象 → 列表 → 等待 → 复制/变换 → 等稳定 → 清理"全流程的 DAG。下面展示 example_s3.py 中的核心依赖链:

chain( # 测试环境准备 test_context, # 测试主体 create_bucket, create_bucket_2, put_tagging, get_tagging, delete_tagging, create_object, create_object_2, read_object, list_prefixes, list_keys, [sensor_one_key, sensor_two_keys, sensor_key_with_function, sensor_key_with_regex], [ sensor_one_key_deferrable, sensor_two_keys_deferrable, sensor_key_with_function_deferrable, sensor_key_with_regex_deferrable, ], copy_object, copy_prefix, file_transform, sensor_keys_unchanged, # 测试清理 delete_objects, delete_bucket, delete_bucket_2, )

几个值得借鉴的编排要点:

  • 并发等待:多个 Sensor 通过列表传入chain(),可并行等待不同条件(键存在、正则匹配、自定义校验)同时满足后再继续。
  • 触发规则收尾:删除类任务(delete_objectsdelete_bucket)显式设置trigger_rule = TriggerRule.ALL_DONE,无论测试主体成败都执行资源清理。
  • 动态命名:桶名与键名基于系统测试上下文env_id生成(如{env_id}-s3-bucket{env_id}-key),避免多租户/多轮测试之间的资源冲突。
  • watcher 收尾:DAG 末尾追加watcher()任务,用于在存在 tearDown 任务(带触发规则)时正确标记系统测试的成功/失败。

该 DAG 定义于DAG_ID = "example_s3",采用schedule="@once"start_date=datetime(2021, 1, 1)catchup=False,可作为自研 S3 管道的编排范式参考。

源码结构速查

  • Operators 定义:operators/s3.py 中依次定义了S3CreateBucketOperator(L45)、S3DeleteBucketOperator(L95)、S3GetBucketTaggingOperator(L138)、S3PutBucketTaggingOperator(L174)、S3DeleteBucketTaggingOperator(L228)、S3CopyObjectOperator(L268)、S3CopyPrefixOperator(L388)、S3CreateObjectOperator(L543)、S3DeleteObjectsOperator(L641)、S3FileTransformOperator(L776)、S3ListOperator(L949)、S3ListPrefixesOperator(L1026)、S3ReadObjectOperator(L1094)。
  • Sensors 定义:sensors/s3.py 中定义S3KeySensorS3KeysUnchangedSensor
  • 底层 Hook:hooks/s3.py 提供S3Hook,封装建桶、删桶、读写对象、标签管理、对象列表等 boto3 操作。
  • 系统测试:example_s3.py 是本文所有示例的权威出处,可通过 pytest 运行(参见 system_tests 文档)。

参考

  • boto3 官方 S3 客户端 API 文档:S3.Client.head_objectS3.Client.list_objects_v2
  • AWS 官方《Amazon S3 User Guide》中关于前缀(prefixes)与 S3 Select 的说明

使用提示:本文涉及的 API 调用均会产生 AWS 费用,建议先在低成本测试桶中验证;对生产环境,务必为所用 IAM 凭证配置最小权限策略,并谨慎使用force_deletereplace等破坏性参数。

【免费下载链接】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 7:58:03

JESD22-B112C标准解析:IC封装翘曲测量与优化

1. JESD22-B112C标准概述JESD22-B112C是由JEDEC固态技术协会制定的表面贴装集成电路(IC)在高温环境下封装翘曲测量的行业标准测试方法。该标准主要针对电子封装行业在回流焊工艺中出现的封装变形问题,提供了一套标准化的测量流程和评估体系。在表面贴装技术(SMT)工艺…

作者头像 李华
网站建设 2026/9/13 7:56:42

Kioxia BG7固态硬盘解析:OEM原厂盘的优势与实测

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

作者头像 李华
网站建设 2026/9/13 7:54:49

金融数据处理中的前导零问题与解决方案

1. 数据处理中的前导零陷阱:为什么股票代码必须作为字符串读取在金融数据处理领域,A股股票代码的处理看似简单却暗藏玄机。许多新手在处理CSV或Excel格式的财务数据时,经常会遇到一个典型问题:以"600519"(贵…

作者头像 李华