Airflow SqlSensor数据库恢复状态检测配置问题求助
Airflow SqlSensor 问题修复方案
问题分析与修复
1. 检测间隔异常(每秒检查而非5分钟)
- 原因:原代码中
timeout=120秒(2分钟)小于poke_interval=300秒(5分钟),导致Sensor在超时前会持续高频轮询,忽略了poke_interval的配置。 - 修复:将
timeout调整为大于poke_interval的合理值,例如1小时(3600秒),确保Sensor有足够时间按设定间隔执行检查。
2. catchup参数未生效,出现回填
- 原因:部分Airflow版本要求在DAG实例化时显式指定
catchup=False才能生效,且原代码schedule_interval设置为10分钟,与需求的5分钟不符。 - 修复:在
DAG初始化参数中显式添加catchup=False,同时将schedule_interval修改为*/5 * * * *(每5分钟执行一次)。
3. 根据SQL布尔返回值标记任务成功/失败
- 原因:SqlSensor默认以查询返回非空结果为成功条件,需要自定义判断逻辑匹配
pg_is_in_recovery()的布尔返回值。注意查询结果是元组列表格式(如[(True,)]),需正确提取布尔值。 - 修复:通过
success_fn参数定义成功条件,当返回值为True时标记任务成功,否则任务失败。
完整修复代码
from airflow.sensors.sql import SqlSensor from airflow import DAG from datetime import datetime, timedelta default_args = { 'start_date': datetime(2023, 7, 15), } dag = DAG( 'database_monitor', default_args=default_args, schedule_interval='*/5 * * * *', # 每5分钟执行一次 catchup=False, # 关闭任务回填 ) check_db_alive = SqlSensor( task_id="check_db_alive", conn_id="evergreen", sql="SELECT pg_is_in_recovery()", success_fn=lambda records: records[0][0] is True, # 返回True时任务成功 poke_interval=60 * 5, # 5分钟检查间隔 timeout=60 * 60, # 超时时间1小时 mode="reschedule", # 等待期间释放worker资源 dag=dag ) check_db_alive
额外说明
mode="reschedule"适合长间隔Sensor任务,能在等待下一次检查时释放worker资源,避免资源占用。- 若需返回
False时立即失败不重试,可在default_args或Sensor参数中添加retries=0。
内容的提问来源于stack exchange,提问作者moth
相关产品推荐
相关产品推荐

