200 lines
8.5 KiB
Python
200 lines
8.5 KiB
Python
# -*- coding: utf-8 -*-
|
|
"""
|
|
调度器 (Celery beat + worker) —— 设计 §10 调度总表
|
|
====================================================
|
|
| 调度 | 时间 | 任务 |
|
|
|-----------|----------------------|---------------------------------------------|
|
|
| 盘前准备 | 交易日 08:50 | T+1 可卖重置 / 参考位取数 / 刹车结算 |
|
|
| 命令轮询 | 每 1 分钟 (全天) | 新命令解析 → 方案生成 → 任务状态机推进 |
|
|
| 盘中执行 | 交易时段每 1 分钟 | 择时出手 + 自主提议扫描 (执行器下一批交付) |
|
|
| 信号消化 | 交易时段每 1 分钟 | 订阅决策系统盘中信号 (下一批交付) |
|
|
| 成交回放 | 交易时段每 5 分钟 | 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
|
|
r = {"materialized": executor.materialize_plans()}
|
|
r.update(executor.run_tick())
|
|
return r
|
|
|
|
|
|
@celery_app.task(name="pms.signal_digest")
|
|
@guard(trade_day=True, session=True)
|
|
def signal_digest():
|
|
"""信号消化: 订阅决策系统盘中信号 (风控 SELL/止盈/反转) → 卖出方案或提议。下一批交付。"""
|
|
return {"consumed": 0, "note": "决策系统信号订阅为下一批交付 (设计 §10 信号消化)"}
|
|
|
|
|
|
@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="*")},
|
|
"replay_fills": {"task": "pms.replay_fills", "schedule": crontab(minute="*/5")},
|
|
"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)},
|
|
}
|