Airflow仪表盘清除运行实例时如何清理历史数据?
Airflow clear操作自动执行清理脚本的实现方案
完全可以实现你的需求,以下是两种生产环境常用的实现方案:
方案1:使用Task的pre_execute回调(通用性最高,无需额外配置)
所有Airflow Operator都支持pre_execute回调参数,该回调会在任务正式执行前自动触发。你可以在回调函数中编写对应日期的历史数据清理逻辑,不管是clear重跑、手动触发还是catch-up调度,都会先执行清理再跑新的INSERT逻辑,完全不需要暂停Airflow进程。
代码示例如下:
from airflow.models import TaskInstance from airflow.providers.postgres.hooks.postgres import PostgresHook # 定义清理回调函数 def clear_history_data(context: dict): # 从上下文中获取当前任务的执行日期 exec_date = context["execution_date"].strftime("%Y-%m-%d") # 编写和INSERT逻辑匹配的DELETE语句,按执行日期过滤数据 delete_sql = f"DELETE FROM your_target_table WHERE dt = '{exec_date}'" # 用对应数据库钩子执行SQL,这里以Postgres为例,可替换为你用的数据库类型 db_hook = PostgresHook(postgres_conn_id="your_db_connection_id") db_hook.run(delete_sql) # 在你的INSERT任务中绑定pre_execute回调 insert_task = PythonOperator( task_id="your_insert_task_id", python_callable=your_actual_insert_function, pre_execute=clear_history_data, dag=dag )
方案2:监听TaskInstanceCleared事件(仅clear操作触发)
如果你只需要在UI点击clear操作时才触发清理,正常每日调度不需要清理,可以使用Airflow的事件监听器机制,仅监听任务清除事件触发清理逻辑。
代码示例如下:
from airflow.listeners import hookimpl from airflow.models import TaskInstance from airflow.providers.postgres.hooks.postgres import PostgresHook @hookimpl def on_task_instance_cleared(task_instance: TaskInstance, **kwargs): # 过滤你需要处理的DAG和任务,避免全局所有DAG的clear操作都触发清理 if task_instance.dag_id == "your_dag_id" and task_instance.task_id == "your_insert_task_id": exec_date = task_instance.execution_date.strftime("%Y-%m-%d") delete_sql = f"DELETE FROM your_target_table WHERE dt = '{exec_date}'" db_hook = PostgresHook(postgres_conn_id="your_db_connection_id") db_hook.run(delete_sql)
Airflow 2.0及以上版本默认支持监听器功能,无需额外开启配置。
注意事项
- 两种方案都需要确保DELETE语句的过滤条件和你的INSERT写入逻辑完全对齐,避免误删其他日期的业务数据
- 如果你的业务逻辑允许,更推荐使用方案1,所有重跑场景都会自动先清理历史数据,天然保证任务幂等
- 正式重跑生产历史数据前,建议先在测试环境验证清理和写入逻辑的正确性
内容的提问来源于stack exchange,提问作者Nicholas Humphrey
相关产品推荐
相关产品推荐

