Files
stock/backend/app/scheduler.py
2026-09-09 15:07:58 +08:00

115 lines
4.7 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.
"""夜间定时任务收盘后自动全市场同步A股 + ETF+ 过期会话清理。
进程内 asyncio 循环单进程部署假设。start_sync 均幂等(已在跑直接返回),
即使多 worker / 手动触发与定时撞车也不会重复跑。关闭NIGHTLY_SYNC_ENABLED=false。
补跑语义:进程启动时若已过触发点、今天是交易日、且当日 candles 尚未落库
(例如定时点机器没开机),立即补跑一次,不让数据断档等到第二天。
"""
from __future__ import annotations
import asyncio
import logging
from datetime import datetime, timedelta
from sqlalchemy import delete, func, or_, select
from .auth import utcnow
from .config import settings
from .models import AuthSession, Candle, TradeCalendar
log = logging.getLogger(__name__)
_TRIGGER_MINUTE = 5 # 触发点 = nightly_sync_hour:05避开整点拥挤纯本地任务习惯
def _seconds_until_next_run() -> float:
now = datetime.now()
target = now.replace(hour=settings.nightly_sync_hour, minute=_TRIGGER_MINUTE,
second=0, microsecond=0)
if target <= now:
target += timedelta(days=1)
return (target - now).total_seconds()
async def cleanup_sessions() -> int:
"""删除过期 / 吊销超 7 天的会话行(登录路径只清本人,长跑进程需要兜底)。"""
from .db import async_session # 延迟导入避免循环
now = utcnow()
async with async_session() as s:
res = await s.execute(delete(AuthSession).where(or_(
AuthSession.expires_at < now,
AuthSession.revoked_at.is_not(None) & (AuthSession.revoked_at < now - timedelta(days=7)),
)))
await s.commit()
return res.rowcount or 0
async def _nightly_routine() -> None:
"""当日例行A 股全市场同步 -> ETF 同步(先后跑,避免两路 tushare 控频挤兑)-> 会话清理。"""
from .data import etf_sync
from .db import async_session
from .screener import market_sync
log.info("夜间任务开始:全市场同步窗口 %d 个交易日", settings.screener_market_days)
try:
async with async_session() as session:
await market_sync.start_sync(session, days=settings.screener_market_days, force=False)
# 等日线同步收尾再触发 ETF轮询模块内状态start_sync 是即发即忘的)
while market_sync._sync_state["running"]:
await asyncio.sleep(30)
except Exception: # noqa: BLE001
log.exception("夜间 A 股同步触发失败")
try:
await etf_sync.start_sync(full=False)
except Exception: # noqa: BLE001
log.exception("夜间 ETF 同步触发失败")
try:
n = await cleanup_sessions()
if n:
log.info("夜间会话清理:%d", n)
except Exception: # noqa: BLE001
log.exception("夜间会话清理失败")
log.info("夜间任务结束")
async def _should_run_on_startup() -> bool:
"""启动补跑判定:已过触发点 + 今天是交易日 + 当日 candles 还没落库。"""
now = datetime.now()
if now.hour * 60 + now.minute < settings.nightly_sync_hour * 60 + _TRIGGER_MINUTE:
return False
today8 = now.strftime("%Y%m%d")
from .db import async_session
async with async_session() as s:
is_trade_day = bool(await s.scalar(
select(TradeCalendar.id).where(TradeCalendar.trade_date == today8).limit(1)))
if not is_trade_day:
return False
# 当日未收盘/数据未生成时 max(ts) < 今日,同步会拉到空——那正是要补跑的信号
latest = await s.scalar(select(func.max(Candle.ts)).where(Candle.timeframe == "1d"))
return latest is None or latest.strftime("%Y%m%d") < today8
async def run_nightly_loop() -> None:
"""每日触发点跑一次;启动时满足补跑条件先补跑。由 main.lifespan 拉起。"""
if not settings.nightly_sync_enabled:
log.info("夜间自动同步未启用NIGHTLY_SYNC_ENABLED=false")
return
try:
if await _should_run_on_startup():
log.info("启动补跑:已过 %02d:%02d 且当日数据未落库", settings.nightly_sync_hour, _TRIGGER_MINUTE)
await _nightly_routine()
except Exception: # noqa: BLE001
log.exception("启动补跑判定失败(跳过,等待每日定时点)")
while True:
try:
await asyncio.sleep(_seconds_until_next_run())
await _nightly_routine()
except asyncio.CancelledError:
raise
except Exception: # noqa: BLE001 —— 循环体不能死
log.exception("夜间任务循环异常60s 后继续")
await asyncio.sleep(60)