Airflow设置任务为Failed状态后无法终止Redshift查询问题
解决Airflow任务终止后Redshift查询未停止的问题
核心问题分析
- Airflow的
ti.set_state()仅更新元数据库中的任务状态,不会主动终止任务的运行进程,因此任务内的Redshift查询会持续执行。 - Airflow默认在发送SIGTERM信号后短时间(通常5秒)内会发送SIGKILL强制杀死进程,导致
cursor.close()这类代码来不及执行——且cursor.close()仅关闭客户端连接,不会主动取消Redshift端的查询。 - 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
相关产品推荐
相关产品推荐

