154 lines
5.4 KiB
Python
154 lines
5.4 KiB
Python
"""
|
|
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
|