ARTICLE DETAIL

资讯详情

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

Apache Airflow 新增 triggerer.batch_trigger_creation_duration 指标:原理、源码实现与监控实践

Apache Airflow 新增 triggerer.batch_trigger_creation_duration 指标:原理、源码实现与监控实践 Apache Airflow 新增 triggerer.batch_trigger_creation_duration 指标原理、源码实现与监控实践【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow本文以 Apache Airflow 仓库中的变更记录airflow-core/newsfragments/68521.feature.rst为线索深入讲解 Triggerer 组件新增的triggerer.batch_trigger_creation_duration指标它衡量的是单个 trigger runner 主循环周期内创建所有待处理pending触发器所花费的总时间。读完本文你将理解该指标在触发器生命周期中的测量位置、源码实现细节、如何借助 metrics.rst 中的配置消费它以及如何用它定位 Triggerer 的延迟瓶颈。一、背景Triggerer 与触发器批量创建在 Apache Airflow 中Triggerer是专门为可延迟任务deferrable tasks服务的独立组件。当一个任务实例通过defer将自身挂起后其对应的触发器Trigger会被交给 Triggerer 进程托管——触发器本质上是基于 asyncio 的异步生成器由 Triggerer 驱动执行在条件满足时向任务实例发送TriggerEvent使其恢复运行。triggerer.batch_trigger_creation_duration指标正是在这一背景下诞生的触发器并非一个一个地立即创建而是由 trigger runner 在每个主循环周期内批量创建。该指标专门测量从一个周期开始到把所有待创建触发器全部实例化完毕所消耗的总时长是评估 Triggerer 创建触发器开销的直接依据。二、指标定义测量什么、何时测量变更记录原文68521.feature.rst给出如下定义Added thetriggerer.batch_trigger_creation_durationmetric, which measures the total time spent creating pending triggers during a single trigger runner cycle.拆解该定义有三个关键语义total time总时间包含一个周期内所有待创建触发器pending triggers的完整创建过程而非单个触发器。a single trigger runner cycle单个 trigger runner 周期指 TriggerRunner 主事件循环的一次迭代。creating pending triggers创建待处理触发器特指将 DB 中请求创建的触发器反序列化、解密参数、实例化并注册为 asyncio 任务的过程不包含运行阶段。一个重要的行为细节是只有当本周期内确实有触发器需要创建to_create队列非空时该指标才会被发射。若一个周期内没有任何待创建触发器则不会产生该指标的数据点避免在空闲周期输出大量无意义的零值样本。三、源码实现create_triggers 与计时逻辑该指标在 triggerer_job_runner.py 的TriggerRunner.create_triggers()方法中实现对应源码 L1354-L1441async def create_triggers(self): Drain the to_create queue and create all new triggers that have been requested in the DB. # Emit batch creation duration only when triggers were processed. has_work bool(self.to_create) if has_work: creation_start time.monotonic() while self.to_create: await asyncio.sleep(0) context: Context | None None workload self.to_create.popleft() trigger_id workload.id if trigger_id in self.triggers: self.log.warning(Trigger %s had insertion attempted twice, trigger_id) continue # ... 加载触发器类、解密并反序列化 kwargs、实例化触发器 ... self.triggers[trigger_id] { task: asyncio.create_task( self.run_trigger(trigger_id, trigger_instance, workload.timeout_after, context), nametrigger_name, ), # ... } if has_work: stats.timing( triggerer.batch_trigger_creation_duration, (time.monotonic() - creation_start) * 1000, tagsprune_dict({team_name: self.team_name}), )从源码可以提炼出以下实现要点计时起点在进入创建循环前通过has_work bool(self.to_create)判断本周期是否有工作仅当有工作时才用time.monotonic()记录creation_start。批量处理to_create是一个deque源码 L1208循环内使用popleft()逐个弹出待创建的触发器 workload 进行处理处理内容包括检查触发器是否已被创建防止重复插入重复时会告警并跳过通过get_trigger_by_classpath(workload.classpath)加载触发器类源码 L1371加载失败则记入failed_triggers解密encrypted_kwargs并反序列化参数Trigger._decrypt_kwargssmart_decode_trigger_kwargs源码 L1382-L1390为资产事件触发器注入AssetStateStoreAccessors源码 L1413-L1416通过asyncio.create_task将每个触发器包装为协程任务并登记进self.triggers源码 L1426-L1434。计时终点循环结束后若has_work为真用time.monotonic()计算差值乘以 1000 转换为毫秒通过stats.timing发射。值得注意的实现细节是循环体内部的await asyncio.sleep(0)源码 L1363、L1379它主动让出事件循环保证加载触发器类这类可能较重的操作不会长时间阻塞同进程内其他正在运行的触发器协程从而也使得本指标测量到的时间更贴近真实的创建开销 事件循环调度开销。四、指标在主循环中的位置create_triggers()并不是独立运行的而是 TriggerRunner 主异步循环arun()的一部分源码 L1247-L1286其调用点位于 L1275while not self.stop: # ... finished_ids await self.cleanup_finished_triggers() # This also loads the triggers we need to create or cancel await self.sync_state_to_supervisor(finished_ids) await self.create_triggers() await self.cancel_triggers() # Sleep for a bit, or exit early if stop is requested. with anyio.move_on_after(1): await stop_event.wait()主循环每个周期依次执行清理已结束触发器 → 与 supervisor 同步状态顺带加载待创建/待取消的触发器→批量创建触发器→ 取消触发器 → 最多等待 1 秒后进入下一周期。因此一个 trigger runner cycle的边界是清晰的create_triggers()每次被调用即代表一个周期的批量创建阶段指标数据点也按周期粒度产生。工作负载的产生路径在build_trigger_workloads/sync_state_to_supervisor等逻辑中源码 L933-L995、L1565supervisor 将新触发器请求写入to_create队列runner 在下一周期消费。该指标与同文件中的triggerer.trigger_queue_delay源码 L1420-L1424形成互补——后者衡量workload 入队时刻到实际开始创建之间的排队延迟前者衡量创建本身的耗时。五、指标标签与配套配置发射该指标时携带的标签由prune_dict({team_name: self.team_name})生成team_name来自 supervisor 下发的StartTriggerer消息源码 L1331。当 Triggerer 以多租户方式按团队team划分运行范围时不同团队的指标可以按team_name维度区分当未配置团队作用域时该标签为空。同一文件中其他指标如triggerer_heartbeat、triggers.succeeded、triggers.failed、triggers.blocked_main_thread也遵循这一模式。指标过滤allow / block 列表如果你不希望把全部指标都发送到后端可以通过 metrics.rst 中描述的 allow/block 列表控制。列表是逗号分隔的正则表达式集合匹配指标名称的任意位置例如[metrics] metrics_allow_list scheduler,executor,dagrun,pool,triggerer,celery[metrics] metrics_block_list scheduler,executor,dagrun,pool,triggerer,celery使用triggerer前缀即可同时放行或屏蔽triggerer.batch_trigger_creation_duration、triggerer.trigger_queue_delay、triggerer_heartbeat等指标。注意若同时设置了两个列表block 列表会被忽略allow 列表优先。若需要更精确的匹配可用^triggerer.batch_trigger_creation_duration$这类锚定正则只针对本指标。指标重命名若下游监控系统需要统一命名规范可以配置stat_name_handlermetrics.rst对指标名做转换def my_custom_stat_name_handler(stat_name: str) - str: return stat_name.lower()[:32]该函数接收原始指标名返回转换后的名称可用于统一大小写、截断或添加前缀。六、测试验证与使用建议仓库单元测试 test_triggerer_job.py 对create_triggers()有大量直接覆盖如 L847-L877 验证BaseEventTrigger的asset_state_store注入等这些测试通过构造RunTriggerworkload 塞入runner.to_create后调用await runner.create_triggers()间接验证了批量创建路径的稳定性也说明该指标的计时逻辑位于经过充分测试的核心代码路径上。在实际运维中建议将triggerer.batch_trigger_creation_duration与以下指标配合观察指标含义观察角度triggerer.batch_trigger_creation_duration单周期批量创建触发器总耗时毫秒创建阶段开销triggerer.trigger_queue_delayworkload 入队到开始创建之间的排队延迟毫秒排队与调度延迟triggerer_heartbeatTriggerer 心跳计数进程存活triggers.succeeded/triggers.failed触发器成功/失败次数触发器质量triggers.blocked_main_thread主线程被阻塞告警次数事件循环健康度若batch_trigger_creation_duration持续偏高相对触发器数量而言可以推断触发器类加载classpath 解析、kwargs 解密与反序列化、或asyncio.create_task注册过程存在瓶颈若伴随triggers.blocked_main_thread上升则更可能是有触发器代码阻塞了事件循环此时可参考 check-health.rst 中描述的 Triggerer 健康检查与多实例心跳信息进一步定位是单实例资源不足还是触发器实现质量问题。七、小结triggerer.batch_trigger_creation_duration是 Airflow 在触发器批量创建路径上新增的可观测性指标它的实现完整落在 triggerer_job_runner.py 的create_triggers()方法中以time.monotonic()记录单周期批量创建总耗时仅在存在待创建触发器时发射单位为毫秒并携带team_name标签以支持多团队维度聚合。配合triggerer.trigger_queue_delay等既有指标以及[metrics]段的 allow/block 列表配置运维团队可以完整刻画触发器从入队、排队到创建完成的端到端耗时从而精准定位 Triggerer 组件的性能瓶颈。【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表