如何终止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文件或重启容器无法清除历史任务记录,需要手动清理元数据:
- 进入数据库容器:
docker exec -it <airflow-db-container-name> bash(替换<airflow-db-container-name>为实际容器名,可通过docker ps查看) - 登录数据库:以PostgreSQL为例,执行
psql -U airflow - 执行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命令无效,正确的彻底清理命令为:- 停止并删除容器、网络、关联卷:
docker compose down -v - 删除未使用的镜像:
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
相关产品推荐
相关产品推荐

