""" scheduler.py — 后台自动调度:定时拉取数据 + 计算因子 + 制度判定。 Python main.py cron → 前台阻塞运行,每 N 分钟一个 tick Web app 启动时自动启动调度器 → 后台线程,不阻塞 Web 请求 """ import threading import logging import time from datetime import datetime, timezone, timedelta from typing import Optional logger = logging.getLogger("chanmacro.scheduler") class MacroScheduler: """后台调度器:定时 fetch + score。""" def __init__(self, interval_minutes: int = 60): self.interval = interval_minutes self._thread: Optional[threading.Thread] = None self._stop = threading.Event() self._last_run: Optional[datetime] = None self._running = False def start(self) -> None: """启动后台线程。""" if self._running: return self._stop.clear() self._thread = threading.Thread(target=self._loop, name="macro-scheduler", daemon=True) self._thread.start() self._running = True logger.info(f"调度器已启动, 每 {self.interval} 分钟执行一次") def stop(self) -> None: """停止后台线程。""" self._stop.set() self._running = False logger.info("调度器已停止") @property def last_run(self) -> Optional[datetime]: return self._last_run def _loop(self) -> None: """后台循环。""" # 首次启动立即跑一次 self._tick() while not self._stop.wait(self.interval * 60): self._tick() def _tick(self) -> None: """执行一次:fetch → score。""" try: from fetchers.ohlcv import OHLCVFetcher from fetchers.breadth import BreadthFetcher from fetchers.derivatives import DerivativesFetcher from database import init_db from datetime import date as Date init_db() today = Date.today() # Fetch ohlcv = OHLCVFetcher() df = ohlcv.fetch() if not df.empty: ohlcv.store_df(df) breadth = BreadthFetcher() record = breadth.fetch() if record: breadth.store(record=record) deriv = DerivativesFetcher() records = deriv.fetch(today) if records: deriv.store(records=records) # Score + Regime (also persisted inside _build_state) from scoring.price_structure import PriceStructureScorer from scoring.breadth_scorer import BreadthScorer from scoring.oi_matrix import OIMatrixScorer from scoring.volatility_regime import VolatilityRegimeScorer from regime_detector import RegimeDetector from models import MarketStateVector from config import config import json from database import get_connection ps = PriceStructureScorer().compute(today) br = BreadthScorer().compute(today) oi = OIMatrixScorer().compute(today) vol = VolatilityRegimeScorer().compute(today) detector = RegimeDetector() detector.load_state(config.db_path) r = detector.detect(ps.score, br.breadth_top50, vol.vol_regime.value, today) conn = get_connection() conn.execute(""" INSERT OR REPLACE INTO regime_history (date, regime, confidence, regime_version, maturity_score, all_scores_json, prior_regime, confirmation_days) VALUES (?, ?, ?, ?, ?, ?, ?, ?) """, ( str(today), r.regime.value, r.confidence, r.regime_version, r.maturity_score, json.dumps(r.all_scores), r.prior_regime.value if r.prior_regime else None, r.confirmation_days, )) conn.commit() conn.close() # 检测新信号(每天运行一次,UTC 0 点后首次触发) now = datetime.now(timezone.utc) if self._last_run is None or now.date() > self._last_run.date(): try: from chan_integration import ChanSignalDetector detector = ChanSignalDetector() # 检测最近 90 天的 4h 信号 count = detector.populate_signal_features( start_date=(today - __import__('datetime').timedelta(days=90)).isoformat(), end_date=today.isoformat(), ) if count > 0: logger.info(f"新增 {count} 条信号记录") except Exception as e: logger.debug(f"信号检测跳过: {e}") self._last_run = now logger.info( f"Tick 完成: regime={r.regime.value} conf={r.confidence:.2f} " f"breadth={br.score:.0f}({br.breadth_bucket.value}) " f"price={ps.score:.0f} oi={oi.oi_state.value} vol={vol.vol_regime.value}" ) except Exception as e: logger.error(f"Tick 失败: {e}", exc_info=True) # 单例 _scheduler: Optional[MacroScheduler] = None def get_scheduler() -> MacroScheduler: global _scheduler if _scheduler is None: _scheduler = MacroScheduler(interval_minutes=60) return _scheduler