You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.20 07:51:34