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

Airflow设置任务为Failed状态后无法终止Redshift查询问题

解决Airflow任务终止后Redshift查询未停止的问题

核心问题分析

  1. Airflow的ti.set_state()仅更新元数据库中的任务状态,不会主动终止任务的运行进程,因此任务内的Redshift查询会持续执行。
  2. Airflow默认在发送SIGTERM信号后短时间(通常5秒)内会发送SIGKILL强制杀死进程,导致cursor.close()这类代码来不及执行——且cursor.close()仅关闭客户端连接,不会主动取消Redshift端的查询。
  3. Redshift查询一旦启动,除非被主动取消,否则会运行至完成,对应的锁也不会释放。

可行解决方案

方案1:主动取消Redshift查询(推荐)

先通过Redshift系统函数终止目标查询,再更新Airflow任务状态:

  • 步骤1:在任务A执行时,将Redshift查询的PID存入XCom(任务上下文的跨任务通信机制):
    def run_redshift_query(**context):
        conn = psycopg2.connect("你的Redshift连接字符串")
        cur = conn.cursor()
        # 获取当前会话的PID
        cur.execute("SELECT pg_backend_pid();")
        query_pid = cur.fetchone()[0]
        context['ti'].xcom_push(key='redshift_query_pid', value=query_pid)
        # 执行大型查询
        cur.execute("你的Redshift查询语句")
        conn.commit()
        cur.close()
        conn.close()
    
  • 步骤2:检测到DAG B运行时,从XCom读取PID,调用Redshift函数取消查询:
    -- 优先尝试优雅取消
    SELECT pg_cancel_backend(<从XCom获取的PID>);
    -- 若优雅取消无效,强制终止会话
    SELECT pg_terminate_backend(<从XCom获取的PID>);
    
  • 步骤3:执行ti.set_state(State.FAILED)和dag.set_state(DagRunState.FAILED)更新Airflow状态。

方案2:延长Airflow任务终止等待时间

修改任务的超时配置,给优雅关闭代码足够执行时间:

  • 在任务A定义中延长执行超时:
    task_a = PythonOperator(
        task_id='task_a',
        python_callable=run_redshift_query,
        execution_timeout=timedelta(minutes=2),  # 延长超时,给取消查询留时间
        ...
    )
    
  • 同时在任务A中捕获SIGTERM信号,主动取消Redshift查询:
    import signal
    import psycopg2
    
    def handle_sigterm(signum, frame):
        # 从上下文或本地变量获取查询PID
        query_pid = ...
        conn = psycopg2.connect("你的Redshift连接字符串")
        cur = conn.cursor()
        cur.execute(f"SELECT pg_cancel_backend({query_pid});")
        conn.commit()
        cur.close()
        conn.close()
        exit(1)
    
    def run_redshift_query(**context):
        signal.signal(signal.SIGTERM, handle_sigterm)
        # 后续查询执行逻辑同方案1
    

方案3:使用Airflow原生的TaskInstance.kill()方法

不要直接调用ti.set_state(),而是用kill()触发任务进程的终止信号,配合信号处理器完成优雅关闭:

from airflow.models import DagRun, TaskInstance
from airflow.utils.state import State
import time

# 检测运行中的DAG B
dag_run = DagRun.find(dag_id='dag_b', state=State.RUNNING)[0]
# 获取任务A的实例
ti = TaskInstance(task=dag_run.get_task_instance('task_a'), run_id=dag_run.run_id)
# 发送SIGTERM信号终止任务
ti.kill()
# 等待一段时间让取消逻辑执行
time.sleep(30)
# 更新任务和DAG状态
ti.set_state(State.FAILED)
dag_run.set_state(DagRunState.FAILED)

关键注意事项

  • 必须先取消Redshift查询,再更新Airflow状态,反向操作无法终止正在运行的查询。
  • pg_cancel_backend会尝试优雅终止查询,pg_terminate_backend直接杀死会话,可根据查询状态选择使用。
  • 执行取消操作的账号需要Redshift的SUPERUSER或pg_signal_backend权限。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.22 05:00:35