241 lines
11 KiB
Python
241 lines
11 KiB
Python
# -*- coding: utf-8 -*-
|
|
"""
|
|
调度器 (Celery beat + worker) —— 设计 §10 调度总表
|
|
====================================================
|
|
| 调度 | 时间 | 任务 |
|
|
|-----------|----------------------|---------------------------------------------|
|
|
| 拉上游计划| 交易日 08:40 | /plan 拉当日选股计划 + theme 灌行业映射表 |
|
|
| 盘前准备 | 交易日 08:50 | T+1 可卖重置 / 参考位取数 / 刹车结算 |
|
|
| 命令轮询 | 每 1 分钟 (全天) | 新命令解析 → 方案生成 → 任务状态机推进 |
|
|
| 盘中执行 | 交易时段每 1 分钟 | 择时出手 + 自主提议扫描 (执行器下一批交付) |
|
|
| 信号消化 | 交易时段每 1 分钟 | 订阅决策系统盘中信号 (下一批交付) |
|
|
| 成交回放 | 交易时段每 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())
|
|
# 清仓完成后清场: 撤该股策略/在途建仓/相关提议 (幂等, 每只闭仓票只清一次)。
|
|
# 挂在这里而非盘中执行位, 是因为它与持仓状态相关、非交易日照样该清干净。
|
|
try:
|
|
r["cleanup"] = command_service.cleanup_exited_positions()
|
|
except Exception as e:
|
|
logger.error("[command_poll] 清场失败 (不影响命令轮询): %s", e)
|
|
r["cleanup_error"] = str(e)
|
|
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 动作引擎 → 规则闸 → 研判闸 → 按档位分流)。
|
|
"""
|
|
from app.services import executor, proposal_service, strategy_runner
|
|
r = {"materialized": executor.materialize_plans()}
|
|
r["proposals"] = proposal_service.scan_and_route()
|
|
r["strategies"] = strategy_runner.tick() # 个股交易方案: 触发即发短窗口指令, 由下面 run_tick 执行
|
|
# 过期收口改为每分钟一次 (原来只在 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():
|
|
"""T 仓强制平回 (14:50): 对挂了做T策略且当日仍有未平腿的票, 立刻发对向平仓腿打平,
|
|
绝不过夜 (设计 §六 rails)。平回后再自证账面 T 仓, 残留的记 WARN 交人工核查。
|
|
|
|
正T 的账面 t0_qty 在平回后仍可能 > 0: 今日买入的那份 T0 批次 T+1 才可卖, 平回卖的是
|
|
底仓存量, 净持仓已打平, 该 T0 批次次日由对账并入底仓 —— 这属正常, 不是没平回。"""
|
|
from app.repo import pms_repo
|
|
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:
|
|
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="*")},
|
|
# 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)},
|
|
}
|