Merge pull request #3016 from InfinityPacer/feature/scheduler

This commit is contained in:
jxxghp
2024-11-06 20:01:14 +08:00
committed by GitHub
+16 -12
View File
@@ -1,4 +1,3 @@
import logging
import threading import threading
import traceback import traceback
from datetime import datetime, timedelta from datetime import datetime, timedelta
@@ -27,12 +26,6 @@ from app.schemas.types import EventType
from app.utils.singleton import Singleton from app.utils.singleton import Singleton
from app.utils.timer import TimerUtils from app.utils.timer import TimerUtils
# 获取 apscheduler 的日志记录器
scheduler_logger = logging.getLogger('apscheduler')
# 设置日志级别为 WARNING
scheduler_logger.setLevel(logging.WARNING)
class SchedulerChain(ChainBase): class SchedulerChain(ChainBase):
pass pass
@@ -436,7 +429,6 @@ class Scheduler(metaclass=Singleton):
try: try:
sid = f"{service['id']}" sid = f"{service['id']}"
job_id = sid.split("|")[0] job_id = sid.split("|")[0]
if job_id not in self._jobs:
self._jobs[job_id] = { self._jobs[job_id] = {
"func": service["func"], "func": service["func"],
"name": service["name"], "name": service["name"],
@@ -450,7 +442,8 @@ class Scheduler(metaclass=Singleton):
id=sid, id=sid,
name=service["name"], name=service["name"],
**service["kwargs"], **service["kwargs"],
kwargs={"job_id": job_id} kwargs={"job_id": job_id},
replace_existing=True
) )
logger.info(f"注册插件{plugin_name}服务:{service['name']} - {service['trigger']}") logger.info(f"注册插件{plugin_name}服务:{service['name']} - {service['trigger']}")
except Exception as e: except Exception as e:
@@ -468,14 +461,25 @@ class Scheduler(metaclass=Singleton):
with self._lock: with self._lock:
# 获取插件名称 # 获取插件名称
plugin_name = PluginManager().get_plugin_attr(pid, "plugin_name") plugin_name = PluginManager().get_plugin_attr(pid, "plugin_name")
for job_id, service in self._jobs.copy().items(): # 先从 _jobs 中查找匹配的服务
jobs_to_remove = [(job_id, service) for job_id, service in self._jobs.items() if service.get("pid") == pid]
if not jobs_to_remove:
return
for job_id, service in jobs_to_remove:
try: try:
if service.get("pid") == pid: # 移除服务
self._jobs.pop(job_id, None) self._jobs.pop(job_id, None)
# 在调度器中查找并移除对应的 job
job_removed = False
for job in list(self._scheduler.get_jobs()):
job_id_from_service = job.id.split("|")[0]
if job_id == job_id_from_service:
try: try:
self._scheduler.remove_job(job_id) self._scheduler.remove_job(job.id)
job_removed = True
except JobLookupError: except JobLookupError:
pass pass
if job_removed:
logger.info(f"移除插件服务({plugin_name}){service.get('name')}") logger.info(f"移除插件服务({plugin_name}){service.get('name')}")
except Exception as e: except Exception as e:
logger.error(f"移除插件服务失败:{str(e)} - {job_id}: {service}") logger.error(f"移除插件服务失败:{str(e)} - {job_id}: {service}")