# -*- coding: utf-8 -*- """ 调度器 (Celery beat + worker) —— 设计 §10 调度总表 ==================================================== | 调度 | 时间 | 任务 | |-----------|----------------------|---------------------------------------------| | 盘前准备 | 交易日 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")} 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.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 动作引擎 → 规则闸 → 研判闸 → 按档位分流)。 """ from app.services import executor, proposal_service r = {"materialized": executor.materialize_plans()} r["proposals"] = proposal_service.scan_and_route() 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为二期上线, 此处先留调度位并自证 T 仓应为 0。""" from app.repo import pms_repo 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] 仍有 T 仓未平回: %s (做T为二期功能, 请人工核查)", left) return {"t0_open": left, "phase": "二期"} @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 = { "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)}, }