ARTICLE DETAIL

资讯详情

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

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

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/airflowAmazon Simple Storage ServiceAmazon S3是面向互联网的对象存储服务可用于随时随地从 Web 存储和检索任意规模的数据。在 Apache Airflow 生态中Amazon Providerapache-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_idcreate_bucket, bucket_namebucket_name, )关键参数bucket_name要创建的存储桶名称必填。bucket_namespace桶的命名空间。设为account-regional可在账号区域级命名空间中创建桶默认在全局命名空间中创建。region_nameAWS 区域不指定时使用 boto3 默认行为。从源码S3CreateBucketOperator.execute可以看到其幂等处理逻辑execute先调用hook.check_for_bucket()判断桶是否已存在不存在才调用hook.create_bucket()创建已存在则仅记录日志跳过避免重复建桶报错。删除 S3 存储桶S3DeleteBucketOperator使用S3DeleteBucketOperator删除存储桶delete_bucket S3DeleteBucketOperator( task_iddelete_bucket, bucket_namebucket_name, force_deleteTrue, )关键参数bucket_name要删除的存储桶名称必填。force_delete设为True时会先强制删除桶内所有对象再删除桶本身适用于清理测试环境。在示例 DAG 中删除任务被设置为TriggerRule.ALL_DONE见 example_s3.py保证无论前面任务成败都会执行清理这是编写系统测试 DAG 时的常用手法。设置/获取/删除桶标签S3PutBucketTaggingOperator、S3GetBucketTaggingOperator、S3DeleteBucketTaggingOperator设置桶标签put_tagging S3PutBucketTaggingOperator( task_idput_tagging, bucket_namebucket_name, keyTAG_KEY, valueTAG_VALUE, )获取桶标签get_tagging S3GetBucketTaggingOperator( task_idget_tagging, bucket_namebucket_name, )删除桶标签delete_tagging S3DeleteBucketTaggingOperator( task_iddelete_tagging, bucket_namebucket_name, )S3PutBucketTaggingOperator通过key与value参数写入单个标签键值对S3GetBucketTaggingOperator读取并返回桶的完整标签集合get_bucket_taggingAPI 结果S3DeleteBucketTaggingOperator删除桶上的全部标签。创建/替换对象S3CreateObjectOperator使用S3CreateObjectOperator向桶中写入或替换对象create_object S3CreateObjectOperator( task_idcreate_object, s3_bucketbucket_name, s3_keykey, dataDATA, replaceTrue, )关键参数s3_bucket/s3_key目标桶与对象键。data要写入的对象内容示例中为一段 CSV 格式文本。replace设为True时允许覆盖已存在的同名对象否则目标已存在时会跳过或报错取决于具体实现。复制对象S3CopyObjectOperator将一个桶中的对象复制到另一个桶copy_object S3CopyObjectOperator( task_idcopy_object, source_bucket_namebucket_name, dest_bucket_namebucket_name_2, source_bucket_keykey, dest_bucket_keykey_2, )使用注意原文明确强调所使用的 S3 连接必须同时具备源桶/源键与目标桶/目标键的访问权限。如果不希望使用目标桶的默认加密密钥可通过 AWS KMS 指定服务端加密使用 KMS 时必须同时提供kms_key_id和kms_encryption_type两个参数并确保角色或用户拥有使用该密钥的权限。按前缀批量复制S3CopyPrefixOperator将某前缀下的所有对象复制到另一桶copy_prefix S3CopyPrefixOperator( task_idcopy_prefix, source_bucket_namebucket_name, source_bucket_prefixf{env_id}-, dest_bucket_namebucket_name_2, dest_bucket_prefixf{env_id}-copied-, )该 Operator 会把source_bucket_prefix下所有键批量复制到dest_bucket_prefix适用于整目录/整前缀迁移场景。同样支持 KMS 服务端加密使用kms_key_id与kms_encryption_type时必须成对提供且连接需同时具备源、目标两侧的访问权限。删除一个或多个对象S3DeleteObjectsOperatordelete_objects S3DeleteObjectsOperator( task_iddelete_objects, bucketbucket_name_2, keyskey_2, )keys参数支持传入单个键或键列表底层通过S3Hook.delete_objects调用 boto3 的delete_objectsAPI 实现批量删除。示例中该任务同样设置了TriggerRule.ALL_DONE作为收尾清理步骤。变换对象S3FileTransformOperator读取源对象的数据经变换脚本处理后写入目标对象file_transform S3FileTransformOperator( task_idfile_transform, source_s3_keyfs3://{bucket_name}/{key}, dest_s3_keyfs3://{bucket_name_2}/{key_2}, # 以 cp 命令作为变换脚本示例 transform_scriptcp, replaceTrue, )关键参数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_idlist_keys, bucketbucket_name, prefixPREFIX, )prefix用于过滤键名以该前缀开头的对象。注意示例中PREFIX 空字符串前缀代表桶根目录关于 S3 前缀的更多说明可参考 AWS 官方文档《Using prefixes》。列出桶内前缀按分隔符分组list_prefixes S3ListPrefixesOperator( task_idlist_prefixes, bucketbucket_name, prefixPREFIX, delimiterDELIMITER, )delimiter参数示例中为/用于把对象按键中的分隔符折叠成前缀与list_objects_v2的Delimiter语义一致——这在模拟目录结构、做分页或分组统计时非常实用。读取对象内容S3ReadObjectOperator将对象内容以字符串形式读回read_object S3ReadObjectOperator( task_idread_object, s3_bucketbucket_name, s3_keykey, )该 Operator 返回对象内容字符串可通过 XCom 传递给下游任务使用常用于在 DAG 内部对 S3 中小文件的实时读取与加工。Sensor 实战等待与感知 S3 状态变化Sensors 与 Operators 不同它们在满足条件之前会持续等待适合做上游数据就绪再触发下游的编排。等待对象出现S3KeySensorS3KeySensor用于等待一个或多个键出现在指定桶中。对每个键它调用 boto3 的head_objectAPI 检查对象是否存在当wildcard_match为True时改用list_objects_v2API 做通配匹配。需要注意每检查一个键就产生一次 API 调用当检查大量键时会产生大量请求需留意成本与限流。检查单个文件# 检查文件是否存在 sensor_one_key S3KeySensor( task_idsensor_one_key, bucket_namebucket_name, bucket_keykey, )检查多个文件同时存在才通过# 检查两个文件是否都存在 sensor_two_keys S3KeySensor( task_idsensor_two_keys, bucket_namebucket_name, bucket_key[key, key_2], )使用正则匹配# 检查是否存在匹配正则表达式的文件 sensor_key_with_regex S3KeySensor( task_idsensor_key_with_regex, bucket_namebucket_name, bucket_keykey_regex_pattern, use_regexTrue )示例中key_regex_pattern .*-key当use_regexTrue时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_idsensor_key_with_function, bucket_namebucket_name, bucket_keykey, check_fncheck_fn, )可延迟deferrable模式将deferrable参数设为True即可让 Sensor 以可延迟模式运行——轮询工作从占用 worker 改为由 triggerer 异步执行从而高效利用 Airflow worker 资源。注意使用该模式要求你的 Airflow 部署中已配置并运行 triggerer 组件。可延迟模式同样支持以上三种用法# 检查单个文件 sensor_one_key_deferrable S3KeySensor( task_idsensor_one_key_deferrable, bucket_namebucket_name, bucket_keykey, deferrableTrue, ) # 检查多个文件 sensor_two_keys_deferrable S3KeySensor( task_idsensor_two_keys_deferrable, bucket_namebucket_name, bucket_key[key, key_2], deferrableTrue, ) # 正则匹配 可延迟 sensor_key_with_regex_deferrable S3KeySensor( task_idsensor_key_with_regex_deferrable, bucket_namebucket_name, bucket_keykey_regex_pattern, use_regexTrue, deferrableTrue, )等待前缀对象数稳定S3KeysUnchangedSensorS3KeysUnchangedSensor用于监听指定前缀下的对象数量变化并持续等待直到超过inactivity_period不活动期单位秒内对象数量不再增加才继续执行。这在等待上游持续写入文件、直到写入结束的场景如等待批处理作业完全落盘中非常有用sensor_keys_unchanged S3KeysUnchangedSensor( task_idsensor_keys_unchanged, bucket_namebucket_name_2, prefixPREFIX, inactivity_period10, # inactivity_period 单位为秒 )原文特别警告该 Sensor 在 reschedule重调度模式下行为不正确因为重调度调用之间会丢失桶内已列出对象的状态对象数量统计无法在多次调用间持续累积因此应避免将其与modereschedule组合使用。与S3KeySensor一样它也可通过deferrableTrue切换为可延迟模式由 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采用scheduleonce、start_datedatetime(2021, 1, 1)、catchupFalse可作为自研 S3 管道的编排范式参考。源码结构速查Operators 定义operators/s3.py 中依次定义了S3CreateBucketOperatorL45、S3DeleteBucketOperatorL95、S3GetBucketTaggingOperatorL138、S3PutBucketTaggingOperatorL174、S3DeleteBucketTaggingOperatorL228、S3CopyObjectOperatorL268、S3CopyPrefixOperatorL388、S3CreateObjectOperatorL543、S3DeleteObjectsOperatorL641、S3FileTransformOperatorL776、S3ListOperatorL949、S3ListPrefixesOperatorL1026、S3ReadObjectOperatorL1094。Sensors 定义sensors/s3.py 中定义S3KeySensor、S3KeysUnchangedSensor。底层 Hookhooks/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),仅供参考
返回列表