ARTICLE DETAIL

资讯详情

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

Apache Airflow Amazon 提供包 S3 运算符与传感器完全指南:从桶管理到对象转换与异步等待

Apache Airflow Amazon 提供包 S3 运算符与传感器完全指南:从桶管理到对象转换与异步等待 Apache Airflow Amazon 提供包 S3 运算符与传感器完全指南从桶管理到对象转换与异步等待【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow本篇技术指南聚焦 Apache Airflow Amazon Providerapache-airflow[amazon]中面向 Amazon S3 的整套 Operators 与 Sensors内容以仓库内官方指南 providers/amazon/docs/operators/s3/s3.rst 为骨架结合系统测试示例 providers/amazon/tests/system/amazon/aws/example_s3.py 与源码实现深入展开。读完本文你将掌握如何在 DAG 中创建/删除 S3 桶、管理桶标签、读写与复制对象、按前缀批量复制、通过 S3 Select 与本地脚本完成对象转换以及如何使用S3KeySensor与S3KeysUnchangedSensor以同步或 deferrable可延迟模式等待 S3 就绪状态。前置任务安装、连接与权限准备官方文档在 prerequisite_tasks.rst 中明确了使用 S3 运算符前的三步准备准备 AWS 资源通过 AWS Console 或 AWS CLI 预先创建所需的 IAM 角色/用户、桶等资源。安装 API 依赖库使用 pip 安装 Amazon Providerpip install apache-airflow[amazon]配置 Airflow Connection在 Airflow 中为 AWS 配置 Connection默认连接 id 为aws_default详见 AWS 连接配置。从源码看所有 S3 运算符与传感器都继承自AwsBaseOperator/AwsBaseSensor其aws_conn_id参数语义为若为None或空则使用默认 boto3 凭据链行为若 Airflow 以分布式方式运行且aws_conn_id为空则需在每个 worker 节点上维护默认 boto3 配置见 operators/s3.py 的 docstring。除aws_conn_id外各运算符还普遍支持region_name、verifySSL 证书校验、botocore_config等通用参数。桶级操作创建与删除S3CreateBucketOperator创建 S3 桶使用airflow.providers.amazon.aws.operators.s3.S3CreateBucketOperator创建桶系统测试示例example_s3.py中的用法如下create_bucket S3CreateBucketOperator( task_idcreate_bucket, bucket_namebucket_name, )关键参数参数说明bucket_name要创建的桶名称必填templatedbucket_namespace桶命名空间设为account-regional则创建在账户区域命名空间下不指定则创建在全局命名空间templatedregion_nameAWS 区域不指定则使用默认 boto3 行为aws_conn_id使用的 Airflow AWS 连接从源码 operators/s3.py 可以看出execute会先调用S3Hook.check_for_bucket检查桶是否已存在已存在则仅记录日志并跳过不存在才调用create_bucket创建——这是幂等设计可在 DAG 中安全地重复运行。底层S3Hook.create_buckethooks/s3.py在构造 boto3 请求时有几个细节值得注意若未指定region_name且连接配置为aws-global区域会抛出AirflowException避免在无区域信息时误建桶仅当区域不是us-east-1时才附带CreateBucketConfiguration.LocationConstraintbucket_namespace会被透传为 boto3create_bucket的BucketNamespace参数。S3DeleteBucketOperator删除 S3 桶使用S3DeleteBucketOperator删除桶example_s3.pydelete_bucket S3DeleteBucketOperator( task_iddelete_bucket, bucket_namebucket_name, force_deleteTrue, )bucket_name要删除的桶名称force_delete是否强制删除桶内所有对象后再删除桶默认False。源码实现operators/s3.py同样先检查桶是否存在避免对不存在的桶报错。底层S3Hook.delete_buckethooks/s3.py在force_deleteTrue时会先清空桶内对象并内置max_retries默认 5 次重试机制以规避边删对象边删桶的竞态条件——S3 要求桶必须为空才能删除。在系统测试示例中删除类任务都显式设置了delete_bucket.trigger_rule TriggerRule.ALL_DONE确保即使前面某个任务失败清理任务仍会执行避免测试残留资源。桶标签管理Put / Get / DeleteS3 桶标签Bucket Tagging可用于成本分摊、资源分类与权限控制。Airflow 提供了三个配套运算符S3PutBucketTaggingOperator设置桶标签example_s3.pyput_tagging S3PutBucketTaggingOperator( task_idput_tagging, bucket_namebucket_name, keyTAG_KEY, valueTAG_VALUE, )S3GetBucketTaggingOperator获取桶标签example_s3.pyget_tagging S3GetBucketTaggingOperator( task_idget_tagging, bucket_namebucket_name, )S3DeleteBucketTaggingOperator删除桶标签example_s3.pydelete_tagging S3DeleteBucketTaggingOperator( task_iddelete_tagging, bucket_namebucket_name, )参数细节来自源码 operators/s3.pyS3PutBucketTaggingOperator支持两种传参方式通过keyvalue设置单个标签二者必须成对出现或通过tag_set传入包含多个标签的 dict 或 key/value 对列表其tag_set字段注册了 JSON 渲染器template_fields_renderers {tag_set: json}便于在 UI 中以 JSON 形式渲染模板化内容三个运算符的execute都会先check_for_bucket桶不存在时仅记录 warning 并返回NoneS3GetBucketTaggingOperator的返回值标签集合会自动推送到 XCom供下游任务消费。对象级操作创建、读取、复制、删除S3CreateObjectOperator创建对象S3CreateObjectOperator将字符串或字节数据直接写入 S3 对象example_s3.pycreate_object S3CreateObjectOperator( task_idcreate_object, s3_bucketbucket_name, s3_keykey, dataDATA, replaceTrue, )系统测试示例中的数据源是一段多行文本DATA apple,0.5 milk,2.5 bread,4.0 参数说明operators/s3.py参数说明s3_bucket目标桶名称templated当s3_key以完整s3://URL 给出时可省略s3_key对象键templated支持完整s3://URL 或相对根目录的路径data要写入的内容str或bytes均可replace若对象已存在是否覆盖默认Falseencrypt是否启用 S3 服务端加密默认Falseacl_policy上传对象的 canned ACL 策略encoding当data为字符串时的编码方式仅字符串数据可用compression压缩方式目前仅支持gzip仅字符串数据可用从源码看execute会调用S3Hook.get_s3_bucket_key统一解析s3://URL 或分离的 bucket/key 写法当data为字符串时走S3Hook.load_string支持 encoding 与 gzip 压缩为字节时走S3Hook.load_bytes。该运算符还实现了 OpenLineage 的get_openlineage_facets_on_start会把目标对象作为 output dataset 上报从而在数据血缘Data Lineage中记录产出关系。S3ReadObjectOperator读取对象S3ReadObjectOperator读取对象内容并以字符串返回example_s3.pyread_object S3ReadObjectOperator( task_idread_object, s3_bucketbucket_name, s3_keykey, )s3_bucket桶名称templateds3_key为完整s3://URL 时可省略s3_key对象键templated。源码 operators/s3.py 通过S3Hook.read_key取回对象体并按 UTF-8 解码为字符串返回。返回值自动推送至 XCom因此下游任务可以通过{{ task_instance.xcom_pull(task_idsread_object) }}直接消费文件内容——非常适合读取小文件并做下游处理的流水线场景。S3CopyObjectOperator单对象复制S3CopyObjectOperator把对象从一个桶复制到另一个桶example_s3.pycopy_object S3CopyObjectOperator( task_idcopy_object, source_bucket_namebucket_name, dest_bucket_namebucket_name_2, source_bucket_keykey, dest_bucket_keykey_2, )参数说明operators/s3.py参数说明source_bucket_key/dest_bucket_key源/目标对象键templated支持完整s3://URL此时对应 bucket 参数可省略source_bucket_name/dest_bucket_name源/目标桶名templatedsource_version_id源对象版本 ID可选需开启桶版本控制acl_policy复制后对象的 canned ACL默认privatemeta_data_directive元数据指令COPY继承源对象元数据或REPLACE使用请求中提供的元数据kms_key_idAWS KMS 密钥的 ARN/id/别名用于目标对象服务端加密kms_encryption_typeKMS 加密类型aws:kms标准 KMS或aws:kms:dsse双屏蔽 KMS注意两个官方强调的约束使用的 S3 连接必须具备源与目标两侧桶/键的访问权限若不想使用目标桶默认密钥而改用 KMS 加密必须同时提供kms_key_id与kms_encryption_type两个参数并确保所用角色/用户具有使用该密钥的权限。这一约束在底层有强校验S3Hook.copy_objecthooks/s3.py中if bool(kms_key_id) ! bool(kms_encryption_type)会抛出ValueError防止只传一个参数导致加密配置不完整。该 Hook 同时通过get_hook_lineage_collector()上报源/目标资产为复制操作提供 OpenLineage 血缘信息运算符自身的get_openlineage_facets_on_start也会把源与目标对象分别注册为 input / output dataset。S3CopyPrefixOperator按前缀批量复制S3CopyPrefixOperator复制源桶某个前缀下的所有对象到目标桶example_s3.pycopy_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-, )参数说明operators/s3.py参数说明source_bucket_prefix/dest_bucket_prefix源/目标前缀templated支持完整s3://URLsource_bucket_name/dest_bucket_name源/目标桶名templatedkms_key_id/kms_encryption_type目标对象 KMS 加密配置同样需成对提供continue_on_failure单个对象复制失败时False立即失败默认True则尝试复制完所有对象后在任务末尾统一报错acl_policy/meta_data_directive同S3CopyObjectOperator源码实现operators/s3.py展示了其底层机制通过 boto3 的get_paginator(list_objects_v2)分页列出源前缀下的所有对象对每个对象计算dest_key dest_bucket_prefix source_key[len(source_bucket_prefix):]即保留相对前缀结构再逐个调用S3Hook.copy_object完成复制并累计成功/失败计数。这意味着一大批文件的复制只消耗一次 API 调度但每个对象仍对应一次copy_object调用。S3DeleteObjectsOperator批量删除对象S3DeleteObjectsOperator支持单次请求删除一个或多个对象example_s3.pydelete_objects S3DeleteObjectsOperator( task_iddelete_objects, bucketbucket_name_2, keyskey_2, )删除目标支持四种模式来自源码 docstringoperators/s3.py参数说明bucket桶名称templated必填keys单个键字符串或键列表对应 S3 的批量删除 APIprefix前缀删除桶内所有匹配该前缀的对象templatedfrom_datetime/to_datetime按LastModified时间范围过滤删除最后修改时间大于from_datetime或小于to_datetime的对象templated对象转换S3FileTransformOperator 与 S3 SelectS3FileTransformOperator是数据管道中非常实用的ETL 胶水运算符它把源 S3 对象下载到本地临时文件运行指定的转换脚本或 S3 Select 表达式进行加工再将结果上传到目标 S3 位置example_s3.pyfile_transform S3FileTransformOperator( task_idfile_transform, source_s3_keyfs3://{bucket_name}/{key}, dest_s3_keyfs3://{bucket_name_2}/{key_2}, # 以 cp 命令作为转换脚本示例 transform_scriptcp, replaceTrue, )参数说明operators/s3.py参数说明source_s3_key/dest_s3_key源/目标 S3 位置templated完整s3://URL 或相对路径transform_script可执行转换脚本的路径脚本接收两个位置参数本地源文件路径、本地目标文件路径select_expressionS3 Select 表达式可从source_s3_key中只取满足条件的部分数据指定后可不提供transform_scriptselect_expr_serialization_config包含 S3 Select 输入/输出序列化配置的 dictscript_args传给转换脚本的额外参数templatedsource_aws_conn_id/dest_aws_conn_id源/目标可分别使用不同的 AWS 连接默认均为aws_default实现跨账号读取source_verify/dest_verify源/目标连接是否校验 SSL 证书可传False或 CA 证书 bundle 路径replace目标键已存在时是否覆盖执行流程与原理源码 operators/s3.py校验transform_script与select_expression至少提供一个否则抛出AirflowException用S3Hook.check_for_key确认源键存在在本地创建两个NamedTemporaryFile若指定了select_expression则调用S3Hook.select_key把 S3 Select 的过滤结果写入源临时文件否则直接download_fileobj下载整个对象若指定transform_script通过subprocess.Popen([transform_script, f_source.name, f_dest.name, *script_args])运行脚本实时把脚本 stdout 写入任务日志非零返回码会触发AirflowException将转换后的本地文件通过dest_s3.load_file上传到目标键replace控制覆盖行为。S3 Select 表达式语法如SELECT * FROM s3object s WHERE s.Name Jane可以直接在服务端完成 CSV/JSON 数据的列与行过滤减少网络传输量——这是官方文档推荐的可选的 Amazon S3 Select 表达式用法目的是从source_s3_key中挑选所需数据再交给后续处理。列出对象与前缀S3ListOperator 与 S3ListPrefixesOperatorS3ListOperator列出桶内对象S3ListOperator返回桶内所有匹配前缀的对象名列表结果可经 XCom 供下游任务使用example_s3.pylist_keys S3ListOperator( task_idlist_keys, bucketbucket_name, prefixPREFIX, )参数bucket桶名templated、prefix只返回名称以该前缀开头的对象templated、delimiter键层次分隔符templated、apply_wildcard是否将*视为通配符而非普通字符。系统测试中PREFIX 表示桶根目录空字符串前缀。文档中的经典示例是s3_file S3ListOperator( task_idlist_3s_files, bucketdata, prefixcustomers/2018/04/, delimiter/, aws_conn_idaws_customers_conn, )该示例配合delimiter/只会列出该前缀下不含子目录的文件与S3ListPrefixesOperator形成互补。S3ListPrefixesOperator列出前缀子目录S3ListPrefixesOperator返回桶内指定前缀下的所有**前缀即子目录**列表example_s3.pylist_prefixes S3ListPrefixesOperator( task_idlist_prefixes, bucketbucket_name, prefixPREFIX, delimiterDELIMITER, )其中DELIMITER /。底层S3Hook.list_prefixeshooks/s3.py使用list_objects_v2分页器并读取响应中的CommonPrefixes字段来收集前缀支持page_size与max_items分页控制且在requester_pays桶上会自动附加RequestPayerrequester参数。等待 S3 就绪S3KeySensorS3KeySensor用于等待一个或多个键在 S3 桶中出现是典型的文件到达即触发下游模式example_s3.py# 检查单个文件是否存在 sensor_one_key S3KeySensor( task_idsensor_one_key, bucket_namebucket_name, bucket_keykey, )核心参数与底层 API参数说明bucket_key等待的键支持字符串或字符串列表也支持完整s3://URL此时bucket_name置为Nonebucket_name桶名仅当bucket_key不是完整s3://URL 时需要wildcard_match是否将bucket_key解释为 Unix 通配符模式use_regex是否使用正则表达式匹配键名check_fn自定义检查函数接收匹配对象的属性列表返回布尔值metadata_keys要收集并传给check_fn的head_object属性列表默认[Size, Key]传*返回全部属性deferrable是否以可延迟deferrable模式运行默认读取配置operators.default_deferrable文档明确说明其底层探测方式对每个键调用head_objectAPI当wildcard_matchTrue时改用list_objects_v2检查是否存在。需要特别注意的是每检查一个键就产生一次 API 调用当检查大量键时请评估 API 成本与配额。从源码sensors/s3.py看_check_key的探测路径有三种wildcard_matchTrue先用正则切出前缀通过S3Hook.iter_file_metadata迭代对象元数据再用fnmatch匹配若未提供check_fn命中第一个匹配即返回True避免不必要的遍历use_regexTrue遍历桶内全部元数据并用re.match匹配默认模式直接head_object探测单个键并从响应中抽取metadata_keys指定属性其中Size兼容映射为ContentLength。poke方法对bucket_key列表执行all(...)语义——所有键都存在才视为满足条件。多键等待bucket_key传入列表即可同时等待多个文件example_s3.py# 检查两个文件是否都存在 sensor_two_keys S3KeySensor( task_idsensor_two_keys, bucket_namebucket_name, bucket_key[key, key_2], )正则匹配等待use_regexTrue让bucket_key按正则语义匹配example_s3.py# 检查是否存在匹配某正则模式的文件 sensor_key_with_regex S3KeySensor( task_idsensor_key_with_regex, bucket_namebucket_name, bucket_keykey_regex_pattern, use_regexTrue )系统测试中的正则模式为key_regex_pattern .*-key用于匹配任意以-key结尾的键。自定义检查函数 check_fn当内置存在性判断不够用时可以定义check_fn它接收匹配 S3 对象属性的列表并返回布尔值——True表示条件满足False表示未满足example_s3.pydef check_fn(files: list, **kwargs) - bool: Example of custom check: check if all files are bigger than 20 bytes :param files: List of S3 object attributes. :return: true if the criteria is met return all(f.get(Size, 0) 20 for f in files)该函数会对bucket_key中传入的每个键各调用一次。之所以函数参数是对象列表当wildcard_matchTrue时一个键模式可能匹配多个文件匹配到的 S3 对象属性列表目前仅包含大小字段格式为[{Size: int}]当前源码的默认metadata_keys为[Size, Key]即默认还会带上键名。使用方式example_s3.py# 检查文件是否存在且满足 check_fn 定义的规则 sensor_key_with_function S3KeySensor( task_idsensor_key_with_function, bucket_namebucket_name, bucket_keykey, check_fncheck_fn, )源码sensors/s3.py对check_fn的签名做了向后兼容处理若函数声明了**kwargs变长关键字参数则以check_fn(files, **context)调用附带 Airflow 上下文否则仅传files。官方 docstring 还给出等待任意大于 1MB 的对象的经典范例any(f.get(Size, 0) 1048576 for f in files)。deferrable 可延迟模式文档强调将deferrableTrue即可让该 Sensor 以可延迟模式运行——轮询工作从 worker 转移到triggerer上异步执行从而高效利用 Airflow worker 资源。前提是 Airflow 部署中必须运行 triggerer 组件。官方为单键、多键、自定义函数、正则四种场景都提供了 deferrable 变体示例example_s3.py# 单键deferrable sensor_one_key_deferrable S3KeySensor( task_idsensor_one_key_deferrable, bucket_namebucket_name, bucket_keykey, deferrableTrue, ) # 多键deferrable sensor_two_keys_deferrable S3KeySensor( task_idsensor_two_keys_deferrable, bucket_namebucket_name, bucket_key[key, key_2], deferrableTrue, ) # check_fn deferrable sensor_key_with_function_deferrable S3KeySensor( task_idsensor_key_with_function_deferrable, bucket_namebucket_name, bucket_keykey, check_fncheck_fn, deferrableTrue, ) # 正则 deferrable sensor_key_with_regex_deferrable S3KeySensor( task_idsensor_key_with_regex_deferrable, bucket_namebucket_name, bucket_keykey_regex_pattern, use_regexTrue, deferrableTrue, )从源码sensors/s3.py看deferrable 模式的执行链路为execute先poke一次未命中则调用self.defer(...)挂起任务并创建S3KeyTriggertrigger 在 triggerer 进程内异步轮询事件返回后经execute_complete恢复——若事件状态为running且check_fn仍未满足会再次_defer()继续等待事件状态为error则抛出AirflowException。等待前缀稳定S3KeysUnchangedSensorS3KeysUnchangedSensor用于等待某个前缀下的对象数量在指定不活动期内不再增长常用来判断一批文件已经全部到达例如等待数据上传完成后再触发下游处理example_s3.pysensor_keys_unchanged S3KeysUnchangedSensor( task_idsensor_keys_unchanged, bucket_namebucket_name_2, prefixPREFIX, inactivity_period10, # inactivity_period 单位为秒 )参数说明sensors/s3.py参数说明bucket_nameS3 桶名prefix等待的前缀相对桶根路径默认桶根inactivity_period判定键保持不变所需的总不活动秒数默认 3600 秒1 小时为负数会抛出ValueError。注意该机制非实时实际返回可能比设定周期晚一个poke_intervalmin_objects前缀下对象数量的最小阈值默认 1不活动期过后对象数仍低于该值则视为失败previous_objects上次 poke 时记录的对象 ID 集合allow_delete两次 poke 之间允许对象被删除True记录 warning 并视为合法行为False则抛出AirflowException报错deferrable是否以 deferrable 模式运行重要限制官方明确警告该 Sensor在 reschedule重调度模式下行为不正确因为重调度调用之间桶内已列出的对象状态会丢失。因此应使用默认的 poke 模式或 deferrable 模式。底层判定逻辑is_keys_unchangedsensors/s3.py展示了其状态机若当前对象集合新增了对象重置last_activity_time与inactivity_seconds更新previous_objects返回False继续等待若期间有对象被删除allow_deleteTrue时重置计时并继续等待allow_deleteFalse时抛出AirflowException若对象无变化且inactivity_seconds inactivity_period对象数达到min_objects则返回True成功否则记录错误并返回False失败。poke通过self.hook.list_keys(bucket_name, prefixprefix)获取当前对象集合传入状态机。同样支持deferrableTrue届时由S3KeysUnchangedTrigger在 triggerer 上维持上述状态并异步等待。完整示例 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] → [四个 deferrable 变体] → copy_object → copy_prefix → file_transform → sensor_keys_unchanged → delete_objects → delete_bucket → delete_bucket_2其中删除类任务delete_objects、delete_bucket、delete_bucket_2均设置trigger_rule TriggerRule.ALL_DONE确保无论测试主体成败都会清理 AWS 资源DAG 末尾还挂载了watcher()任务以在包含 tearDown 任务使用 trigger rule时正确标记整体成功/失败。该示例同时展示了s3://URL 与分离 bucket/key 两种参数书写风格的混用如S3FileTransformOperator使用完整 URL其他运算符使用分离参数。小结与最佳实践幂等设计S3CreateBucketOperator、S3DeleteBucketOperator及三个标签运算符在执行前都会先探测桶是否存在天然适合在周期性 DAG 中安全重跑权限最小化复制类运算符要求连接同时具备源与目标两侧权限使用 KMS 加密时必须成对提供kms_key_id与kms_encryption_type且 IAM 需授予密钥使用权限成本意识S3KeySensor每键一次 API 调用批量等待优先考虑wildcard_matchcheck_fn的组合以减少探测次数worker 效率长轮询类等待优先开启deferrableTrue将轮询负担转移到 triggerer前提是部署中已运行 triggerer状态保持S3KeysUnchangedSensor依赖跨 poke 维护对象集合状态切勿使用 reschedule 模式测试驱动仓库中 example_s3.py 是一份可运行的完整参考实现官方所有运算符与传感器代码集中在 operators/s3.py 与 sensors/s3.py底层 API 封装位于 hooks/s3.py需要更细粒度控制如分页、requester-pays、跨连接读写时可直接复用S3Hook。参考官方操作指南providers/amazon/docs/operators/s3/s3.rst系统测试示例providers/amazon/tests/system/amazon/aws/example_s3.py前置任务说明providers/amazon/docs/_partials/prerequisite_tasks.rst运算符源码providers/amazon/src/airflow/providers/amazon/aws/operators/s3.py传感器源码providers/amazon/src/airflow/providers/amazon/aws/sensors/s3.pyS3 Hook 源码providers/amazon/src/airflow/providers/amazon/aws/hooks/s3.py同目录其他主题providers/amazon/docs/operators/s3/glacier.rstS3 Glacier 归档相关运算符与传感器【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表