Files
2026-07-13 12:34:57 +08:00

628 lines
21 KiB
Python
Raw Permalink Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
"""
任务调度器模块
提供定时任务调度功能,基于 APScheduler 实现
"""
import uuid
from datetime import datetime
from typing import Optional, Callable, Any, TYPE_CHECKING, Union
from loguru import logger
from dynaconf import Dynaconf
from apscheduler.schedulers.background import BackgroundScheduler
from apscheduler.triggers.cron import CronTrigger
from apscheduler.triggers.interval import IntervalTrigger
from apscheduler.triggers.date import DateTrigger
from apscheduler.executors.pool import ThreadPoolExecutor
from apscheduler.jobstores.memory import MemoryJobStore
from .config import get_config, get_settings
from ..exceptions import SchedulerError
if TYPE_CHECKING:
from ..jobs.scheduled_job import ScheduledJob
class Scheduler:
"""应用任务调度器(基于 APScheduler"""
def __init__(self, config: Optional[Union[str, Dynaconf]] = None, enabled: Optional[bool] = None):
"""
初始化调度器
Args:
config: 配置文件路径或配置对象(Dynaconf),如果为 None 则使用默认配置
enabled: 是否启用调度器,如果为 None 则从配置文件读取
多 workers 模式下,默认只在 primary worker 启用
"""
self._logger = logger.bind(name="scheduler")
self._scheduled_jobs = {} # 存储 ScheduledJob 对象
# 获取配置对象
if isinstance(config, Dynaconf):
# 如果传入的是配置对象,直接使用
self._config = config
else:
# 如果传入的是配置文件路径或 None,则获取配置对象
self._config = get_settings(config)
# 从配置对象读取调度器配置,enabled 参数可覆盖配置
config_enabled = self._config.get('scheduler.enabled', True)
self._enabled = enabled if enabled is not None else config_enabled
timezone_str = self._config.get('scheduler.timezone', 'UTC')
self._timezone = self._parse_timezone(timezone_str)
max_workers = self._config.get('scheduler.max_workers', 10)
# 配置 APScheduler
jobstores = {
'default': MemoryJobStore()
}
executors = {
'default': ThreadPoolExecutor(max_workers=max_workers)
}
job_defaults = {
'coalesce': False, # 不合并错过的任务
'max_instances': 3, # 最大并发实例数
'misfire_grace_time': 30 # 错过执行的宽限时间(秒)
}
# 创建 APScheduler 实例
self._scheduler = BackgroundScheduler(
jobstores=jobstores,
executors=executors,
job_defaults=job_defaults,
timezone=self._timezone
)
self._logger.debug(f"调度器初始化完成 - enabled: {self._enabled}, timezone: {timezone_str}, max_workers: {max_workers}")
def _parse_timezone(self, timezone_str: str):
"""
解析时区字符串
APScheduler 会自动使用此时区处理所有任务,无需手动转换
"""
try:
import pytz
return pytz.timezone(timezone_str)
except ImportError:
self._logger.warning("未安装 pytz,使用系统时区")
return None
except Exception as e:
self._logger.warning(f"解析时区失败: {e},使用系统时区")
return None
def _resolve_job_id(self, job_id: Optional[str], prefix: str, func: Callable) -> str:
"""解析并校验任务 ID
默认 ID 为 ``{prefix}_{模块名}.{限定名}``(限定名含类名),
不同类中的同名方法不会冲突(issue #14)。
注册前显式查重——APScheduler 未启动时不校验 pending 队列中的
ID 冲突,问题会延迟到 ``start()`` 落库时才爆发,这里提前拦截:
- 显式传入的 ID 已存在:抛 ``SchedulerError``
- 自动生成的 ID 已存在(同一函数注册多次):追加短 uid 消歧并告警
"""
existing_ids = set(self.list_jobs())
if job_id is not None:
if job_id in existing_ids:
raise SchedulerError(
f"任务 ID 已存在: {job_id},请使用唯一的 job_id",
scheduler="scheduler",
)
return job_id
module = getattr(func, '__module__', 'unknown')
qualname = getattr(func, '__qualname__', getattr(func, '__name__', str(func)))
auto_id = f"{prefix}_{module}.{qualname}"
if auto_id in existing_ids:
suffix = uuid.uuid4().hex[:8]
self._logger.warning(
f"自动生成的任务 ID 冲突: {auto_id}(同一函数注册多次),"
f"追加后缀消歧: {auto_id}_{suffix}"
)
auto_id = f"{auto_id}_{suffix}"
return auto_id
def add_cron_job(
self,
func: Callable,
cron: str,
job_id: Optional[str] = None,
**kwargs
) -> str:
"""
添加 Cron 任务
Args:
func: 要执行的函数
cron: Cron 表达式(标准5位格式:分 时 日 月 周,或使用 CronTrigger 支持的格式)
job_id: 任务ID,如果为 None 则自动生成
**kwargs: 其他任务参数
Returns:
str: 任务ID
"""
job_id = self._resolve_job_id(job_id, "cron", func)
try:
# 解析 cron 表达式
trigger = self._parse_cron(cron)
# 添加到 APScheduler
self._scheduler.add_job(
func=func,
trigger=trigger,
id=job_id,
**kwargs
)
self._logger.info(f"已添加 Cron 任务: {job_id} - {cron}")
return job_id
except Exception as e:
self._logger.error(f"添加 Cron 任务失败: {e}", exc_info=True)
raise
def _parse_cron(self, cron_expr: str) -> CronTrigger:
"""
解析 Cron 表达式
Args:
cron_expr: Cron 表达式字符串
Returns:
CronTrigger 对象(使用调度器配置的时区)
"""
# 尝试使用 CronTrigger.from_crontab 解析标准格式
try:
return CronTrigger.from_crontab(cron_expr, timezone=self._timezone)
except (ValueError, TypeError):
# 如果不是标准格式,尝试手动解析
parts = cron_expr.split()
if len(parts) == 5:
# 标准5位格式:分 时 日 月 周
minute, hour, day, month, day_of_week = parts
return CronTrigger(
minute=minute,
hour=hour,
day=day,
month=month,
day_of_week=day_of_week,
timezone=self._timezone
)
elif len(parts) == 6:
# 6位格式:秒 分 时 日 月 周(兼容旧格式)
second, minute, hour, day, month, day_of_week = parts
return CronTrigger(
second=second,
minute=minute,
hour=hour,
day=day,
month=month,
day_of_week=day_of_week,
timezone=self._timezone
)
else:
raise ValueError(f"无效的 Cron 表达式格式: {cron_expr},应为5位或6位")
def add_interval_job(
self,
func: Callable,
interval: int,
job_id: Optional[str] = None,
**kwargs
) -> str:
"""
添加间隔任务
Args:
func: 要执行的函数
interval: 间隔时间(秒)
job_id: 任务ID,如果为 None 则自动生成
**kwargs: 其他任务参数
Returns:
str: 任务ID
"""
job_id = self._resolve_job_id(job_id, "interval", func)
try:
# 创建间隔触发器(使用调度器配置的时区)
trigger = IntervalTrigger(seconds=interval, timezone=self._timezone)
# 添加到 APScheduler
self._scheduler.add_job(
func=func,
trigger=trigger,
id=job_id,
**kwargs
)
self._logger.info(f"已添加间隔任务: {job_id} - {interval}秒")
return job_id
except Exception as e:
self._logger.error(f"添加间隔任务失败: {e}", exc_info=True)
raise
def add_date_job(
self,
func: Callable,
run_date: str,
job_id: Optional[str] = None,
**kwargs
) -> str:
"""
添加一次性任务
Args:
func: 要执行的函数
run_date: 运行日期时间字符串,格式: 'YYYY-MM-DD HH:MM:SS' 或 'YYYY-MM-DD HH:MM'
job_id: 任务ID,如果为 None 则自动生成
**kwargs: 其他任务参数
Returns:
str: 任务ID
"""
job_id = self._resolve_job_id(job_id, "date", func)
try:
# 解析日期时间(返回 naive datetime
run_datetime = self._parse_run_date(run_date)
# 创建日期触发器(使用调度器配置的时区)
trigger = DateTrigger(run_date=run_datetime, timezone=self._timezone)
# 添加到 APScheduler
self._scheduler.add_job(
func=func,
trigger=trigger,
id=job_id,
**kwargs
)
self._logger.info(f"已添加一次性任务: {job_id} - {run_date}")
return job_id
except Exception as e:
self._logger.error(f"添加一次性任务失败: {e}", exc_info=True)
raise
def _parse_run_date(self, run_date: str) -> datetime:
"""
解析运行日期时间字符串
Args:
run_date: 日期时间字符串
Returns:
naive datetime 对象(APScheduler 会自动应用调度器的全局时区)
"""
# 尝试解析不同的日期格式
formats = [
'%Y-%m-%d %H:%M:%S',
'%Y-%m-%d %H:%M',
'%Y-%m-%d',
]
for fmt in formats:
try:
# 直接返回 naive datetimeAPScheduler 会自动应用全局时区
return datetime.strptime(run_date, fmt)
except ValueError:
continue
raise ValueError(f"无效的日期格式: {run_date},支持的格式: 'YYYY-MM-DD HH:MM:SS', 'YYYY-MM-DD HH:MM', 'YYYY-MM-DD'")
def remove_job(self, job_id: str) -> bool:
"""
移除任务
Args:
job_id: 任务ID
Returns:
bool: 是否成功移除
"""
try:
if self._scheduler.get_job(job_id):
self._scheduler.remove_job(job_id)
self._logger.info(f"已移除任务: {job_id}")
return True
return False
except Exception as e:
self._logger.error(f"移除任务失败: {e}", exc_info=True)
return False
def get_job(self, job_id: str) -> Optional[Any]:
"""
获取任务
Args:
job_id: 任务ID
Returns:
APScheduler Job 对象,如果不存在则返回 None
"""
return self._scheduler.get_job(job_id)
def list_jobs(self) -> list:
"""
列出所有任务ID
Returns:
list: 任务ID列表
"""
return [job.id for job in self._scheduler.get_jobs()]
def start(self) -> None:
"""启动调度器"""
# 检查是否启用
if not self._enabled:
self._logger.info("调度器已禁用,不启动")
return
if not self._scheduler.running:
try:
self._scheduler.start()
self._logger.info(f"任务调度器已启动 (时区: {self._timezone}, 任务数: {len(self._scheduler.get_jobs())})")
except Exception as e:
self._logger.error(f"启动调度器失败: {e}", exc_info=True)
raise
def stop(self) -> None:
"""停止调度器"""
if self._scheduler.running:
try:
self._scheduler.shutdown(wait=True)
self._logger.info("任务调度器已停止")
except Exception as e:
self._logger.error(f"停止调度器失败: {e}", exc_info=True)
def is_running(self) -> bool:
"""检查调度器是否运行中"""
return self._scheduler.running
def has_jobs(self) -> bool:
"""检查是否有任务"""
return len(self._scheduler.get_jobs()) > 0
def get_job_info(self, job_id: str) -> Optional[dict]:
"""
获取任务信息
Args:
job_id: 任务ID
Returns:
dict: 任务信息字典
"""
job = self._scheduler.get_job(job_id)
if not job:
return None
# 未启动调度器时 pending Job 对象没有 next_run_time 属性,用 getattr 防御
next_run_time = getattr(job, 'next_run_time', None)
info = {
'job_id': job_id,
'func_name': job.func.__name__ if hasattr(job.func, '__name__') else str(job.func),
'next_run_time': next_run_time.isoformat() if next_run_time else None,
}
# 根据触发器类型添加特定信息
trigger = job.trigger
if isinstance(trigger, CronTrigger):
info['type'] = 'cron'
info['cron'] = str(trigger)
elif isinstance(trigger, IntervalTrigger):
info['type'] = 'interval'
info['interval'] = trigger.interval.total_seconds()
elif isinstance(trigger, DateTrigger):
info['type'] = 'date'
info['run_date'] = trigger.run_date.isoformat() if trigger.run_date else None
return info
def list_all_jobs(self) -> list:
"""
列出所有任务的详细信息
Returns:
list: 任务信息列表
"""
return [self.get_job_info(job.id) for job in self._scheduler.get_jobs()]
def enable(self) -> None:
"""启用调度器
任务级 all_workers 支持:非 primary worker 默认调度器禁用,
当注册门控放行了 all_workers=True 的任务时调用本方法启用调度器,
使其能在 lifespan 中真正启动。
"""
if not self._enabled:
self._enabled = True
self._logger.debug("调度器已启用(all_workers 任务触发)")
def is_enabled(self) -> bool:
"""检查调度器是否启用"""
return self._enabled
def get_config(self) -> dict:
"""获取调度器配置"""
return {
'enabled': self._enabled,
'timezone': str(self._timezone) if self._timezone else 'system',
'running': self._scheduler.running,
'job_count': len(self._scheduler.get_jobs()),
'scheduled_job_count': len(self._scheduled_jobs)
}
def add_scheduled_job(
self,
job: 'ScheduledJob',
job_id: Optional[str] = None
) -> str:
"""
添加 ScheduledJob 对象到调度器
Args:
job: ScheduledJob 对象
job_id: 任务ID,如果为 None 则自动生成
Returns:
str: 任务ID
"""
from ..jobs.scheduled_job import ScheduledJob as ScheduledJobClass
if not isinstance(job, ScheduledJobClass):
raise TypeError(f"job 必须是 ScheduledJob 的实例,当前类型: {type(job)}")
if job.trigger is None:
raise ValueError("ScheduledJob 必须设置 trigger 属性")
# 生成任务ID
if job_id is None:
job_id = f"scheduled_{job.name}_{id(job)}"
# 保存 ScheduledJob 对象
self._scheduled_jobs[job_id] = job
job.job_id = job_id
# 转换触发器格式并添加到调度器
trigger = job.trigger
if isinstance(trigger, str):
# 字符串视为 cron 表达式
self.add_cron_job(
func=job.execute,
cron=trigger,
job_id=job_id,
name=job.name
)
elif isinstance(trigger, dict):
trigger_type = trigger.get('type')
if trigger_type == 'cron':
self.add_cron_job(
func=job.execute,
cron=trigger['cron'],
job_id=job_id,
name=job.name
)
elif trigger_type == 'interval':
interval = trigger.get('seconds', 0) or 0
interval += (trigger.get('minutes', 0) or 0) * 60
interval += (trigger.get('hours', 0) or 0) * 3600
interval += (trigger.get('days', 0) or 0) * 86400
if interval <= 0:
raise ValueError("间隔时间必须大于 0")
self.add_interval_job(
func=job.execute,
interval=interval,
job_id=job_id,
name=job.name
)
elif trigger_type == 'date':
self.add_date_job(
func=job.execute,
run_date=trigger['run_date'],
job_id=job_id,
name=job.name
)
else:
raise ValueError(f"不支持的触发器类型: {trigger_type}")
else:
raise ValueError(f"不支持的触发器类型: {type(trigger)}")
self._logger.info(f"已添加 ScheduledJob: {job.name} (ID: {job_id})")
return job_id
def get_scheduled_job(self, job_id: str) -> Optional['ScheduledJob']:
"""
获取 ScheduledJob 对象
Args:
job_id: 任务ID
Returns:
ScheduledJob 对象,如果不存在则返回 None
"""
return self._scheduled_jobs.get(job_id)
def get_all_scheduled_jobs(self) -> list:
"""
获取所有 ScheduledJob 对象
Returns:
ScheduledJob 对象列表
"""
return list(self._scheduled_jobs.values())
def add_job_object(
self,
job: 'ScheduledJob',
job_id: Optional[str] = None
) -> str:
"""
添加 ScheduledJob 对象到调度器(不设置触发器,用于非定时任务)
此方法用于添加不需要定时执行的任务,仅用于状态跟踪和管理。
如果需要定时执行,请使用 add_scheduled_job() 并设置 trigger。
Args:
job: ScheduledJob 对象
job_id: 任务ID,如果为 None 则自动生成
Returns:
str: 任务ID
"""
from ..jobs.scheduled_job import ScheduledJob as ScheduledJobClass
if not isinstance(job, ScheduledJobClass):
raise TypeError(f"job 必须是 ScheduledJob 的实例,当前类型: {type(job)}")
# 生成任务ID
if job_id is None:
job_id = f"job_{job.name}_{id(job)}"
# 保存 ScheduledJob 对象
self._scheduled_jobs[job_id] = job
job.job_id = job_id
self._logger.info(f"已添加任务对象: {job.name} (ID: {job_id})")
return job_id
def remove_scheduled_job(self, job_id: str) -> bool:
"""
移除 ScheduledJob
Args:
job_id: 任务ID
Returns:
是否成功移除
"""
if job_id in self._scheduled_jobs:
# 同时从调度器中移除(如果存在)
self.remove_job(job_id)
del self._scheduled_jobs[job_id]
self._logger.info(f"已移除 ScheduledJob: {job_id}")
return True
return False
# 全局调度器实例
_scheduler: Optional[Scheduler] = None
def get_scheduler() -> Scheduler:
"""获取调度器实例"""
global _scheduler
if _scheduler is None:
_scheduler = Scheduler()
return _scheduler