72 lines
2.3 KiB
Python
72 lines
2.3 KiB
Python
from contextlib import asynccontextmanager
|
|
from datetime import datetime, timedelta
|
|
|
|
from apscheduler.schedulers.asyncio import AsyncIOScheduler
|
|
from fastapi import FastAPI
|
|
|
|
from app.client.embedding_client_manager import embedding_client
|
|
from app.client.milvus_client_manager import milvus_client
|
|
from app.client.mysql_client_manager import db_assistant_mysql_client_manager
|
|
from app.client.redis_client_manager import redis_client_manager
|
|
from app.core.log import logger
|
|
from app.summary.service import build
|
|
|
|
"""
|
|
生命周期
|
|
容器启动时执行
|
|
yield
|
|
容器停止时执行
|
|
"""
|
|
scheduler = AsyncIOScheduler(timezone='Asia/Shanghai')
|
|
|
|
async def generate_daily_summary():
|
|
yesterday = datetime.now() - timedelta(days=1)
|
|
date = datetime(yesterday.year, yesterday.month, yesterday.day)
|
|
date_format = f"{date.year}-{date.month}-{date.day}"
|
|
|
|
lock_key = f"daily_summary_lock_{date_format}"
|
|
lock_expire = 60
|
|
|
|
try:
|
|
acquired = await redis_client_manager.client.set(lock_key, "1", ex=lock_expire, nx=True)
|
|
if not acquired:
|
|
logger.info(f"{date_format} 消息摘要任务已在其他实例执行,跳过")
|
|
return
|
|
|
|
logger.info(f"{date_format} 消息摘要任务获取锁成功,开始执行")
|
|
await build(date_format)
|
|
logger.info(f"{date_format} 消息摘要生成完成")
|
|
except Exception as e:
|
|
logger.error(f"{date_format} 消息摘要任务执行失败: {str(e)}")
|
|
finally:
|
|
try:
|
|
await redis_client_manager.client.delete(lock_key)
|
|
except Exception as e:
|
|
logger.warning(f"释放锁失败: {e}")
|
|
|
|
@asynccontextmanager
|
|
async def lifespan(app: FastAPI):
|
|
db_assistant_mysql_client_manager.init()
|
|
embedding_client.init()
|
|
milvus_client.init()
|
|
redis_client_manager.init()
|
|
|
|
scheduler.add_job(
|
|
generate_daily_summary,
|
|
trigger='cron',
|
|
hour=2,
|
|
minute=0,
|
|
second=0,
|
|
id='daily_summary',
|
|
name='每日消息摘要任务',
|
|
replace_existing=True
|
|
)
|
|
scheduler.start()
|
|
logger.info("定时任务调度器已开启")
|
|
yield
|
|
await db_assistant_mysql_client_manager.close()
|
|
await embedding_client.close()
|
|
milvus_client.close()
|
|
await redis_client_manager.close()
|
|
scheduler.shutdown()
|
|
logger.info("定时任务调度器已关闭") |