tradingSystem/app/scheduler.py

234 lines
11 KiB
Python
Raw Normal View History

# -*- coding: utf-8 -*-
"""
调度器 (Celery beat + worker) 设计 §10 调度总表
====================================================
| 调度 | 时间 | 任务 |
|-----------|----------------------|---------------------------------------------|
| 拉上游计划| 交易日 08:40 | /plan 拉当日选股计划 + theme 灌行业映射表 |
| 盘前准备 | 交易日 08:50 | T+1 可卖重置 / 参考位取数 / 刹车结算 |
| 命令轮询 | 1 分钟 (全天) | 新命令解析 方案生成 任务状态机推进 |
| 盘中执行 | 交易时段每 1 分钟 | 择时出手 + 自主提议扫描 (执行器下一批交付) |
| 信号消化 | 交易时段每 1 分钟 | 订阅决策系统盘中信号 (下一批交付) |
2026-07-29 10:49:56 +08:00
| 成交回放 | 交易时段每 1 分钟 | ws 逐笔入账 + trading_order 增量 + 轻对账 |
| T 仓平回 | 14:50 (二期) | 做T强制平回 |
| 日终结算 | 15:10 | 全量对账 / 除权 / 安全垫 / 命令进度日结 |
| 运营日报 | 15:30 | 关注区 + 全量统计, 页面可查 |
三条守卫:
1. 交易日守卫 非交易日任务直接返回 (chinesecalendar; 未装则退化为周一至周五)
2. 故障即守成 任务内异常一律吞掉并记 ERROR 日志, 绝不因调度异常产生新指令
3. 全局暂停执行 (休假模式) 除对账与日报外的任务全部跳过
启用: docker compose --profile sched up -d
"""
from __future__ import annotations
import functools
import logging
from datetime import datetime
from celery import Celery
from celery.schedules import crontab
from config.settings import settings
from app.core import tradedays as td
logging.basicConfig(level=logging.INFO,
format="%(asctime)s %(levelname)s [%(name)s] %(message)s")
logger = logging.getLogger("pms.sched")
celery_app = Celery("pms", broker=settings.PMS_REDIS_URL, backend=settings.PMS_REDIS_URL)
celery_app.conf.update(
timezone="Asia/Shanghai", enable_utc=False,
task_serializer="json", result_serializer="json", accept_content=["json"],
result_expires=3600, worker_max_tasks_per_child=200,
task_soft_time_limit=240, task_time_limit=300,
broker_connection_retry_on_startup=True,
)
# 交易时段 (含集合竞价前 5 分钟余量)
SESSIONS = ((9, 25, 11, 30), (13, 0, 15, 5))
def in_session(now=None) -> bool:
now = now or datetime.now()
hm = now.hour * 60 + now.minute
return any(h1 * 60 + m1 <= hm <= h2 * 60 + m2 for h1, m1, h2, m2 in SESSIONS)
def guard(*, trade_day=True, session=False, respect_exec_halt=True):
"""任务守卫: 交易日 / 交易时段 / 全局暂停执行; 并统一吞异常 (守成)。"""
def deco(fn):
@functools.wraps(fn)
def wrapper(*a, **kw):
name = fn.__name__
if trade_day and not td.is_trade_day():
logger.info("[%s] 非交易日, 跳过", name)
return {"skipped": "not_trade_day"}
if session and not in_session():
return {"skipped": "not_in_session"}
if respect_exec_halt:
try:
from app.services import param_store
if param_store.get_bool("PMS_GLOBAL_EXEC_HALT", False):
logger.info("[%s] 全局暂停执行 (休假模式), 跳过", name)
return {"skipped": "exec_halt"}
except Exception as e:
logger.warning("[%s] 读取暂停开关失败 (按未暂停继续): %s", name, e)
t0 = datetime.now()
try:
r = fn(*a, **kw)
logger.info("[%s] 完成 %.1fs · %s", name, (datetime.now() - t0).total_seconds(),
_brief(r))
return r
except Exception as e:
logger.exception("[%s] 异常 (守成: 不产生新指令): %s", name, e)
return {"error": f"{type(e).__name__}: {e}"}
return wrapper
return deco
def _brief(r):
if not isinstance(r, dict):
return str(r)[:200]
keep = {k: v for k, v in r.items()
if k in ("ok", "fills", "actions", "planned", "failed", "diffs", "errors",
"skipped", "avail_reset", "refs", "cursor", "date", "returned")}
if isinstance(keep.get("diffs"), list):
keep["diffs"] = len(keep["diffs"])
if isinstance(keep.get("errors"), list):
keep["errors"] = keep["errors"][:2]
return str(keep)[:300]
# ================================================================ 任务
@celery_app.task(name="pms.plan_pull")
@guard(trade_day=True, respect_exec_halt=False) # 取数动作, 休假模式照跑 (只更新映射不产生指令)
def plan_pull():
"""盘前拉上游选股计划: 刷候选池缓存 + evidence.theme 灌进 pms_industry_map。
放在 premarket (08:50) 之前, 是因为盘前那一跳要用行业映射算集中度
拉失败照旧只记 ERROR 候选池当天就是空的, 升仓类命令会明确报"无票可选",
比拿昨天的榜静默买进去好
"""
from app.services import plan_feed
plan_feed.invalidate()
plan = plan_feed.get_plan(force=True)
return {"ok": True, "date": plan["date"], "age_tdays": plan.get("age_tdays"),
"returned": plan["returned"], "theme_sync": plan.get("theme_sync")}
@celery_app.task(name="pms.premarket")
@guard(trade_day=True)
def premarket():
from app.services import ledger_service
return ledger_service.premarket()
@celery_app.task(name="pms.command_poll")
@guard(trade_day=False, respect_exec_halt=True)
def command_poll():
"""命令轮询 (全天, 含非交易日 —— 用户随时可下命令, 方案先生成好待开盘执行)。"""
from app.services import command_service
r = command_service.plan_pending()
r.update(command_service.refresh_progress())
return r
@celery_app.task(name="pms.replay_fills")
@guard(trade_day=True, session=True)
def replay_fills():
from app.services import ledger_service
r = ledger_service.replay_fills()
light = ledger_service.reconcile(apply_fix=False) # 盘中轻对账: 只看差异不改账
r["light_recon"] = {"diffs": len(light.get("diffs") or []),
"severity": light.get("severity")}
return r
@celery_app.task(name="pms.intraday_exec")
@guard(trade_day=True, session=True)
def intraday_exec():
"""盘中执行: 方案转指令 → 择时出手 (规则闸终检 → 下发 → 记子单)。
含自主提议扫描 (FILL/ADD/DCA/TRIM 动作引擎 规则闸 研判闸 按档位分流)
"""
2026-08-11 10:35:03 +08:00
from app.services import executor, proposal_service, strategy_runner
r = {"materialized": executor.materialize_plans()}
r["proposals"] = proposal_service.scan_and_route()
2026-08-11 10:35:03 +08:00
r["strategies"] = strategy_runner.tick() # 个股交易方案: 触发即发短窗口指令, 由下面 run_tick 执行
2026-08-12 11:42:43 +08:00
# 过期收口改为每分钟一次 (原来只在 15:10 日结跑) —— 让过了窗口/deadline 的单当即作废,
# 不再多活一整天堵在「今日在办」、也不会隔天还替它试单 (2026-08-12 用户定的"实时"口径)。
# 先收口再 run_tick: 确保 run_tick 不会再去撮合一张本该已作废的单。
r["swept"] = executor.sweep_windows()
r.update(executor.run_tick())
return r
@celery_app.task(name="pms.signal_digest")
@guard(trade_day=True, session=True)
def signal_digest():
"""信号消化: 订阅决策系统盘中信号 (db2 盘中广播 + db3 风控 LLM 卖出) → 卖出指令或提议。"""
from app.services import signal_service
return signal_service.consume()
@celery_app.task(name="pms.t0_close")
@guard(trade_day=True)
def t0_close():
2026-08-11 11:58:16 +08:00
"""T 仓强制平回 (14:50): 对挂了做T策略且当日仍有未平腿的票, 立刻发对向平仓腿打平,
绝不过夜 (设计 § rails)平回后再自证账面 T , 残留的记 WARN 交人工核查
正T 的账面 t0_qty 在平回后仍可能 > 0: 今日买入的那份 T0 批次 T+1 才可卖, 平回卖的是
底仓存量, 净持仓已打平, T0 批次次日由对账并入底仓 这属正常, 不是没平回"""
from app.repo import pms_repo
2026-08-11 11:58:16 +08:00
from app.services import strategy_runner
r = strategy_runner.force_t0_close()
left = [p["ts_code"] for p in pms_repo.list_positions(only_open=True)
if int(p.get("t0_qty") or 0) > 0]
if left:
2026-08-11 11:58:16 +08:00
logger.warning("[t0_close] 平回后账面仍有 T0 批次: %s —— 多为今日买入次日才可卖(T+1)"
"或平仓腿尚未成交, 请人工核查", left)
return {"forced": r, "t0_open_after": left}
@celery_app.task(name="pms.daily_settle")
@guard(trade_day=True, respect_exec_halt=False) # 对账属守成动作, 休假模式下照跑
def daily_settle():
from app.services import executor, ledger_service
out = ledger_service.daily_settle()
try:
out["steps"]["windows"] = executor.sweep_windows() # 窗口耗尽收口 + 方案成交回写
except Exception as e:
out.setdefault("errors", []).append(f"窗口收口失败: {e}")
out["ok"] = False
try:
out["steps"]["expired_instructions"] = ledger_service.expire_stale_instructions()
except Exception as e:
out.setdefault("errors", []).append(f"指令过期处理失败: {e}")
return out
@celery_app.task(name="pms.daily_report")
@guard(trade_day=True, respect_exec_halt=False)
def daily_report():
from app.services import ledger_service
r = ledger_service.build_daily_report()
return {"ymd": r.get("ymd"), "attention": len(r.get("attention") or [])}
# ================================================================ beat 调度表
celery_app.conf.beat_schedule = {
"plan_pull": {"task": "pms.plan_pull", "schedule": crontab(hour=8, minute=40)},
"premarket": {"task": "pms.premarket", "schedule": crontab(hour=8, minute=50)},
"command_poll": {"task": "pms.command_poll", "schedule": crontab(minute="*")},
2026-07-29 10:49:56 +08:00
# ws 通道的成交是推过来的, 落 inbox 后没必要再等 5 分钟才入账 —— 改成每分钟。
# 影子期这一跳基本是空转 (inbox 为空 + 下游无新成交), 成本可以忽略。
"replay_fills": {"task": "pms.replay_fills", "schedule": crontab(minute="*")},
"intraday_exec": {"task": "pms.intraday_exec", "schedule": crontab(minute="*")},
"signal_digest": {"task": "pms.signal_digest", "schedule": crontab(minute="*")},
"t0_close": {"task": "pms.t0_close", "schedule": crontab(hour=14, minute=50)},
"daily_settle": {"task": "pms.daily_settle", "schedule": crontab(hour=15, minute=10)},
"daily_report": {"task": "pms.daily_report", "schedule": crontab(hour=15, minute=30)},
}