Airflow:传感器设为soft_fail时如何发送邮件告警
解决Airflow传感器soft_fail=True时的告警需求
当传感器设置soft_fail=True时,检测失败会被标记为任务成功,常规的失败回调无法触发告警。以下是几种可行的实现方案:
方法一:自定义传感器,检测失败时主动触发告警
继承Airflow内置传感器,重写poke方法,在检测到目标不存在时先发送告警邮件,再返回检测结果。由于soft_fail=True,任务会被标记为成功并跳过下游,同时告警已触发。
from airflow.sensors.filesystem import FileSensor from airflow.utils.email import send_email class AlertingFileSensor(FileSensor): def poke(self, context): # 执行原生检测逻辑 target_exists = super().poke(context) if not target_exists: # 发送告警邮件 send_email( to='your-alert-email@example.com', subject='Sensor A 检测失败:目标未找到', html_content='<p>Sensor A未检测到目标文件,下游任务B、C已被跳过。</p>' ) return target_exists # 在DAG中使用自定义传感器 sensor_a = AlertingFileSensor( task_id='sensor_a', filepath='/path/to/target/file', soft_fail=True, dag=dag ) sensor_a >> task_b >> task_c
方法二:利用XCom+并行告警任务
保持原传感器配置,在传感器中推送检测结果到XCom,新增一个并行任务读取XCom结果,若检测失败则发送告警。该告警任务设置trigger_rule='all_done',确保无论传感器最终状态如何都能执行。
from airflow.sensors.filesystem import FileSensor from airflow.operators.python import PythonOperator from airflow.utils.email import send_email def sensor_poke_with_xcom(context): # 执行检测并推送结果到XCom sensor = FileSensor(task_id='sensor_a', filepath='/path/to/target/file') target_exists = sensor.poke(context) context['ti'].xcom_push(key='sensor_a_detection_result', value=target_exists) return target_exists def check_sensor_and_send_alert(**context): # 从XCom读取检测结果 detection_result = context['ti'].xcom_pull( key='sensor_a_detection_result', task_ids='sensor_a' ) if not detection_result: send_email( to='your-alert-email@example.com', subject='Sensor A 检测失败:目标未找到', html_content='<p>Sensor A未检测到目标文件,下游任务B、C已被跳过。</p>' ) # 定义传感器任务 sensor_a = FileSensor( task_id='sensor_a', filepath='/path/to/target/file', soft_fail=True, poke_function=sensor_poke_with_xcom, dag=dag ) # 定义并行告警任务 alert_task = PythonOperator( task_id='alert_on_sensor_failure', python_callable=check_sensor_and_send_alert, provide_context=True, trigger_rule='all_done', dag=dag ) # 设置任务依赖 sensor_a >> task_b >> task_c sensor_a >> alert_task
方法三:成功回调中判断实际检测结果
利用on_success_callback,因为soft_fail=True时任务会标记为成功,在成功回调里读取XCom中存储的检测结果,若为失败则触发告警。
from airflow.sensors.filesystem import FileSensor from airflow.utils.email import send_email def sensor_success_callback(context): # 读取XCom中的检测结果 detection_result = context['ti'].xcom_pull( key='sensor_a_detection_result', task_ids='sensor_a' ) if not detection_result: send_email( to='your-alert-email@example.com', subject='Sensor A 软失败:目标未找到', html_content='<p>Sensor A未检测到目标文件,下游任务B、C已被跳过。</p>' ) def sensor_poke_with_xcom(context): target_exists = super(FileSensor, self).poke(context) context['ti'].xcom_push(key='sensor_a_detection_result', value=target_exists) return target_exists sensor_a = FileSensor( task_id='sensor_a', filepath='/path/to/target/file', soft_fail=True, poke_function=sensor_poke_with_xcom, on_success_callback=sensor_success_callback, dag=dag ) sensor_a >> task_b >> task_c
内容的提问来源于stack exchange,提问作者user3735871
相关产品推荐
相关产品推荐

