From 73e25050df422a8e8986f7543a8c072e2a82c707 Mon Sep 17 00:00:00 2001 From: "34047007@qq.com" <34047007@qq.com> Date: Fri, 7 Aug 2026 21:47:38 +0800 Subject: [PATCH] =?UTF-8?q?checkpoint:=20=E8=B0=83=E5=BA=A6=E4=BB=BB?= =?UTF-8?q?=E5=8A=A1=E5=AD=98=E5=82=A8=E8=90=BD=20PG=E2=80=94=E2=80=94jobs?= =?UTF-8?q?tore=20=E5=BD=92=E4=B8=80=20default=20+=20execute=20=E5=86=99?= =?UTF-8?q?=E5=9B=9E=E7=9C=9F=E5=AE=9E=E8=A7=A6=E5=8F=91=20+=20=E5=89=8D?= =?UTF-8?q?=E7=AB=AF=20store=20=E5=AD=97=E5=85=B8=E6=9B=B4=E6=96=B0?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - ap_scheduler: 默认分组改 SQLAlchemyJobStore on PG,删 RedisJobStore 与 redis_instance;注释/告警口径同步 - node/service: jobstore 归一 default(兼容存量 sqlalchemy 名);cron/interval/date 写回 task_node.trigger/trigger_args(前端列表可见真实触发);run_job_now 用独立临时 id,不再覆盖节点既有排程 - enums: 删 APSCHEDULER_LOCK_KEY(jobstore 落 PG 后死锁) - sys_dict_data: sys_job_store 字典归一 default,删 sqlalchemy 档 - 前端: 任务表单 store 选项与 default 归一 --- .../v1/module_task/cronjob/node/service.py | 40 ++++++++++++++----- backend/app/common/enums.py | 1 - backend/app/core/ap_scheduler.py | 21 +++------- backend/app/init_app.py | 6 +-- backend/sql/data/sys_dict_data.json | 16 +------- .../views/module_task/cronjob/node/index.vue | 23 ++++++++++- 6 files changed, 61 insertions(+), 46 deletions(-) diff --git a/backend/app/api/v1/module_task/cronjob/node/service.py b/backend/app/api/v1/module_task/cronjob/node/service.py index 07720be..b7c2865 100644 --- a/backend/app/api/v1/module_task/cronjob/node/service.py +++ b/backend/app/api/v1/module_task/cronjob/node/service.py @@ -156,6 +156,12 @@ class NodeService: else: raise CustomException(msg=f"不支持的触发方式: {trigger}") + # 非一次性触发写回 task_node,任务列表可展示真实触发配置; + # apscheduler_jobs 表已有完整状态,此处仅补齐业务表字段 + if trigger != "now": + obj.trigger = trigger + obj.trigger_args = trigger_args + return {"job_id": id, "status": "executed", "trigger": trigger} async def batch_set_status(self, ids: list[int], status: int) -> None: @@ -171,14 +177,21 @@ class NodeService: # ── NodeModel 封装的任务添加方法 ──────────────────────────── -def _add_job_with_trigger(job_info: NodeModel, trigger) -> Job: - """将 NodeModel 封装的任务添加到 APScheduler 调度器。""" +def _add_job_with_trigger(job_info: NodeModel, trigger, job_id: str | None = None) -> Job: + """将 NodeModel 封装的任务添加到 APScheduler 调度器。 + + job_id 为 None 时用节点 id 作调度任务 id(cron/interval/date 复用节点 id,便于原地更新); + 立即执行(now)传独立临时 id,避免覆盖节点已有的定时排程。 + """ code_block = job_info.func if not code_block or not code_block.strip(): raise ValueError("任务代码块不能为空") - jobstore = job_info.jobstore or "sqlalchemy" + jobstore = job_info.jobstore or "default" + if jobstore == "sqlalchemy": + jobstore = "default" # 存量节点 jobstore="sqlalchemy"(旧持久化 store 名)归一为统一 PG store executor = job_info.executor or "threadpool" + job_id = job_id or str(job_info.id) job_args = [] if job_info.args: @@ -195,7 +208,7 @@ def _add_job_with_trigger(job_info: NodeModel, trigger) -> Job: except json.JSONDecodeError: raise ValueError(f"关键字参数JSON格式无效: {kwargs_str}") - SchedulerUtil.job_name_cache[str(job_info.id)] = job_info.name or "" + SchedulerUtil.job_name_cache[job_id] = job_info.name or "" try: job = scheduler.add_job( @@ -203,37 +216,42 @@ def _add_job_with_trigger(job_info: NodeModel, trigger) -> Job: trigger=trigger, args=[str(job_info.id), code_block, *job_args], kwargs=job_kwargs, - id=str(job_info.id), + id=job_id, name=job_info.name, coalesce=job_info.coalesce, max_instances=1, jobstore=jobstore, executor=executor, ) - logger.info(f"任务 {job_info.id} 添加到 {jobstore} 存储器成功") + logger.info(f"任务 {job_id} 添加到 {jobstore} 存储器成功") return job except ConflictingIdError: - scheduler.remove_job(job_id=str(job_info.id), jobstore=jobstore) + scheduler.remove_job(job_id=job_id, jobstore=jobstore) job = scheduler.add_job( func=SchedulerUtil._task_wrapper, trigger=trigger, args=[str(job_info.id), code_block, *job_args], kwargs=job_kwargs, - id=str(job_info.id), + id=job_id, name=job_info.name, coalesce=job_info.coalesce, max_instances=1, jobstore=jobstore, executor=executor, ) - logger.info(f"任务 {job_info.id} 已存在,已移除旧任务并重新添加") + logger.info(f"任务 {job_id} 已存在,已移除旧任务并重新添加") return job def add_and_run_job_now(job_info: NodeModel) -> Job: - """立即执行任务(加入调度器并尽快触发一次)。""" + """立即执行任务(加入调度器并尽快触发一次)。 + + 用独立临时 id({节点id}_run_now_{时间戳})添加一次性 date 任务: + 只触发一次、跑完自动移除,不覆盖该节点已有的 cron/interval 定时排程。 + """ + temp_id = f"{job_info.id}_run_now_{datetime.now().timestamp()}" trigger = DateTrigger(run_date=datetime.now() + timedelta(seconds=0.1)) - return _add_job_with_trigger(job_info, trigger) + return _add_job_with_trigger(job_info, trigger, job_id=temp_id) def add_cron_job( diff --git a/backend/app/common/enums.py b/backend/app/common/enums.py index 8189a38..3b241c4 100644 --- a/backend/app/common/enums.py +++ b/backend/app/common/enums.py @@ -48,7 +48,6 @@ class RedisInitKeyConfig(Enum): EMAIL_CODES = {"key": "email_codes", "remark": "邮箱验证码"} SYSTEM_CONFIG = {"key": "system_config", "remark": "系统配置"} SYSTEM_DICT = {"key": "system_dict", "remark": "数据字典"} - APSCHEDULER_LOCK_KEY = {"key": "scheduler_job_lock", "remark": "定时任务初始化锁"} AI_MODEL_CONFIG = {"key": "ai_model_config", "remark": "用户AI模型配置"} @property diff --git a/backend/app/core/ap_scheduler.py b/backend/app/core/ap_scheduler.py index 1e5108b..76b5dd9 100644 --- a/backend/app/core/ap_scheduler.py +++ b/backend/app/core/ap_scheduler.py @@ -18,13 +18,11 @@ from apscheduler.executors.asyncio import AsyncIOExecutor from apscheduler.executors.pool import ProcessPoolExecutor, ThreadPoolExecutor from apscheduler.job import Job from apscheduler.jobstores.memory import MemoryJobStore -from apscheduler.jobstores.redis import RedisJobStore from apscheduler.jobstores.sqlalchemy import SQLAlchemyJobStore from apscheduler.schedulers.asyncio import AsyncIOScheduler from apscheduler.triggers.cron import CronTrigger from apscheduler.triggers.date import DateTrigger from apscheduler.triggers.interval import IntervalTrigger -from redis.asyncio import Redis from sqlalchemy.orm import Session from app.config.setting import settings @@ -37,15 +35,9 @@ JOB_STATUS_FAILED = 3 scheduler = AsyncIOScheduler() scheduler.configure( jobstores={ - "default": RedisJobStore( - host=settings.REDIS_HOST, - port=int(settings.REDIS_PORT), - username=settings.REDIS_USER or None, - password=settings.REDIS_PASSWORD or None, - db=int(settings.REDIS_DB_NAME), - protocol=2, - ), - "sqlalchemy": SQLAlchemyJobStore(url=settings.DB_URI, engine=engine), + # 默认分组 = PostgreSQL(apscheduler_jobs 表):任务持久化,随 pg_data 每日备份; + # 与业务主库同生共死,不再依赖可降级/无备份的 Redis(多 Worker 场景本系统不适用,见运营手册) + "default": SQLAlchemyJobStore(url=settings.DB_URI, engine=engine), "memory": MemoryJobStore(), }, executors={ @@ -64,7 +56,6 @@ scheduler.configure( class SchedulerUtil: """定时任务 SDK — 仅封装 APScheduler 核心操作,不含业务逻辑(无 ORM/实体引用)。""" - redis_instance: Redis | None = None job_name_cache: dict[str, str | tuple[str, str]] = {} @classmethod @@ -129,15 +120,13 @@ class SchedulerUtil: logger.info("所有任务已从调度器中移除") @classmethod - async def init_scheduler(cls, redis: Redis | None = None) -> None: + async def init_scheduler(cls) -> None: """应用启动时初始化定时任务调度器(含系统级周期任务注册)。 返回: - None """ try: - if redis: - cls.redis_instance = redis scheduler.start() scheduler.add_listener(cls._dispatch_job_event, EVENT_ALL) scheduler.resume() @@ -177,7 +166,7 @@ class SchedulerUtil: @classmethod def shutdown(cls, wait: bool = False) -> None: - # Redis 不可达时 start()/add_job 可能失败导致调度器未真正运行: + # jobstore 已落 PG:数据库不可达时 start()/add_job 可能失败导致调度器未真正运行: # 对未运行实例调用 shutdown 会抛 SchedulerNotRunningError,须先判状态 if not scheduler.running: return diff --git a/backend/app/init_app.py b/backend/app/init_app.py index aac9ee7..abf0b38 100644 --- a/backend/app/init_app.py +++ b/backend/app/init_app.py @@ -40,7 +40,7 @@ async def lifespan(app: FastAPI) -> AsyncGenerator[Any, Any]: else: logger.warning( "⚠️ Redis 不可用,以降级模式启动:参数/字典读接口回源数据库;" - "登录会话/在线监控/AI配置/调度任务接口将返回 503" + "登录会话/在线监控/AI配置接口将返回 503(调度任务 jobstore 已落 PG,不受影响)" ) from app.utils.dict_util import preload_value2label_cache @@ -50,11 +50,11 @@ async def lifespan(app: FastAPI) -> AsyncGenerator[Any, Any]: scheduler_ready = False try: - await SchedulerUtil.init_scheduler(redis=app.state.redis) + await SchedulerUtil.init_scheduler() scheduler_ready = SchedulerUtil.is_running() logger.info("✅ 定时任务调度器初始化完成") except Exception as e: - # Redis 不可达时默认 jobstore 连不上:不阻断启动,任务接口降级为 503 + # 调度器初始化失败(jobstore 已落 PG,与主库同生共死):不阻断启动 SchedulerUtil.shutdown(wait=False) logger.error("❌ 定时任务调度器初始化失败(继续启动,任务接口将返回 503): {}", e) diff --git a/backend/sql/data/sys_dict_data.json b/backend/sql/data/sys_dict_data.json index 45a5eb1..d925663 100644 --- a/backend/sql/data/sys_dict_data.json +++ b/backend/sql/data/sys_dict_data.json @@ -229,7 +229,7 @@ }, { "dict_sort": 1, - "dict_label": "默认(Redis)", + "dict_label": "默认(数据库)", "dict_value": "default", "dict_type": "sys_job_store", "dict_type_id": 6, @@ -237,22 +237,10 @@ "list_class": null, "is_default": true, "status": 0, - "description": "默认分组,支持多Worker部署" + "description": "默认分组(PostgreSQL 持久化,随主库每日备份)" }, { "dict_sort": 2, - "dict_label": "数据库(Sqlalchemy)", - "dict_value": "sqlalchemy", - "dict_type": "sys_job_store", - "dict_type_id": 6, - "css_class": "", - "list_class": null, - "is_default": false, - "status": 0, - "description": "数据库分组,持久化存储" - }, - { - "dict_sort": 3, "dict_label": "内存(Memory)", "dict_value": "memory", "dict_type": "sys_job_store", diff --git a/frontend/web/src/views/module_task/cronjob/node/index.vue b/frontend/web/src/views/module_task/cronjob/node/index.vue index e2e968a..4fd2f20 100644 --- a/frontend/web/src/views/module_task/cronjob/node/index.vue +++ b/frontend/web/src/views/module_task/cronjob/node/index.vue @@ -502,6 +502,27 @@ const { label: "执行器", minWidth: 80, }, + { + prop: "trigger", + label: "触发方式", + minWidth: 100, + formatter: (row: NodeTable) => { + const map: Record = { + now: "立即执行", + cron: "Cron", + interval: "间隔", + date: "定时", + }; + return row.trigger ? (map[row.trigger] ?? row.trigger) : "-"; + }, + }, + { + prop: "trigger_args", + label: "触发配置", + minWidth: 160, + showOverflowTooltip: true, + formatter: (row: NodeTable) => row.trigger_args || "-", + }, { prop: "created_time", label: "创建时间", @@ -794,7 +815,7 @@ const initialFormData: Partial = { id: undefined, name: "", code: undefined, - jobstore: "sqlalchemy", + jobstore: "default", executor: "default", func: defaultCodeBlock, args: undefined,