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

如何基于指定条件将PostgresOperator类型的Airflow DAG任务标记为失败

实现方案说明

首先回答你的疑问:不需要额外做DAG执行结果的回传动作,有两种成熟的实现方案可以直接满足需求:

方案1:直接修改SQL语句,无需调整任务类型

这是成本最低的方案,直接调整你要执行的SQL,统计结果为0时主动抛出数据库异常,PostgresOperator检测到SQL执行报错会自动标记任务失败。
修改后的SQL如下:

DO $$
DECLARE
    insert_count integer;
BEGIN
    SELECT count(*) INTO insert_count FROM my_table WHERE date_insert = '{{ ds }}';
    IF insert_count = 0 THEN
        RAISE EXCEPTION 'Date % has no inserted data in my_table', '{{ ds }}';
    END IF;
END $$;

适用场景

仅需要做0值校验、没有后续逻辑依赖统计结果的场景,无需修改任何Python代码即可实现需求。

方案2:使用PostgresHook结合PythonOperator实现

如果后续逻辑需要用到统计的数值,你可以替换原有的PostgresOperator为PythonOperator,通过PostgresHook获取查询结果后做判断,需要的代码如下:

  1. 首先导入依赖包
from airflow.providers.postgres.hooks.postgres import PostgresHook
from airflow.operators.python import PythonOperator
  1. 定义校验逻辑函数
def check_my_table_count(**context):
    # 替换为你自己的Postgres连接ID
    pg_hook = PostgresHook(postgres_conn_id="your_postgres_conn_id")
    exec_date = context["ds"]
    # 获取统计结果
    query_result = pg_hook.get_first(f"SELECT count(*) FROM my_table WHERE date_insert = '{exec_date}'")
    count = query_result[0]
    if count == 0:
        # 抛出异常即可标记任务失败
        raise ValueError(f"No data found in my_table for date {exec_date}, count is {count}")
    # 可选:把统计结果推送到XCom给下游任务使用
    context["ti"].xcom_push(key="my_table_daily_count", value=count)
  1. 在DAG中定义任务
check_data_task = PythonOperator(
    task_id="check_my_table_daily_data",
    python_callable=check_my_table_count,
    provide_context=True,
    dag=dag # 替换为你自己的DAG实例名
)

适用场景

需要用到统计结果做后续自定义逻辑(比如数值阈值告警、下游分支判断等)的场景,灵活度更高。

内容的提问来源于stack exchange,提问作者Alireza

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.24 05:06:08