如何基于指定条件将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获取查询结果后做判断,需要的代码如下:
- 首先导入依赖包
from airflow.providers.postgres.hooks.postgres import PostgresHook from airflow.operators.python import PythonOperator
- 定义校验逻辑函数
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)
- 在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
相关产品推荐
相关产品推荐

