如何计算Apache Airflow Sensor的总执行时间
sensor_job = PythonSensor( task_id='sensor_id', python_callable=call_jobsensor, poke_interval=10, timeout=7 * 60, mode='reschedule', ) def call_jobsensor(): # 探测业务逻辑 # 执行条件检查 return True/False
需求说明
需要统计上述Sensor任务从首次启动到最终结束的全周期总时长,任务结束包含两种判定场景:
- 探测条件满足,返回
True正常标记为任务成功 - 运行时长达到
timeout配置阈值,抛出超时异常标记为任务失败
注意:不要尝试在call_jobsensor内部加计时逻辑,reschedule模式下Sensor每次poke失败后会释放worker资源、等待间隔时间后重新调度,单次callable执行只覆盖单轮探测的耗时,且每次poke都是独立进程运行,内存变量无法跨轮次保留,完全无法统计多轮poke+等待间隔的全周期时长。
方案1:基于Airflow原生StatsD能力实现(推荐)
不需要修改任何业务探测逻辑,Airflow默认的StatsD指标已经完全覆盖需求:
- 先确认Airflow配置开启StatsD:在
airflow.cfg中设置statsd_on = True,并配置好StatsD服务的地址、端口即可。 - 直接使用原生上报的任务时长指标
airflow.tasks.duration.<dag_id>.<task_id>,这个指标统计的就是任务实例从首次启动执行,到最终进入成功/失败终态的完整墙钟时长,自动覆盖reschedule模式下的poke等待间隔、多轮调度的开销,不管是探测成功还是超时终止,都会准确上报总时长。 - 如果需要区分两种结束状态的时长,搭配
airflow.tasks.finish.<dag_id>.<task_id>.<status>指标关联即可:status=success对应探测成功场景,status=failed对应超时(超时会抛出AirflowSensorTimeout异常,任务标记为失败)场景,按任务实例维度聚合就能拆分统计两类场景的时长数据。
方案2:自定义回调实现(无StatsD场景适用)
如果没有部署StatsD链路,可以用Airflow的任务生命周期回调实现,同样不需要侵入探测逻辑:
- 给Sensor配置三类生命周期回调:
on_execute_callback:任务第一次启动执行时触发,在回调里把当前时间戳写入XCom,key设为sensor_start_ts即可on_success_callback:探测成功、任务进入成功终态时触发,从XCom读取之前存的启动时间戳,和当前时间做差得到总时长on_failure_callback:任务超时、进入失败终态时触发,同样读取XCom里的启动时间戳计算总时长
- 计算得到的总时长可以直接打日志、写入业务库或者做自定义上报。
- 如果是Airflow 2.x版本,回调入参里的
task_instance对象自带duration属性,这个值就是Airflow已经计算好的任务全周期运行时长,连手动算时间差的步骤都可以省掉,直接取值即可。
避坑提示
不要在pre_execute/post_execute方法里写计时逻辑,reschedule模式下每一轮poke都会触发这两个方法,最终拿到的还是单轮探测的分段耗时,不是全周期总时长。
内容的提问来源于stack exchange,提问作者Gitesh kumar Jha
相关产品推荐
相关产品推荐

