# -*- coding: utf-8 -*- """ 调度器 (Celery beat + worker) —— 设计 §10 调度总表 ==================================================== | 调度 | 时间 | 任务 | |-----------|----------------------|---------------------------------------------| | 拉上游计划| 交易日 08:40 | /plan 拉当日选股计划 + theme 灌行业映射表 | | 盘前准备 | 交易日 08:50 | T+1 可卖重置 / 参考位取数 / 刹车结算 | | 命令轮询 | 每 1 分钟 (全天) | 新命令解析 → 方案生成 → 任务状态机推进 | | 盘中执行 | 交易时段每 1 分钟 | 择时出手 + 自主提议扫描 (执行器下一批交付) | | 信号消化 | 交易时段每 1 分钟 | 订阅决策系统盘中信号 (下一批交付) | | 宏观择时 | 交易日 09:35 | 股汇对冲指数 → 升降仓命令或建议 + 宏观闸 | | 策略挂载 | 交易日 09:40 | 个股打法状态机: 吸筹挂网格/高热挂止盈/接力 | | 成交回放 | 交易时段每 1 分钟 | ws 逐笔入账 + trading_order 增量 + 轻对账 | | T 仓平回 | 14:50 (二期) | 做T强制平回 | | 日终结算 | 15:10 | 全量对账 / 除权 / 安全垫 / 命令进度日结 | | 净值快照 | 15:20 | 公示净值落一行 (导出表的净值序列) | | 运营日报 | 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 _exclusive(lease_key: str, ttl_sec: int = 115): """同名任务互斥 (2026-08-28 审查加): worker 双并发 + 每分钟一发 + 任务可跑到 300 秒, 重叠是常态配置 —— 而成交入账、方案物化这些路径不是可重入的 (重叠 = 双入账 / 双下单)。 租约存参数表 (INFRA_ 前缀不进页面): 拿到租约才跑, 没拿到直接返回 skipped。 写后回读校验缩小并发窗口; 租约超时 (ttl) 自动失效, 崩溃不会永久卡死。 读写参数表失败按"照常执行"降级 —— 互斥是保护, 不能反过来把任务停摆。""" def deco(fn): @functools.wraps(fn) def wrapper(*a, **kw): import json as _json import time as _time import uuid as _uuid from app.repo import pms_repo tok = _uuid.uuid4().hex[:12] try: raw = pms_repo.get_param(lease_key) if raw: d = _json.loads(raw) if d.get("tok") and _time.time() - float(d.get("ts") or 0) < ttl_sec: logger.info("[%s] 另一跳仍在跑 (租约 %.0fs 内), 本跳跳过", fn.__name__, ttl_sec) return {"skipped": "another_run_inflight"} pms_repo.set_param(lease_key, _json.dumps({"ts": _time.time(), "tok": tok}), "system") back = _json.loads(pms_repo.get_param(lease_key) or "{}") if back.get("tok") != tok: return {"skipped": "lease_lost"} except Exception as e: logger.warning("[%s] 租约获取失败 (照常执行): %s", fn.__name__, e) try: return fn(*a, **kw) finally: try: cur = _json.loads(pms_repo.get_param(lease_key) or "{}") if cur.get("tok") == tok: pms_repo.set_param(lease_key, "", "system") except Exception: pass 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.tech_pull") @guard(trade_day=True, respect_exec_halt=False) # 取数动作, 休假模式照跑 (只落读数更新映射, 不产生指令) def tech_pull(): """盘前很早 (06:30) 拉决策系统的全市场技术面读数, 落 pms_tech_daily 并建映射。 决策系统夜扫 23:20 起、实测 00:38 前完成, 06:30 留足余量。此刻今天的选股计划还没拉 (plan_pull 在 08:40), 映射先覆盖在持与待拍板; 08:40 拉完计划后 plan_pull 会补拉一次并 重建映射, 让覆盖面带上当天的主榜与观察档。拉失败只记, 当天按无读数, 不拦任何动作。 """ from app.services import tech_service return tech_service.pull_and_map() @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) out = {"ok": True, "date": plan["date"], "age_tdays": plan.get("age_tdays"), "returned": plan["returned"], "theme_sync": plan.get("theme_sync")} # 在持票的逻辑状态 (2026-09-07 第三件): 拉完计划顺带查一次, 写进运行参数并按结果停或恢复策略 # 买入腿。失败只记录, 不影响拉计划 —— 取不到就留空, 扫描层按没有读数处理。 try: from app.services import logic_state_service out["logic_state"] = logic_state_service.pull_for_held() except Exception as e: # noqa: BLE001 logger.error("[plan_pull] 逻辑状态取回失败: %s", e) out["logic_state"] = {"ok": False, "error": f"{type(e).__name__}: {e}"} # 技术面补拉并重建映射 (2026-09-11): 06:30 已拉过一次, 这里拉完计划后再拉一次并重建, # 让映射覆盖当天主榜与观察档的新票。失败只记, 不影响 plan_pull (取不到就按无读数)。 try: from app.services import tech_service out["tech"] = tech_service.pull_and_map() except Exception as e: # noqa: BLE001 logger.error("[plan_pull] 技术面补拉失败: %s", e) out["tech"] = {"ok": False, "error": f"{type(e).__name__}: {e}"} return out @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) @_exclusive("INFRA_LOCK_COMMAND_POLL", ttl_sec=115) 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) @_exclusive("INFRA_LOCK_REPLAY_FILLS", ttl_sec=115) 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) @_exclusive("INFRA_LOCK_INTRADAY_EXEC", ttl_sec=230) 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) @_exclusive("INFRA_LOCK_SIGNAL_DIGEST", ttl_sec=115) 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.strategy_attach") @guard(trade_day=True, respect_exec_halt=True) def strategy_attach(): """策略自动挂载 (09:40): 个股打法状态机 —— 吸筹挂网格 / 高热挂止盈 / 转派发暂停买入 / 接力。 定在 09:40: 在宏观扫描 (09:35) 之后 —— 若宏观当天要降仓, 先让降仓命令占住在途, 有在途的票本轮自动挂载会主动让路 (缓到下一个扫描日); 也避开开盘前 15 分钟的 竞价噪声, 此刻 positions_view 已有实时价, 网格区间锚得住。只做配置动作不下单, 真正买卖由 strategy_runner 逐笔过闸。详见 STRATEGY_AUTO_ATTACH_PLAN.md。 """ from app.services import strategy_advisor return strategy_advisor.scan() @celery_app.task(name="pms.t0_close") @guard(trade_day=True, respect_exec_halt=False) # 平回是守成收敛动作: 盘中开休假模式 def t0_close(): # 也不能让当日 T 仓裸奔过夜 (2026-08-28) """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 # 先把尾窗成交入账 (2026-08-28 审查修): replay_fills 15:05 停跑而这里 15:10 全量对账, # 15:05~15:10 落进 inbox 的成交若不先消化, 会被对账当数量差异按 RECON 补一次, # 次日开盘回放又正常入账一次 —— 同一笔账两份。休假模式下 replay 任务被闸住, # 这一步同样兜住 (入账是守成, 不产生新指令)。 pre = {} try: pre = ledger_service.replay_fills() except Exception as e: logger.error("[daily_settle] 尾窗成交入账失败 (对账将按差异吸收): %s", e) out = ledger_service.daily_settle() out.setdefault("steps", {})["pre_replay"] = { "ok": pre.get("ok"), "actions": pre.get("actions"), "ws_actions": (pre.get("ws") or {}).get("actions"), "errors": (pre.get("errors") or [])[:3]} 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.nav_snapshot") @guard(trade_day=True, respect_exec_halt=False) # 记账类动作, 休假模式照跑 def nav_snapshot(): """每日净值快照 (15:20): 按公示口径记一行净值 (pms_nav_daily), 供公示表导出的 净值序列使用。排在日终结算 (15:10) 之后 —— 尾窗成交已入账、收盘价已落定; 在运营日报 (15:30) 之前。同日重跑覆盖, 以最后一次为准。历史不回填 (2026-08-28 拍板: 启用日之前的净值继续查人工表格)。""" from app.services import publish_export return publish_export.record_nav_snapshot() @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 = { # 技术面取数 (2026-09-11): 决策系统夜扫实测 00:38 完, 06:30 留足余量; 在计划拉取 (08:40) 之前。 "tech_pull": {"task": "pms.tech_pull", "schedule": crontab(hour=6, minute=30)}, "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="*"), "options": {"expires": 50}}, # ws 通道的成交是推过来的, 落 inbox 后没必要再等 5 分钟才入账 —— 改成每分钟。 # 影子期这一跳基本是空转 (inbox 为空 + 下游无新成交), 成本可以忽略。 # expires=50: 队列积压时过期的火直接作废, 不叠着补跑 (配合各任务的互斥租约)。 "replay_fills": {"task": "pms.replay_fills", "schedule": crontab(minute="*"), "options": {"expires": 50}}, "intraday_exec": {"task": "pms.intraday_exec", "schedule": crontab(minute="*"), "options": {"expires": 50}}, "signal_digest": {"task": "pms.signal_digest", "schedule": crontab(minute="*"), "options": {"expires": 50}}, "macro_scan": {"task": "pms.macro_scan", "schedule": crontab(hour=9, minute=35)}, "strategy_attach": {"task": "pms.strategy_attach", "schedule": crontab(hour=9, minute=40)}, "t0_close": {"task": "pms.t0_close", "schedule": crontab(hour=14, minute=50)}, "daily_settle": {"task": "pms.daily_settle", "schedule": crontab(hour=15, minute=10)}, "nav_snapshot": {"task": "pms.nav_snapshot", "schedule": crontab(hour=15, minute=20)}, "daily_report": {"task": "pms.daily_report", "schedule": crontab(hour=15, minute=30)}, }