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

如何在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.28 21:15:05