
Apache Airflow Tasks 深度指南从任务依赖到重试策略与超时控制【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow本篇技术指南以 Apache Airflow 官方核心概念文档 tasks.rst 为骨架系统讲解 Airflow 中任务Task这一基本执行单元如何定义任务、声明上下游依赖、理解 Task Instance 生命周期与状态机、配置超时与重试策略以及如何借助特殊异常、executor_config 等机制精确控制任务行为。读完本文你将掌握编写健壮、可控、可运维 DAG 任务所需的全部核心知识并了解其底层实现依据对应源码见 state.py、retry_policy.py。什么是 TaskAirflow 中的基本执行单元在 Airflow 中Task 是最基本的执行单元。Task 被组织进 DAG有向无环图中并通过在任务之间设置 upstream上游与 downstream下游依赖来表达它们应当执行的顺序。Airflow 中一共有三种基本的 TaskOperators操作符预定义好的任务模板可以快速拼接出 DAG 的大部分组成部分例如执行 Bash 命令、传输文件、调用云服务 API 等详见 operators.rst。Sensors传感器Operator 的一个特殊子类其全部职责就是等待某个外部事件发生例如等待某个文件出现、等待某个 API 可用详见 sensors.rst。task装饰的 TaskFlow 任务把一个自定义的 Python 函数包装成 Task详见 taskflow.rst。从实现角度看这三种形式在内部全部都是BaseOperator的子类因此 Task 与 Operator 的概念在某种程度上可以互换。但更准确的思考方式是Operator 和 Sensor 是模板template当你在 DAG 文件中调用它们一次就产生了一个 Task任务实例化的定义。关系Relationships如何声明任务依赖使用 Task 的关键在于定义它们彼此之间的关系即依赖dependenciesAirflow 中称为upstream上游与downstream下游任务。通常的做法是先声明所有 Task再声明它们的依赖关系。注意所谓 upstream 任务是指直接排在另一个任务前面的那个任务旧文档中曾称之为 parent task。需要留意的是这一概念并不描述任务层级中更高层级的任务即并非该任务的间接祖先。downstream 任务的定义同理它必须是另一个任务的直接后继。两种声明依赖的方式Airflow 提供了两种声明依赖的语法方式一位运算操作符与推荐first_task second_task [third_task, fourth_task]方式二显式方法set_upstream与set_downstreamfirst_task.set_downstream(second_task) third_task.set_upstream(second_task)这两种写法实现的效果完全相同但官方推荐优先使用位运算操作符因为它在大多数场景下可读性更强——箭头方向直观地表达了数据/执行的流动方向。默认依赖语义与高级控制默认情况下一个 Task 会在其所有upstream 任务成功后运行。但 Airflow 提供了大量方式修改这一行为例如引入分支branching只等待部分 upstream 任务例如通过TriggerRule.ONE_SUCCESS、ALL_DONE等触发规则根据当前运行在历史中的位置如是否为回填、是否为补数据运行改变行为。这些高级控制方式参见 dags.rst 中的控制流章节trigger rules以及 backfill.rst。此外需要注意Task 默认不向彼此传递信息彼此完全独立运行。如果需要在任务之间传递数据应当使用 XComs跨任务通信机制详见 xcoms.rst。Task Instance任务的一次具体运行正如 DAG 每次运行会被实例化为Dag RunDAG 下的每个 Task 也会被实例化为Task Instance任务实例。一个 Task Instance 是该任务在给定 DAG因而也是给定数据区间 data interval下的一次具体运行。它同时也是**拥有状态state**的任务表示反映它处于生命周期的哪个阶段。当任何自定义 TaskOperator运行时它都会获得一份 Task Instance 的拷贝除了可以检查任务元数据外Task Instance 还携带诸如 XComs 读写之类的方法。Task Instance 的完整状态集Task Instance 可能的状态如下对应源码枚举定义见 state.py 中的TaskInstanceState其中IntermediateTIState表示尚未进入终态/运行态的中间状态TerminalTIState表示终态状态含义none任务尚未被排队执行其依赖尚未满足由调度器创建但尚未运行时使用None表示scheduled调度器已判定任务的依赖满足应当运行IntermediateTIState.SCHEDULEDqueued任务已分配给某个 Executor正在等待 workerIntermediateTIState.QUEUEDrunning任务正在 worker 上运行或在 local/synchronous executor 上运行success任务运行完毕且无错误终态TerminalTIState.SUCCESSrestarting任务在运行时被外部请求重启例如运行时被 clearIntermediateTIState.RESTARTINGfailed任务执行期间出错而运行失败终态TerminalTIState.FAILEDskipped任务因分支branching、LatestOnly 等机制被跳过终态TerminalTIState.SKIPPEDupstream_failed某个上游任务失败且触发规则Trigger Rule要求等待它终态TerminalTIState.UPSTREAM_FAILEDup_for_retry任务失败但还有重试次数将被重新调度IntermediateTIState.UP_FOR_RETRYup_for_reschedule任务是一个处于reschedule模式的 SensorIntermediateTIState.UP_FOR_RESCHEDULEdeferred任务已**延迟deferred**到一个触发器trigger等待异步事件IntermediateTIState.DEFERRED详见 deferring.rstawaiting_input任务是一个Human-in-the-loop人机交互任务等待人工响应由调度器管理既不占用 worker 槽位也不占用 triggererIntermediateTIState.AWAITING_INPUT详见 hitl.rstremoved自运行开始后任务已从 DAG 中消失终态TerminalTIState.REMOVED在理想情况下任务应当按如下路径流转none→scheduled→queued→running→success见上图生命周期示意。关系术语upstream/downstream 与 previous/next 的区别对于任意一个 Task Instance它与其它实例之间存在两类关系第一类upstream 与 downstream 任务task1 task2 task3当 DAG 运行时会为这些彼此互为上下游的任务创建实例这些实例共享同一个数据区间。第二类同一任务在不同数据区间上的实例这些实例来自同一 DAG 的其他运行Airflow 称之为previous上一个与next下一个——这与 upstream/downstream 是完全不同的关系维度注意一些较老版本的 Airflow 文档可能仍用 previous 表示 upstream。如果发现这类用法可以协助社区修正文档。超时控制Timeouts如果希望给任务设置最大运行时长可以设置任务的execution_timeout属性值为一个datetime.timedelta。该设置适用于所有Airflow 任务包括 Sensor。execution_timeout控制的是每一次执行所允许的最大时间一旦超时任务会超时并抛出AirflowTaskTimeout异常。此外Sensor 还有一个额外的timeout参数仅对reschedule模式下的 Sensor 有效。timeout控制的是Sensor 最终成功所允许的最大总时间一旦超过将抛出AirflowSensorTimeoutSensor立即失败且不再重试。综合示例SFTPSensor 的超时与重试组合以下SFTPSensor示例清晰地说明了这两层超时的配合代码取自原文档参数逐条解释如下sensor SFTPSensor( task_idsensor, path/root/test, execution_timeouttimedelta(seconds60), timeout3600, retries2, modereschedule, )该 Sensor 处于reschedule模式即周期性地执行、重新调度直到成功为止各参数行为如下每次探测pokeSFTP 服务器最多允许 60 秒execution_timeout。如果一次探测超过 60 秒抛出AirflowTaskTimeout此时允许重试最多重试 2 次retries。从第一次执行开始到最终成功即文件root/test出现为止总共最多允许 3600 秒timeout。如果 3600 秒内文件始终未出现抛出AirflowSensorTimeout此时不再重试。如果在 3600 秒窗口内因其他原因失败如网络中断仍可最多重试 2 次retries。重试不会重置timeout——它依然总共只有 3600 秒用于成功。SLAs 的历史变更Airflow 2 中基于 SLA 的功能在Airflow 3.0 中已被移除并在 Airflow 3.1 中被Deadlines Alerts截止时间告警取代。当前项目中如需使用类似能力请参考 deadline-alerts.rst。特殊异常Special Exceptions从任务代码内部控制状态如果希望从自定义 Task/Operator 代码内部控制任务的最终状态Airflow 提供了两个可主动抛出的特殊异常定义于 exceptions.py导出名AirflowSkipException、AirflowFailExceptionAirflowSkipException将当前任务标记为skipped跳过。AirflowFailException将当前任务标记为failed失败并且忽略剩余的任何重试次数。这两个异常非常适合代码对自身环境有额外认知、希望更快地失败/跳过的场景例如已知本次没有可用数据时直接跳过任务检测到 API Key 无效时快速失败因为重试不会修复密钥问题没必要浪费重试次数。重试策略Retry Policies按异常类型精细化控制重试默认情况下Airflow 以固定的次数和固定的延迟重试失败任务而不区分错误类型。Retry Policy重试策略允许你以任务或 Operator 上的一个参数按异常类型配置逐异常per-exception的重试行为无需修改任务代码。定义一个策略并应用到任务策略由将异常类型映射到动作的规则组成完整可运行的示例见 example_retry_policy.py。策略定义与使用如下# 定义策略 from airflow.sdk import DAG, ExceptionRetryPolicy, RetryAction, RetryRule, task API_RETRY_POLICY ExceptionRetryPolicy( rules[ RetryRule( exceptionrequests.exceptions.HTTPError, actionRetryAction.RETRY, retry_delaytimedelta(minutes5), reasonRate limit, backing off, ), RetryRule( exceptiongoogle.auth.exceptions.RefreshError, actionRetryAction.FAIL, reasonAuth failure, not retryable, ), RetryRule( exceptionConnectionError, actionRetryAction.RETRY, retry_delaytimedelta(seconds30), ), ], ) # 应用到任务 with DAG( dag_idexample_retry_policy, scheduleNone, catchupFalse, tags[example, retry_policy], ): task(retries5, retry_delaytimedelta(minutes1), retry_policyAPI_RETRY_POLICY) def call_external_api(): import requests response requests.get(https://api.example.com/data) response.raise_for_status() return response.json() call_external_api()工作原理策略在 worker 进程中求值策略运行在任务 worker 进程中绝不在调度器中位于捕获异常与决定任务下一状态之间。每次策略决策都会记录在任务日志中格式为Retry policy decision actionaction reasonreason当任务失败时策略对异常求值并返回三种动作之一对应源码 retry_policy.py 中的RetryAction枚举RETRY重试任务可选择自定义延迟以覆盖retry_delay。重试仍受任务的retries计数约束——策略可以让任务更早失败但不能超出配置的最大次数。FAIL立即失败跳过剩余的重试。DEFAULT回退到标准重试逻辑即retries次数与retry_delay。规则按顺序求值第一个匹配的规则生效若没有规则匹配策略返回DEFAULT即标准重试行为。异常匹配规则异常类型既可以指定为 Python 类也可以指定为点分导入路径字符串例如requests.exceptions.HTTPError。字符串路径会在DAG 解析时进行校验不含点号的路径会立即抛出ValueError无法解析的路径会产生警告见 retry_policy.py 中RetryRule.__post_init__的校验逻辑。默认情况下规则使用isinstance匹配因此针对OSError的规则也会匹配ConnectionError其子类。可以通过match_subclassesFalse改为精确类型匹配RetryRule(exceptionOSError, match_subclassesFalse) # 仅匹配 OSError 本身不含子类另外RetryRule.exception还支持传入列表此时该规则对列表中任一异常匹配即生效源码 retry_policy.py 中的RetryRule文档字符串有明确示例。与既有参数的组合行为当任务设置了retry_policy时各既有参数的行为如下参数设置retry_policy后的行为retries仍然是最大重试次数。策略可以更早失败但不能超过该上限。retry_delay/retry_exponential_backoff/max_retry_delay当策略返回 DEFAULT 或RetryDecision.retry_delay为 None 时使用。on_retry_callback在所有重试包括策略驱动的重试上都会触发。AirflowFailException始终具有最高优先级。该异常从不经过策略求值直接失败。注AirflowSensorTimeout同样总是使任务立即失败重试策略也不会被调用见 retry_policy.py 中evaluate的文档说明。复用策略跨 DAG 共享策略可以定义一次然后通过default_args或共享模块在多个 DAG 间复用# policies.py -- 在任何 DAG 中导入 STANDARD_RETRY_POLICY ExceptionRetryPolicy( rules[ RetryRule(exceptionrequests.exceptions.HTTPError, actionRetryAction.FAIL), RetryRule(exceptionConnectionError, retry_delaytimedelta(seconds10)), ], )动态任务映射Mapped Tasks下的策略策略与动态任务映射dynamic task mapping通过.partial()配合使用。策略按每个映射出的任务实例分别生效——例如 10 个映射实例中的第 2 个命中 FAIL其余 9 个仍独立继续task.partial(retry_policymy_policy).expand(input[1, 2, 3]) def my_mapped_task(input): ...策略在任务级别通过.partial()设置所有映射实例共享同一个策略.expand()不支持按索引变化。但策略的evaluate()方法会接收到异常、try_number和完整 context因此如果需要可以在策略内部实现按索引per-index的分支逻辑。自定义重试策略子类化RetryPolicy对于高级场景可以子类化airflow.sdk.definitions.retry_policy.RetryPolicy并实现evaluate()。当需要检查异常属性状态码、响应头、响应体而声明式的ExceptionRetryPolicy规则无法覆盖时子类化是正确的选择。evaluate()方法接收异常、尝试次数、最大尝试次数以及完整的 Airflow contextdag_run、params等并返回一个RetryDecision该数据类及便捷构造方法fail/retry/default见 retry_policy.py。模式一按 HTTP 状态码路由重试决策并响应 429 的Retry-After头from datetime import timedelta import requests from airflow.sdk import RetryDecision, RetryPolicy class HTTPStatusRetryPolicy(RetryPolicy): 按 HTTP 状态码路由重试决策429 时遵循 Retry-After 头。 def evaluate(self, exception, try_number, max_tries, contextNone): if isinstance(exception, requests.HTTPError) and exception.response is not None: status exception.response.status_code if status 429: # 被限流 -- 遵循 Retry-After 头 retry_after int(exception.response.headers.get(Retry-After, 60)) return RetryDecision.retry(retry_delaytimedelta(secondsretry_after)) if 500 status 600: # 服务器错误 -- 值得重试 return RetryDecision.retry() if 400 status 500: # 客户端错误 -- 不可重试 return RetryDecision.fail(reasonfHTTP {status}) return RetryDecision.default()模式二利用运行 context 决策例如回填时快速失败from airflow.sdk import RetryDecision, RetryPolicy class BackfillAwareRetryPolicy(RetryPolicy): 回填期间快速失败让历史错误立即暴露出来。 def evaluate(self, exception, try_number, max_tries, contextNone): if context and context[dag_run].run_type backfill: return RetryDecision.fail(reasonBackfill run -- not retrying) return RetryDecision.default()提示自定义策略需要支持 DAG 序列化serialize()/deserialize()内置的ExceptionRetryPolicy已提供相应实现子类化时需自行实现参见 retry_policy.py 的接口定义。Task Instance 心跳超时Heartbeat Timeout处理僵尸任务没有哪个系统是完美运行的Task Instance 偶尔也会死亡。Task Instance 可能停留在running状态但其关联的 job 已经不活跃了例如该 TaskInstance 的 worker 内存耗尽被杀。这类任务在旧版中被称为僵尸任务zombie tasks。Airflow 会定期发现它们、进行清理并将 TaskInstance 标记为失败如果还有可用重试次数则触发重试。TaskInstance 心跳超时的常见诱因包括Airflow worker 内存耗尽被 OOMKilledworker 未通过存活探针liveness probe导致系统如 Kubernetes重启了 worker系统如 Kubernetes缩容将 worker 从一个节点迁移到另一个节点。本地复现心跳超时供开发/测试如果需要在本地产复现心跳超时可按下述步骤操作步骤 1设置环境变量也可以等价地修改airflow.cfg中的对应配置项export AIRFLOW__SCHEDULER__TASK_INSTANCE_HEARTBEAT_SEC600 export AIRFLOW__SCHEDULER__TASK_INSTANCE_HEARTBEAT_TIMEOUT2 export AIRFLOW__SCHEDULER__TASK_INSTANCE_HEARTBEAT_TIMEOUT_DETECTION_INTERVAL5三个变量分别含义为TASK_INSTANCE_HEARTBEAT_SEC为任务心跳间隔这里设成 600 秒模拟长时间无心跳TASK_INSTANCE_HEARTBEAT_TIMEOUT为心跳超时阈值设为 2 秒让检测快速触发TASK_INSTANCE_HEARTBEAT_TIMEOUT_DETECTION_INTERVAL为检测间隔每 5 秒扫描一次。步骤 2准备一个耗时约 10 分钟的任务例如from airflow.sdk import dag from airflow.providers.standard.operators.bash import BashOperator from datetime import datetime dag(start_datedatetime(2021, 1, 1), scheduleonce, catchupFalse) def sleep_dag(): t1 BashOperator( task_idsleep_10_minutes, bash_commandsleep 600, ) sleep_dag()运行上述 DAG 并等待一段时间后TaskInstance 将在约task_instance_heartbeat_timeout秒后被标记为失败。Executor 级配置executor_config按任务定制执行环境部分 Executor 允许可选的按任务配置per-task configuration。例如KubernetesExecutor允许为某个任务指定运行它的镜像。这是通过 Task 或 Operator 的executor_config参数实现的。以下示例为将在KubernetesExecutor上运行的任务设置 Docker 镜像MyOperator(..., executor_config{ KubernetesExecutor: {image: myCustomDockerImage} } )executor_config中可设置的项因 Executor 而异请阅读各个 Executor 的专属文档见 executor/index.rst了解可配置项。小结与最佳实践本文围绕 Airflow Task 的完整生命周期梳理了以下核心知识点Task 的三种形式Operator模板、Sensor等待外部事件、taskTaskFlow 函数底层均为BaseOperator子类依赖声明推荐使用/位运算操作符依赖关系基于直接上游/下游Task Instance 状态机从none到scheduled→queued→running→success的理想路径以及up_for_retry、deferred、awaiting_input、removed等全部 15 种状态源码依据见 state.py双层超时execution_timeout单次执行上限超时抛AirflowTaskTimeout可重试与 Sensor 的timeout总成功时限超时抛AirflowSensorTimeout立即失败特殊异常AirflowSkipException跳过与AirflowFailException失败且不重试Retry Policy以ExceptionRetryPolicyRetryRule声明式配置按异常重试支持 RETRY / FAIL / DEFAULT 三动作、isinstance/精确类型匹配、字符串导入路径校验、映射任务支持及自定义RetryPolicy子类示例见 example_retry_policy.py心跳超时僵尸任务的检测与清理机制及其本地复现方法executor_config按任务定制执行环境如 Kubernetes 镜像。在实际编写 DAG 时建议优先采用位运算操作符声明依赖、为所有外部调用类任务显式设置execution_timeout、把可重试/不可重试的异常区分写入 Retry Policy从而让任务具备可预测、可自愈、可运维的高质量行为。【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考