scheduler: auto fetch+score every 60min, integrated into web and CLI
This commit is contained in:
+19
-1
@@ -5,6 +5,7 @@ cli.py — Command-line interface for ChanMacro.
|
||||
import argparse
|
||||
import json
|
||||
import logging
|
||||
import time
|
||||
from datetime import date as Date, datetime, timedelta
|
||||
|
||||
logging.basicConfig(
|
||||
@@ -331,6 +332,8 @@ def main():
|
||||
# validate
|
||||
p_validate = sub.add_parser("validate", help="Run validation framework")
|
||||
|
||||
# cron
|
||||
p_cron = sub.add_parser("cron", help="Run scheduled fetch+score loop")
|
||||
# serve
|
||||
p_serve = sub.add_parser("serve", help="Start web dashboard")
|
||||
|
||||
@@ -353,7 +356,22 @@ def main():
|
||||
report = ValidationReporter().run_all()
|
||||
print(report)
|
||||
elif args.command == "serve":
|
||||
logger.info("Web dashboard not yet implemented (Phase 7)")
|
||||
from scheduler import get_scheduler
|
||||
get_scheduler().start()
|
||||
logger.info("启动 Web Dashboard: http://127.0.0.1:8124")
|
||||
from web.app import app
|
||||
app.run(host="0.0.0.0", port=8124, debug=False)
|
||||
elif args.command == "cron":
|
||||
from scheduler import get_scheduler
|
||||
logger.info("启动后台调度器 (Ctrl+C 停止)")
|
||||
s = get_scheduler()
|
||||
s.start()
|
||||
try:
|
||||
while True:
|
||||
time.sleep(60)
|
||||
except KeyboardInterrupt:
|
||||
s.stop()
|
||||
logger.info("调度器已停止")
|
||||
else:
|
||||
parser.print_help()
|
||||
|
||||
|
||||
@@ -0,0 +1,137 @@
|
||||
"""
|
||||
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()
|
||||
|
||||
self._last_run = datetime.now(timezone.utc)
|
||||
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
|
||||
@@ -173,4 +173,6 @@ def api_expectancy():
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
from scheduler import get_scheduler
|
||||
get_scheduler().start()
|
||||
app.run(host="0.0.0.0", port=8124, debug=True)
|
||||
|
||||
Reference in New Issue
Block a user