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.s3与airflow.providers.amazon.aws.sensors.s3模块,为 S3 的创建、标记、读写、复制、变换、删除与等待等场景提供了开箱即用的任务组件。本文以 providers/amazon/docs/operators/s3/s3.rst 为核心骨架,结合 源码实现 与 系统测试 DAG,完整讲解每一个 Operator 与 Sensor 的用法、关键参数及底层实现,帮助你直接在 Airflow DAG 中编排 S3 数据管道。
前置准备:使用 S3 组件前的必备条件
要使用本文介绍的所有组件,需要完成以下准备工作(参见 prerequisite_tasks.rst):
创建必要的 AWS 资源:可以通过 AWS Console 或 AWS CLI 预先创建 IAM 用户/角色、S3 存储桶等资源,并确保所用凭证具备对应操作的权限。
安装 Amazon Provider:通过 pip 安装
apache-airflow[amazon]:pip install 'apache-airflow[amazon]'更详细的安装说明可参考 apache-airflow 安装文档。
配置 AWS Connection:在 Airflow 中建立名为
aws_default的 AWS 连接,用于提供访问凭证。
从源码看,所有 S3 Operator 均继承自AwsBaseOperator[S3Hook](见 operators/s3.py),底层通过 S3Hook 封装 boto3 客户端完成实际操作。S3Hook 对 boto3 进行了薄封装,支持模板化字段(template_fields),因此bucket_name、s3_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通过key与value参数写入单个标签键值对;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_id和kms_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_id与kms_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_v2的Delimiter语义一致——这在模拟目录结构、做分页或分组统计时非常实用。
读取对象内容: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_match为True时改用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=True时bucket_key按正则模式匹配。
自定义校验函数check_fn:可以定义一个接收"匹配到的 S3 对象属性列表"、返回布尔值的函数——返回True表示条件满足,False表示不满足。该函数会对bucket_key中传入的每个键分别调用。之所以入参是对象列表,是因为当wildcard_match为True时,一个键模式可能匹配多个文件。列表中的对象属性目前只包含大小,格式为:
[{"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_objects、delete_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 中定义
S3KeySensor、S3KeysUnchangedSensor。 - 底层 Hook:hooks/s3.py 提供
S3Hook,封装建桶、删桶、读写对象、标签管理、对象列表等 boto3 操作。 - 系统测试:example_s3.py 是本文所有示例的权威出处,可通过 pytest 运行(参见 system_tests 文档)。
参考
- boto3 官方 S3 客户端 API 文档:
S3.Client.head_object、S3.Client.list_objects_v2等 - AWS 官方《Amazon S3 User Guide》中关于前缀(prefixes)与 S3 Select 的说明
使用提示:本文涉及的 API 调用均会产生 AWS 费用,建议先在低成本测试桶中验证;对生产环境,务必为所用 IAM 凭证配置最小权限策略,并谨慎使用
force_delete与replace等破坏性参数。
【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考