如何在Airflow任务内部感知终止状态并杀死运行中的Neo4j查询
手动标记Airflow任务状态时同步终止Neo4j查询的实现方案
核心思路
- Airflow在GUI手动标记任务为Failed/Completed时,会向运行中的任务进程发送
SIGTERM信号,捕获该信号即可感知任务终止事件 - 每次提交Neo4j查询时先获取当前事务的唯一ID,收到终止信号后调用Neo4j内置的终止事务语句杀掉对应查询即可
可直接复用的代码实现
import signal from neo4j import GraphDatabase from airflow.exceptions import AirflowTaskTerminated # 全局存储当前运行的Neo4j事务ID和会话实例 running_tx_id = None current_session = None # 全局复用Neo4j Driver,不要每次查询重复创建 neo4j_driver = GraphDatabase.driver("bolt://你的Neo4j地址:7687", auth=("账号", "密码")) # SIGTERM信号处理函数:终止当前运行的Neo4j查询后退出任务 def terminate_handler(signum, frame): global running_tx_id, current_session if running_tx_id and current_session: # 执行Neo4j事务终止语句 current_session.run("TERMINATE TRANSACTIONS $tx_id", tx_id=running_tx_id) current_session.close() raise AirflowTaskTerminated("任务被手动终止,关联Neo4j查询已关闭") # 注册SIGTERM信号监听 signal.signal(signal.SIGTERM, terminate_handler) # 业务查询执行函数 def exec_neo4j_query(query: str): global running_tx_id, current_session current_session = neo4j_driver.session() # 获取当前会话的事务ID(适配Neo4j 4.0+版本) tx_id_res = current_session.run("CALL tx.getMetaData() YIELD transactionId RETURN transactionId as tx_id") running_tx_id = tx_id_res.single()["tx_id"] # 执行实际业务查询 query_res = current_session.run(query) # 这里添加你的查询结果处理逻辑 # ... # 执行完成后清空状态 current_session.close() running_tx_id = None
注意事项
- Neo4j 4.0以下版本无
tx.getMetaData()内置过程,可以通过dbms.listTransactions()筛选当前连接用户、客户端地址对应的事务ID来使用 - Airflow 2.x版本发送
SIGTERM后,默认会等待30秒再发送SIGKILL强制杀掉进程,可调整任务的killed_task_cleanup_time参数预留足够的查询终止执行时间 - 若你当前使用Airflow官方
Neo4jOperator实现查询,只需继承该Operator重写execute方法,将上述信号处理和事务ID获取逻辑嵌入即可,无需大幅修改现有任务代码
内容的提问来源于stack exchange,提问作者M Hrytsenia
相关产品推荐
相关产品推荐

