十年匠心定制 · 商业建站与技术教学双线并行 咨询热线:400-886-1026 service@lmnt.cn
ARTICLE DETAIL

资讯详情

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

Apache Airflow Scheduler 调度性能优化:触发规则上游任务实例计数的单次查询与复用

Apache Airflow Scheduler 调度性能优化:触发规则上游任务实例计数的单次查询与复用 Apache Airflow Scheduler 调度性能优化触发规则上游任务实例计数的单次查询与复用【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow导读在 Apache Airflow 中调度器Scheduler每次调度心跳都会评估大量任务实例Task Instance的依赖是否满足其中「触发规则」Trigger Rule评估需要统计每个下游任务的上游任务实例数量。当 DAG 中出现「一个被映射展开mapped / expand的上游任务喂养大量下游任务」的扇出fan-out结构时旧实现会对每个下游任务各发一次上游计数 SQL 查询形成可观的数据库往返开销。本仓库中编号为 67672 的改进条目 记录了该问题的优化在同一个调度 pass 内共享相同直接上游的任务不再各自查询计数结果只计算一次并在多个下游之间复用。读完本文你将理解该优化的实现机制、缓存键设计与正确性保障并能判断自己的 DAG 能否从中受益。优化背景一个映射上游喂饱大量下游的典型场景在 Airflow 3 中映射任务mapped task通过expand()或partial().expand()展开成大量任务实例。一个常见的高扇出结构是from airflow.decorators import task from airflow.models.baseoperator import chain task def fetch_ids() - list[int]: return list(range(1000)) task(trigger_ruleall_success) def process_one(item_id: int) - None: ... with DAG(fan_out_dag, scheduleNone, start_datedatetime(2024, 1, 1)) as dag: ids fetch_ids() process_one.expand(item_idids) # 或者 ids process_one.expand(...)这里fetch_ids是映射上游process_one展开成上千个实例。调度器评估每一个process_one实例的all_success触发规则时都需要回答一个问题“上游fetch_ids在当前 DagRun 中到底有多少个实例”——因为all_success只有在success upstream时才算满足见下文状态统计逻辑。在优化之前这个计数查询会为每一个下游任务实例各执行一次对于ids p1 p2 ... pN这类结构或「一个映射上游直接 fan-out 到 N 个下游」的 DAG调度一个 pass 就会产生 N 次内容完全相同的SELECT count(...) GROUP BY task_id查询。这正是 newsfragment 中描述的“mapped upstream feeds many downstream tasks”场景数据库往返次数随下游数量线性增长。触发规则评估机制回顾触发规则评估由 TriggerRuleDep 完成。它是调度器在判断“上游任务是否允许当前任务实例运行”时使用的任务依赖之一类注释明确写道Determines if a tasks upstream tasks are in a state that allows a given task instance to run.评估流程大致如下trigger_rule_dep.py若任务没有上游upstream_task_ids为空直接判定通过若触发规则是always同样直接通过否则进入_evaluate_trigger_rule统计上游任务实例的各状态数量再结合触发规则推导当前实例应该运行、跳过还是标记upstream_failed。上游状态的汇总结构是_UpstreamTIStatestrigger_rule_dep.py它用Counter统计出success、skipped、failed、upstream_failed、removed、done以及针对 setup 任务的success_setup/skipped_setup。随后的规则推导trigger_rule_dep.py覆盖all_success、all_failed、one_success、one_failed、one_done、none_failed、none_failed_min_one_success、none_skipped、all_skipped、all_done_min_one_success、all_done_setup_success等触发规则。关键在于上游实例数量upstream的获取方式trigger_rule_dep.py若所有直接上游都是“简单任务”不涉及任务级或任务组级映射即get_needs_expansion()均为 False则upstream len(upstream_tasks)根本不需要查数据库否则执行一条聚合查询SELECT task_instance.task_id, count(task_instance.task_id) ... GROUP BY task_instance.task_id按task_id统计当前dag_id run_id下每个上游任务的实例数再求和得到upstream。这条按task_id分组的计数 SQL 就是本次优化针对的目标查询。优化核心单调度 pass 内的去重缓存旧实现对每个下游任务实例都会独立执行上面的计数 SQL。但源码注释trigger_rule_dep.py指出一个关键观察In the simple case,_iter_upstream_conditionsemits exactlytask_id IN (upstream_task_ids)... That predicate, and therefore the resulting counts, are identical for every downstream that shares the same set of direct upstreams, so we memoize them on the DepContext and run the query once per pass instead of once per downstream.也就是说当任务不在映射任务组中时计数查询的谓词恰好是task_id IN (直接上游任务ID集合)这个谓词对“共享同一组直接上游”的所有下游完全一致因此结果可以在一次调度 pass 内被这些下游复用。具体实现分三步第一步定义缓存键并命中查询trigger_rule_dep.pycache_key: tuple[str, str, frozenset[str]] | None None task_id_counts: list[tuple[str, int]] | None None if task.get_closest_mapped_task_group() is None: cache_key (ti.dag_id, ti.run_id, frozenset(upstream_tasks)) task_id_counts dep_context.upstream_task_id_counts.get(cache_key) if task_id_counts is None: task_id_counts [ ... session.execute(select(TaskInstance.task_id, func.count(...))...group_by(...)) ... ] if cache_key is not None: dep_context.upstream_task_id_counts[cache_key] task_id_counts缓存键由三元组构成(dag_id, run_id, frozenset(直接上游任务ID集合))。由于键中包含完整的上游集合上游集合不同的下游会得到不同的缓存键各自独立查询、互不串扰。第二步缓存放在 DepContext 上生命周期为一个调度 pass。DepContext 新增字段upstream_task_id_counts: dict[tuple[str, str, frozenset[str]], list[tuple[str, int]]] attr.ib(factorydict, reprFalse)字段文档说明了三个设计要点生命周期与finished_tis一致都是一次调度 pass——调度器每一轮调度循环会构造新的 DepContext旧缓存随之丢弃天然避免跨 pass 读到过期数据不是调用方传入的快照而是在依赖评估过程中逐步填充刻意保持initTrue因为 TaskInstance.are_dependencies_met 对每个UP_FOR_RESCHEDULE的任务实例都会用attrs.evolve重建 DepContext而attrs.evolve只携带__init__接受的字段。如果这个字段设为initFalse每个重调度实例都会拿到一个全新的空字典——既读不到已缓存结果也无法为其他实例预热缓存优化就形同虚设了。这与测试test_memo_survives_the_up_for_reschedule_dep_context_evolve验证的场景直接对应reschedule 模式传感器sensor在映射上游背后大量扇出正是这个缓存想要消除的查询形态。第三步mapped task group 场景不缓存。如果当前任务位于映射任务组中get_closest_mapped_task_group() is not None每个实例依赖的上游 map index 各不相同计数谓词是逐实例构造的见_iter_upstream_conditions中按 map index 范围/集合/精确值生成的多种条件无法安全复用因此该分支被显式排除在缓存之外。缓存失效正确性高于性能缓存要正确就必须在“计数结果可能变化”的瞬间失效。源码注释trigger_rule_dep.py明确了两类情况需要失效映射任务的实例数量在 pass 中途发生变化——包括「展开尚未展开的映射任务」和「_revise_map_indexes_if_mapped使已展开的映射任务实例数增长」。因为行数变了后面评估的下游必须看到新计数无需失效仅状态变化任务完成、实例被标记REMOVED不改变行数缓存可继续复用。失效的落点有两个都位于 DagRun._get_ready_tis调度 pass 的核心入口之一在_expand_mapped_task_if_needed展开映射任务产生新实例后调用dep_context.invalidate_upstream_task_id_counts()在_revise_map_indexes_if_mapped返回新实例后同样调用该方法。对应的实现即 DepContext.invalidate_upstream_task_id_countsdef invalidate_upstream_task_id_counts(self) - None: Drop the memoized trigger-rule upstream counts... self.upstream_task_id_counts.clear()这套“先缓存、后失效”的机制在测试test_revise_growing_a_mapped_upstream_clears_memo_within_pass中被精确验证驱动_get_ready_tis按固定顺序[d1, 映射实例, d2]执行d1 基于增长前的实例填充缓存映射任务随后被 revise 并增长d2 必须重算——因此计数查询总共执行 2 次。如果没有_get_ready_tis中的缓存清理d2 会读到陈旧值且查询只执行 1 次看似更“快”实则结果错误。测试验证查询次数从 N 降到 1该优化配套的测试集中在 TestTriggerRuleUpstreamCountMemo 类中测试辅助_count_upstream_count_queriestest_trigger_rule_dep.py通过监听 SQLAlchemy 引擎的after_cursor_execute事件精确匹配形如SELECT task_instance.task_id, count(task_instance.task_id) ... GROUP BY task_instance.task_id的查询并计数——只统计触发规则计数 SQL不误伤其他查询。各测试用例对应的行为契约测试场景断言test_memoized_across_downstreams_sharing_upstream4 个普通下游共享同一个映射上游计数查询只执行1 次test_memoized_count_value_is_correct3 个上游实例只有 2 个成功1 个 RUNNINGall_success必须不满足查询 1 次且缓存值是真实计数不是“存在即正确”test_distinct_upstream_sets_are_not_collapsed两个下游的上游集合不同各查各的共2 次test_revise_growing_a_mapped_upstream_clears_memo_within_pass映射上游 pass 中途 revise 增长查询2 次增长后强制重算test_memo_survives_the_up_for_reschedule_dep_context_evolve4 个UP_FOR_RESCHEDULE下游经attrs.evolve重建 DepContext查询1 次test_evolved_dep_context_shares_the_memo_object固定attrs.evolve的字段透传机制evolve 后共享同一缓存对象这些测试同时覆盖了优化的两个目标面性能面共享上游 → 查询去重与正确性面不同上游不串扰、缓存值真实、增长后强制失效、重调度上下文透传缓存。收益与适用场景从实现看该优化的收益主要体现在减少数据库往返在“一个映射上游或一组相同的直接上游喂养 N 个下游”的 DAG 中触发规则计数 SQL 从 N 次降为 1 次配合 DepContext.ensure_finished_tis 在同一 pass 内只加载一次已完成上游实例列表的机制调度循环对相同输入不再重复扫描同一批数据对 reschedule 传感器扇出场景尤其有效UP_FOR_RESCHEDULE状态的实例在调度循环中反复进入依赖评估缓存让它们共享同一份计数结果无需用户改 DAG这是调度器内部优化任何expand()扇出结构自动受益。需要注意的限制条件均来自源码注释与测试仅当任务不在映射任务组中时才启用缓存简单 case映射任务组内的逐实例 map index 谓词不缓存缓存键区分dag_id run_id 上游集合上游集合不同的下游不会合并查询这是正确性取舍缓存生命周期为一个调度 pass不跨 pass 复用映射任务实例数在 pass 内增长时会触发缓存失效并重算因此这类场景的查询次数仍可能多于 1 次——但换来的是计数永远正确。如果你正在排查调度器数据库压力且 DAG 中存在大量“映射上游 → 多下游”的扇出结构本优化见 67672.improvement.rst正是针对该形态的定向改进相关实现与测试可作为理解调度器依赖评估与查询去重模式的参考样例。【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表