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

如何终止Airflow中已删除DAG仍运行的幽灵任务?

问题

我正在学习Airflow的工作机制,遇到了一个问题:我创建了一个DAG用于检查Flask API是否正常运行,若异常则发送邮件,代码如下:

def check_api_prod(id):
    data = ""
   
    url_siv = "http://111.11.11.111:5315/api/v1/search/"+str(id)
    
    body = {
        
    }
    res = requests.get(url_siv, json=body).json()
    
    
    
    if str(res['TAG']['id'])==str(id):
        return True
    else:
        raise AirflowException('l API est injoignable')
        

########################################################
#######################################################

def send_mail(receiver,object,Message ):

    msg = MIMEText(Message)

    msg["Subject"] = str(object)

    msg['From'] = "mp.name@gmail.com"
    

    with smtplib.SMTP('mail.gmx.com', 25) as smtp:
        smtp.ehlo()
        smtp.starttls() 
        smtp.ehlo()

        smtp.login("mp.name@gmail.com", "PASS_wrds?")


        try:

            smtp.sendmail("mp.name@gmail.com",  receiver, msg.as_string())
            print("mail sended")
        except Exception as er:
            print("error", er)



def error_api_email_function(self):
    
    try:
        
    
        receiver=["ch.mp@gmail.com", "fm@dgmail.fr"]
        
                 
        object="nouvelle version du message"
        Message_Echech=f" NOUVELLE VERSION du message "     
        now = datetime.now()   
        dt_string = now.strftime("%d/%m/%Y %H:%M:%S")
        Message_Echech=Message_Echech+ " "+dt_string
        envoi_mail(receiver,object,Message_Echech)
        print("mail envoyé")
     
       
    except Exception as er:
        print("impossible d'envoyer le mail à cause", er)



#################################################################
with DAG('Check_api', description='api 2', start_date=dt.datetime(year=2023, month=5, day=1), catchup=False,schedule_interval=dt.timedelta(minutes=5)) as dag:
    verifier_api=PythonOperator(
        task_id='check_api',
        python_callable=check_api_prod ,
        op_args=[ "GGMT34O"],
        on_failure_callback= error_api_email_function
       
    )

代码运行正常,但当我想要停止该DAG时,已暂停DAG但任务仍在运行,且持续收到邮件。我更换了Airflow Executor,甚至停止并清理了Docker容器及镜像(执行命令docker compose up --all),删除DAG文件后,重启Docker时任务仍会重新启动!请问有解决办法吗?

解决方法
  • 清理Airflow元数据库:Airflow的任务调度、运行状态全部存在元数据库中,仅删除DAG文件或重启容器无法清除历史任务记录,需要手动清理元数据:

    1. 进入数据库容器:docker exec -it <airflow-db-container-name> bash(替换<airflow-db-container-name>为实际容器名,可通过docker ps查看)
    2. 登录数据库:以PostgreSQL为例,执行psql -U airflow
    3. 执行SQL删除目标DAG的所有关联记录:
      DELETE FROM dag WHERE dag_id = 'Check_api';
      DELETE FROM dag_run WHERE dag_id = 'Check_api';
      DELETE FROM task_instance WHERE dag_id = 'Check_api';
      DELETE FROM job WHERE dag_id = 'Check_api';
      

    执行完成后退出数据库和容器。

  • 彻底清理Docker持久化资源:之前使用的docker compose up --all命令无效,正确的彻底清理命令为:

    1. 停止并删除容器、网络、关联卷:docker compose down -v
    2. 删除未使用的镜像:docker image prune -a
      注意:-v参数会删除Airflow的元数据卷,所有历史调度记录会被清除,操作前确认无需备份。
  • 关闭任务重试机制:若任务因API异常触发失败回调,默认重试机制(若配置过)会导致任务反复执行、持续发邮件。可在PythonOperator中明确关闭重试:

    verifier_api=PythonOperator(
        task_id='check_api',
        python_callable=check_api_prod ,
        op_args=[ "GGMT34O"],
        on_failure_callback= error_api_email_function,
        retries=0
    )
    
  • 手动终止运行中任务:

    • 通过Airflow UI:进入目标DAG的任务实例页面,找到正在运行的任务,点击「Kill」按钮强制终止。
    • 通过命令行:执行airflow tasks kill Check_api check_api <execution-date>,替换<execution-date>为任务的执行日期(格式如2023-05-01T00:00:00)。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.22 10:25:45