# -*- coding: utf-8 -*- """ 调度器 (Celery beat + worker) —— 设计 §10 调度总表 ==================================================== | 调度 | 时间 | 任务 | |-----------|----------------------|---------------------------------------------| | 拉上游计划| 交易日 08:40 | /plan 拉当日选股计划 + theme 灌行业映射表 | | 盘前准备 | 交易日 08:50 | T+1 可卖重置 / 参考位取数 / 刹车结算 | | 命令轮询 | 每 1 分钟 (全天) | 新命令解析 → 方案生成 → 任务状态机推进 | | 盘中执行 | 交易时段每 1 分钟 | 择时出手 + 自主提议扫描 (执行器下一批交付) | | 信号消化 | 交易时段每 1 分钟 | 订阅决策系统盘中信号 (下一批交付) | | 宏观择时 | 交易日 09:35 | 股汇对冲指数 → 升降仓命令或建议 + 宏观闸 | | 成交回放 | 交易时段每 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.macro_scan") @guard(trade_day=True, respect_exec_halt=True) def macro_scan(): """宏观择时 (09:35): 算股汇对冲指数 → 区域判定 → 升降仓命令或建议 + 个股宏观闸。 定在 09:35 而非盘前: 盘前无实时价, positions_view 会把持仓标 price_ok=False, planner._usable 排除这类票, 盘前下的降仓命令会被排成零方案直接取消。 09:35 行情已就位, 也在 08:40 拉候选池之后 (升仓有票可选)。指数用 T-1 日终数据, 几点计算数值都一样。详见 MACRO_TIMING_PLAN.md 第六节。 """ from app.services import macro_service return macro_service.scan() @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="*")}, "macro_scan": {"task": "pms.macro_scan", "schedule": crontab(hour=9, minute=35)}, "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)}, }