kangda/app/services/scheduler.py

89 lines
3.6 KiB
Python

import asyncio
from datetime import datetime
import logging
from sqlalchemy.exc import OperationalError, SQLAlchemyError
from app.core.database import async_session
from app.services.robot_sync_service import robot_sync_service
# 配置日志
logging.basicConfig(
level=logging.INFO,
format='%(asctime)s - %(levelname)s - %(message)s'
)
logger = logging.getLogger(__name__)
class Scheduler:
def __init__(self):
self.is_running = False
self.last_sync_time = None
self.retry_count = 0
self.max_retries = 3
self.retry_delay = 60 # 重试延迟(秒)
self.sync_interval = 300 # 同步间隔(秒)
async def start(self):
"""启动定时任务"""
if self.is_running:
logger.info("定时任务已经在运行中")
return
self.is_running = True
logger.info("定时任务启动")
while self.is_running:
try:
current_time = datetime.now()
if self.last_sync_time:
time_diff = (current_time - self.last_sync_time).total_seconds()
logger.info(f"距离上次同步已过去 {time_diff}")
# 同步机器人数据
async with async_session() as session:
try:
# 同步分组信息
logger.info("开始同步机器人数据...")
success = await robot_sync_service.sync_robot_data(session)
if not success:
raise Exception("同步机器人数据失败")
# 同步机器人任务信息
logger.info("开始同步机器人任务...")
await robot_sync_service.sync_robot_task(session)
self.last_sync_time = current_time
self.retry_count = 0 # 重置重试计数
logger.info("同步完成,等待下一次同步")
except OperationalError as e:
logger.error(f"数据库连接错误: {str(e)}")
if self.retry_count < self.max_retries:
self.retry_count += 1
logger.info(f"将在 {self.retry_delay} 秒后进行第 {self.retry_count} 次重试")
await asyncio.sleep(self.retry_delay)
continue
else:
logger.error("达到最大重试次数,等待下一次定时同步")
self.retry_count = 0 # 重置重试计数
except SQLAlchemyError as e:
logger.error(f"数据库操作错误: {str(e)}")
await asyncio.sleep(self.retry_delay)
except Exception as e:
logger.error(f"同步过程发生错误: {str(e)}")
await asyncio.sleep(self.retry_delay)
# 等待下一次同步
await asyncio.sleep(self.sync_interval)
except Exception as e:
logger.error(f"定时任务执行失败: {str(e)}")
await asyncio.sleep(self.retry_delay)
def stop(self):
"""停止定时任务"""
self.is_running = False
logger.info("定时任务已停止")
# 创建单例
scheduler = Scheduler()