169 lines
5.0 KiB
Python
169 lines
5.0 KiB
Python
from datetime import datetime
|
|
from typing import List
|
|
|
|
import structlog
|
|
from apscheduler.job import Job
|
|
from apscheduler.schedulers.background import BackgroundScheduler
|
|
from apscheduler.schedulers.base import STATE_PAUSED, STATE_STOPPED, STATE_RUNNING
|
|
from apscheduler.triggers.cron import CronTrigger
|
|
from apscheduler.events import EVENT_JOB_EXECUTED, EVENT_JOB_ERROR
|
|
|
|
from core.task.dlt_spider import DLTSpider
|
|
|
|
logger = structlog.getLogger(__name__)
|
|
|
|
def test_fun():
|
|
logger.info(f"当前时间---{datetime.now()}")
|
|
|
|
|
|
class LotteryScheduler:
|
|
"""彩票定时任务调度器"""
|
|
|
|
def __init__(self):
|
|
self.scheduler = BackgroundScheduler()
|
|
self.spider = DLTSpider()
|
|
self._setup_job_listeners()
|
|
|
|
def _setup_job_listeners(self):
|
|
"""设置任务监听器"""
|
|
|
|
def job_listener(event):
|
|
if event.exception:
|
|
logger.error(f"任务执行失败: {event.job_id}, 错误: {event.exception}")
|
|
else:
|
|
logger.info(f"任务执行成功: {event.job_id}")
|
|
|
|
self.scheduler.add_listener(job_listener, EVENT_JOB_EXECUTED | EVENT_JOB_ERROR)
|
|
|
|
def fetch_latest_job(self):
|
|
"""获取最新数据的定时任务"""
|
|
logger.info(f"开始执行定时任务: 获取最新大乐透数据")
|
|
try:
|
|
self.spider.fetch_latest()
|
|
logger.info("定时任务执行完成")
|
|
except Exception as e:
|
|
logger.error(f"定时任务执行异常: {e}", exc_info=True)
|
|
|
|
def add_job(self, *args, **kwargs):
|
|
"""添加任务"""
|
|
job = self.scheduler.add_job(*args, **kwargs)
|
|
logger.info(f"任务已添加: {job.id}, 下次执行: {job.next_run_time}")
|
|
return job
|
|
|
|
def remove_job(self, job_id: str):
|
|
"""移除任务"""
|
|
try:
|
|
self.scheduler.remove_job(job_id)
|
|
logger.info(f"任务已移除: {job_id}")
|
|
return True
|
|
except Exception as e:
|
|
logger.error(f"移除任务失败: {e}")
|
|
return False
|
|
|
|
def start(self):
|
|
"""启动或恢复调度器"""
|
|
if self.scheduler.state == STATE_RUNNING:
|
|
logger.info("调度器已在运行")
|
|
return
|
|
|
|
# 暂停状态:直接恢复
|
|
if self.scheduler.state == STATE_PAUSED:
|
|
self.scheduler.resume()
|
|
logger.info("调度器已恢复运行")
|
|
return
|
|
|
|
# 首次启动:注册任务
|
|
self.scheduler.add_job(
|
|
self.fetch_latest_job,
|
|
trigger=CronTrigger(hour=21, minute=30),
|
|
id='fetch_latest_dlt',
|
|
name='获取最新大乐透数据',
|
|
replace_existing=True
|
|
)
|
|
|
|
# 启动时立即执行一次(可选)
|
|
# self.scheduler.add_job(
|
|
# self.fetch_latest_job,
|
|
# trigger='date',
|
|
# run_date=datetime.now(),
|
|
# id='startup_fetch',
|
|
# name='启动时获取最新数据'
|
|
# )
|
|
|
|
self.scheduler.start()
|
|
logger.info("定时调度器已启动,将在每晚 21:30 执行")
|
|
|
|
# 打印所有任务
|
|
self.print_jobs()
|
|
|
|
def pause(self) -> bool:
|
|
"""暂停调度器"""
|
|
if self.scheduler.state == STATE_RUNNING:
|
|
self.scheduler.pause()
|
|
logger.info("定时调度器已暂停")
|
|
return True
|
|
else:
|
|
logger.info("调度器未在运行")
|
|
return False
|
|
|
|
def resume(self) -> bool:
|
|
"""恢复调度器"""
|
|
if self.scheduler.state == STATE_PAUSED:
|
|
self.scheduler.resume()
|
|
logger.info("定时调度器已恢复运行")
|
|
return True
|
|
else:
|
|
logger.info("定时调度器已未在运行状态")
|
|
return False
|
|
|
|
def shutdown(self):
|
|
"""彻底关闭调度器(应用退出时调用)"""
|
|
if self.scheduler.state != STATE_STOPPED:
|
|
self.scheduler.shutdown(wait=False)
|
|
logger.info("定时调度器已关闭")
|
|
|
|
def print_jobs(self):
|
|
"""打印所有定时任务"""
|
|
jobs = self.scheduler.get_jobs()
|
|
if jobs:
|
|
logger.info("当前定时任务:")
|
|
for job in jobs:
|
|
logger.info(f" - {job.id}: {job.name}, 下次执行: {job.next_run_time}")
|
|
else:
|
|
logger.info("当前没有定时任务")
|
|
|
|
def get_jobs(self) -> List[Job]:
|
|
"""获取所有定时任务"""
|
|
jobs = self.scheduler.get_jobs()
|
|
if jobs:
|
|
return jobs
|
|
else:
|
|
return []
|
|
|
|
|
|
# 创建全局调度器实例
|
|
scheduler = LotteryScheduler()
|
|
|
|
|
|
def start_scheduler():
|
|
"""启动调度器的入口函数"""
|
|
scheduler.start()
|
|
return scheduler
|
|
|
|
|
|
if __name__ == "__main__":
|
|
# 测试运行
|
|
scheduler = LotteryScheduler()
|
|
|
|
# 启动调度器
|
|
scheduler.start()
|
|
|
|
# 保持运行
|
|
try:
|
|
import time
|
|
|
|
while True:
|
|
time.sleep(60)
|
|
except KeyboardInterrupt:
|
|
logger.info("收到停止信号,正在关闭调度器...")
|
|
scheduler.shutdown() |