任务中心与进程模型
概述
任务中心回答三个问题:
- 现在在跑什么、跑到哪了(运行中任务的实时进度分子/分母)
- 上次为什么失败、失败在哪一步(步骤级留痕 + 自动关联的 WARNING 日志与 traceback)
- 过去 30 天稳不稳定(
app_task_run历史聚合、失败率 Top、平均耗时)
它与进程模型是同一件事的两面:调度器只有跑在独立 worker 进程里,长耗时的全市场融合
链路才不会被 web 侧热重载腰斩;而拆进程之后,"现在在跑什么"这件事必须从数据库读,
不能再读 web 进程的内存态——这就是三层运行记录(app_task_run / app_task_step /
app_task_event)与命令队列存在的理由。
页面入口:/tasks(需管理员权限)。写入侧 backend/app/services/task_center/,
读取侧 backend/app/api/tasks.py,ORM backend/app/models/task_models.py。
进程模型:XUANGU_ROLE 三角色
定义见 backend/app/core/roles.py。
| 角色 | 谁在用 | 行为 |
|---|---|---|
all(默认) | 单进程部署 / 测试 | API + 调度器同进程(= 历史行为) |
worker | backend/worker.py(dev 由 run.py 拉起;prod 独立 systemd unit) | 只跑调度器 + 命令队列 + 心跳,不起 HTTP |
web | uvicorn 进程 | 只服务 API,不启动调度器、不参与 leader 竞选 |
current_role() / scheduler_in_process() / api_in_process() 每次都读环境变量。
测试用 monkeypatch.setenv 切角色,uvicorn 子进程与 worker 子进程的环境变量各自独立;
一旦缓存成模块常量,"改环境变量即换角色"就失效了。生产不设该变量即 all。
为什么要拆:uvicorn --reload 每次保存代码都杀子进程,15 分钟级的融合任务跑在 web
进程里必被腰斩(收尾 finally 不执行 → app_task_run 卡在 running),且热重载会重复触发
整点任务。拆开后 web 侧随便重载,任务照跑。代价:worker 不随热重载,改调度器/任务代码
须重启 run.py。
prod 下不经 run.py:systemd 成对拉起 xuangu.service(web)与 xuangu-worker.service
(worker),见 deploy/ 目录。
调度循环
实现:backend/app/services/scheduler.py(纯标准库线程,零外部依赖)。
| 周期常量 | 值 | 作用 |
|---|---|---|
_TICK_INTERVAL | 30s | 主循环唤醒间隔(含 leader 竞选重试) |
_WARMUP_INTERVAL | 60s | 交易时段预取 6 类市场聚合数据 → 写缓存 + 快照 |
_BASEINFO_INTERVAL | 600s | 全市场基础信息刷新(保障涨跌幅分布/行业 TOP5 新鲜) |
_TDX_SYNC_INTERVAL | 600s | 交易时段 TDX 增量同步(五档采样 + 资金流轮转) |
| 命令轮询 | 2s | 领取跨进程手动触发命令 |
设计要点:单个常驻事件循环(run_until_complete),避免反复 asyncio.run 导致
market_data_service 里缓存的 httpx.AsyncClient 绑定到已关闭的循环;全部任务
best-effort,DB 不可用或异常都安全跳过,不拖垮主服务。
leader 锁与「谁在跑」的三态答案
app/core/leader_lock.py 用 MySQL GET_LOCK 做 leader 租约,解决 uvicorn --workers 4
时调度器在每个 worker 各跑一份的问题:
- 锁绑定在独立 NullPool 连接上——连接关闭即释放锁,绝不污染业务连接池 (业务池归还的连接若自带 GET_LOCK 结果,会让其他 worker 误判为 leader)
- 非 MySQL / 锁不可用(SQLite 测试、本地开发)时降级为全部执行,单 worker 无竞争
但 owns() 是进程内状态:拆进程后 web 里 owns("scheduler") 恒为 False,
前端徽章就恒显"由其他实例承载"。故调度器存活改按 worker 心跳判断
(task_center/heartbeat.py):
| SysConfig 键 | 值 |
|---|---|
runtime.worker.heartbeat_at | 本地时间 ISO(worker 与读数方同机部署,不跨时区) |
runtime.worker.pid | 进程号(展示 + 兜底排查) |
runtime.worker.worker_id | hostname:pid(判活用,与 app_task_run.worker_id 同口径) |
runtime.worker.role | 角色(worker / all) |
get_scheduler_status() 返回 host 三态,前端徽章必须按 host 分支:
host | 含义 |
|---|---|
local | 本进程持有 MySQL GET_LOCK |
remote | 心跳新鲜(≤ FRESH_SECONDS=90),由别的进程承载 |
none | 真停了(leader = host != "none") |
🔴 threadAlive 只描述本进程的调度线程,XUANGU_ROLE=web 下恒为 False——
拿它下结论必然误报"调度器已停止"。判活是两级:心跳新鲜 → 在跑;心跳过期但
owner_alive(worker_id) 明确为 False → 立刻判停(不必等满 90s);跨主机/判不出
→ 只看新鲜度(保守方向:不误报"已停止")。
任务注册表与补跑
67 个任务的完整时点表由注册表生成,见数据流转 · 定时任务时序。
唯一权威来源是 services/scheduler_tasks.py::SCHEDULED_TASKS,触发时间可在运行时经
SysConfig sched.{key}.trigger_hm 覆盖,PUT /datamgr/schedule/{key} 修改后下个 tick 即生效。
任务成功后写入 sched.{key}.last_success_date 的必须是 _last_expected_date(key, started_at)
(数据日)。凌晨/盘前补跑执行的是前一个应执行日的数据,若记成"今天",当日正式时点
会被 last_success >= expected 判成"已完成"而静默跳过——实测 09-21 00:55 为 09-18 数据
补跑 tactic_signals,当晚 18:20 的正式任务被跳过,整日无信号快照且无任何告警。
调度去重的 last_run_date(进程内 + 启动回填)同样必须按应执行日写。
启动回填 seed_last_run_date_from_db 只认有产出的状态(success / partial /
skipped),failed 刻意不回填——失败的数据缺口还在,回填后当日正式时点会被
"当日已跑过就不再跑"静默跳过。
依赖门控:depends_on
2026-10-07 引入。实现见
services/scheduler_tasks.py::_depends_blocked(), 由run_due_tasks()在 due-gate 与last_success检查之后调用。
调度表里的触发时点只是近似顺序,不是约束。全市场链路的实际耗时随数据量与
上游可用性漂移:fox_kline_sync(15:40 起)实测常跑到 17:45 之后,而按死时点
17:30 起跑的 quant_calc 会撞上「当日 K 线只写了一半」并记成 partial
(2026-09-30 实锤)。依赖门控把「谁等谁」从注释里的约定变成注册表里的声明。
声明方式与现有字段同级:
"quant_calc": {
"trigger_hm": 1730,
"async_task": True,
"depends_on": ["fox_kline_sync"], # ← 上游未交付本次应执行日,本轮不起跑
...
}
当前登记的依赖边:
| 任务 | 依赖 | 为什么 |
|---|---|---|
quant_calc | fox_kline_sync | 读 fox_kline_daily 当日数据 |
unified_selection_sync | quant_calc | 指标列富化 |
factor_score_track | quant_calc | 因子分留痕 |
strategy_auto_run | selection_sync、quant_calc | 用户策略读宽表/指标 |
tactic_signals | fox_kline_sync、quant_calc | 战法扫描要当日 K 线与形态列 |
skill_cache_warmup | fox_wide_sync | 技能 DSL 预热读宽表 |
preset_track | fox_wide_sync、quant_calc | 14 预设留痕 |
builtin_track | dragon_score_sync | VPA 候选池取自当日龙头评分 |
逐个依赖 dep 判,顺序即语义:
last_success_date(dep) >= 本次应执行日→ 放行。此时 dep 若还在跑, 必是补跑/回填历史,不该把下游无限按住;- dep 正在跑且未在应执行日成功 → 按住,等它收尾;
- dep 没在跑但当日槽位还没执行过(进程刚重启、补跑尚未轮到)→ 也按住。
这一支封的是穿透路径:quant_calc 被它的上游按住时,若只看"dep 在跑",
unified_selection_sync会拿昨天的指标抢跑,抢跑成功后当日不再重算; - dep 当日已执行过但没成功(真失败/重试耗尽)→ 不拦。让下游照跑,由各自的 新鲜度/覆盖度守卫显性拒绝。饿死一整天是静默停摆,失败看得见。
依赖任务未注册或被禁用 → 跳过不拦;跳过时不写 last_run_date(与 due-gate 的
blocked 分支同纪律,否则当日槽位会被误判成"已跑过"而静默跳过);下轮 tick(30s)
自然重试。上游为 partial 时 last_success_date 不推进(run_task 只对
success/skipped 记账),所以"在跑就等"天然覆盖 partial 的收尾窗口。
- 上游的
trigger_hm必须 ≤ 下游。run_due_tasks的遍历顺序按 trigger_hm 升序 (_tick_order()),这样同一轮 tick 内上游会先被尝试;写反了会把下游按到第二天。 注册表测试test_depends_on_registry_invariants锁死这一条(连同"不许依赖未注册的 key""不许自己依赖自己")。 - 手动触发不经此门:命令队列与
trigger_task_now代表管理员的显式意图, 顺序由人负责。 - 长耗时任务要配
async_task:门控让下游"晚几十秒起跑"成为常态,若上游是同步 执行,一次 tick 会被它堵满,下游等的是 tick 而不是上游。quant_calc(实测均值 12min)与skill_cache_warmup(51min)已异步化。
心跳除调度 tick 外还由命令轮询线程每 2s heartbeat.touch()
(task_center/commands.py::start_poller),因此即便某个 51min 的任务把 _tick
堵满,任务中心的 host 徽章也不会误红。不需要额外的 watchdog。
三层运行记录
app_task_run 一次执行(trigger_type: scheduled/manual/api/backfill)
├─ app_task_step 子步骤(step_key=fox 的 data_type / 其他步骤名;stage;seq)
└─ app_task_event 日志事件(level / logger / message / traceback)
| 表 | 关键字段 | 说明 |
|---|---|---|
app_task_run | task_key task_label trigger_type status attempt started_at finished_at duration_ms total success failed progress_done progress_total message detail_json worker_id | status 取 running/success/failed/skipped/interrupted;progress_* 是实时进度分子分母 |
app_task_step | run_id step_key stage seq status progress_done/total/failed message source_task_id | stage 取 raw_heal/stg_load/dwd_fuse/collect/calc/report;source_task_id 关联 fox_sync_task.id |
app_task_event | run_id(可空) step_id(可空) task_key level logger message traceback created_at | 错误日志面板的数据源 |
为什么不复用旧表:app_sync_log 只有扁平单行,缺层级(同一时刻跑的 N 条 fox 同步无法
归属到某次 fox_daily_sync)、缺元信息(无触发方式/worker/attempt);fox_sync_task
粒度最细但只覆盖 fox 融合。app_sync_log 保留双写做过渡兼容
(/datamgr/sync-logs 与 LogsPanel 仍依赖),但新语义一律以 app_task_run 为准。
写入侧组件
| 模块 | 职责 |
|---|---|
task_center/run_scope.py | task_run_scope(key, trigger_type, label) 包住一次执行;RunRecorder.step() / task_step_scope 包住子步骤。无论是否抛异常,进入插 running 行、退出回填终态+耗时+计数 |
task_center/context.py | 用 ContextVar(非 threading.local)把当前 run/step id 挂在线程本地,融合器内部取用而不必层层传参。⚠️ ContextVar 不会自动传播到新建线程,故 setter 在 scope 内部调用 |
task_center/log_handler.py | 挂到 root logger,任务执行期间的日志自动关联落 app_task_event——成百上千处 logger.warning 一行都不用改 |
task_center/fox_bridge.py | fox 融合器 ↔ 任务中心的薄胶水:一个 data_type 一次融合 = 一步(带实时进度),并把单只票的失败写进 fox_data_record.success=0 + error_msg(此前两处写入点硬编码 success=1,这张表只能回答"谁提供了数据"、答不了"哪里失败了") |
task_center/store.py | 落账层 |
log_handler 的三个防御性设计(改动务必保留):
- 脱敏:
SensitiveFilter是按 handler 单独挂的,新 handler 必须显式挂,否则 API Key / 密码 / JWT 会明文入库 - 不递归:writer 线程自己产生的日志一律跳过(线程身份识别),否则自己的失败日志进自己的队列死循环
- 不阻塞:业务线程只做
queue.put,落库全在单 writer 线程批量完成(默认只落 WARNING+,队列上限 5000 满则丢最新,单批 200 条 / 2s 刷写,单次执行 500 条事件熔断)
store.py 里所有函数都吞异常(返回值用 None/0/[] 表达失败)、每次写用独立短事务
(融合器长事务里混 INSERT 会放大锁竞争)、DB 不可用时整体降级为 no-op。
连续写失败 3 次进入 60s 冷却,避免 DB 故障时刷屏。
跨进程手动触发:SysConfig 命令队列
web 进程里的 trigger_task_now() 改不了另一个进程的内存态,也就无法去重、无法保证
"点了就真的跑"。改为命令落库 + worker 认领(task_center/commands.py):
POST /datamgr/schedule/{key}/run worker 命令轮询线程(2s)
└─ submit() 写 pending 行 ─────────▶ └─ claim() CAS 认领(pending→running)
└─ 独立线程 run_task(trigger_type="manual")
└─ complete() 写回执 done/failed/rejected
- CAS 认领:
UPDATE ... WHERE config_value = <读到的旧值>,rowcount==1 才算抢到 ——多 worker / 重复轮询并发时只有一个能拿到,天然幂等 - 去重判据在 worker 侧(
is_task_running):任务是否在跑只有调度器进程知道。 运行中则回执rejected,不排队等(手动触发语义是"现在就跑") - 命令键
runtime.sched.cmd.<id>(14 位时间戳 + 6 位随机 hex) - 回收:终态命令保留 24h;超过 24h 仍 pending/running 记为
expired(worker 离线一天了,再执行一个陈旧的手动触发没有意义)
同理,GET /datamgr/schedule/status 的 running / lastRunAt 一律读库
(store.latest_run_by_task()),不读 web 进程内存。
running 行的三重收敛
一条 running 记录必须能自己收尾,否则页面永远显示"进行中":
| 层 | 由谁负责 | 触发时机 |
|---|---|---|
| ① | task_run_scope 的 finally | 正常收尾 |
| ② | store.sweep_orphan_runs() 的两级判活 | 启动期 + 每次 create_run 顺带 |
| ③ | 父 run 已终态却仍 running 的步骤行无条件收敛 | 随 ② |
🔴 ②的判据是两级(fox 侧 sweep_zombie_tasks 同构):
- 执行进程已退出(
app/lib/proc_id.py::owner_alive()返回False:worker_id 可解析、 同主机、pid 已消失)→ 立即收敛为interrupted - 只有判不出死活(
owner_alive()返回None:跨主机 / 旧格式 worker_id / 探测不可用) → 才回落 2h 时长阈值
owner_alive() is True 不得进时长分支——确认存活却按 2h 收敛,会让 fox_kline_sync /
fox_daily_sync 这类常态跑 1.5~2h+ 的全市场链路刚好卡在 2h 处被误标 interrupted
并连带关闭其子步骤行(实测曾在 2h00m11s 处把两个仍在推进的 run 收敛掉)。也不能按
"worker_id != 本进程"判定:多 worker 下别的 worker 是活的。
Windows 下 proc.terminate() 是硬杀,worker 的 finally 不执行(锁随 MySQL 连接断开
自动释放),残留 running 行正是由这条快通道收敛为 interrupted——任务没失败,
是被外部杀掉的。Linux/systemd 的 SIGTERM 才走优雅收敛。
API 契约
读取侧 backend/app/api/tasks.py(前缀 /api/v1/tasks,均需管理员):
| 方法 | 路径 | 说明 |
|---|---|---|
| GET | /tasks/overview | 首屏总览:scheduler(含 leader 与 host)、today、running[](带 steps 与估算 elapsed_ms)、tasks[](各任务最近一次)、failureTop、todayAvgDurationMs、serverTime |
| GET | /tasks/runs | 执行历史分页(task_key / status / trigger_type / from / to / page) |
| GET | /tasks/runs/{run_id} | 单次执行详情 |
| GET | /tasks/runs/{run_id}/events | 该次执行的日志事件 |
| GET | /tasks/events | 跨任务日志事件检索 |
| GET | /tasks/stats/daily | 近 N 天按日聚合(稳定性视图) |
| GET | /tasks/stream | SSE 实时事件流 |
| DELETE | /tasks/runs/purge | 历史清理 |
| GET | /tasks/health | 任务中心自身健康(handler 状态、队列水位、单 run 上限) |
配套的管理侧端点在 datamgr:GET /datamgr/schedule(注册表 + 运行时配置)、
GET /datamgr/schedule/status、PUT /datamgr/schedule/{key}(改开关/触发时间)、
POST /datamgr/schedule/{key}/run(走命令队列,202 受理)、
GET /datamgr/schedule/product-freshness。
前端约定
页面在 src/pages/TaskCenter/(panels/ 五个面板:OverviewPanel / TaskListPanel /
RunHistoryPanel / LogStreamPanel / OpsActionsPanel),API 封装在 src/lib/api.ts 的
tasks 段。
useTradingQuery任务在休市/夜间照跑,useTradingQuery 的 isTradingHours 会在盘后停掉轮询——
总览正是最需要看盘后链路的时候。改用条件轮询(有任务在跑 3s、空闲 30s)
- SSE
/tasks/stream。
SSE 用 fetch 流式读取而非 EventSource:EventSource 无法带 Authorization 头,
而任务中心全部端点需管理员鉴权。
排查清单
| 症状 | 先看 | 判据 |
|---|---|---|
| 徽章显示"调度器已停止" | runtime.worker.heartbeat_at | 90s 内新鲜即活着;host=local/remote/none 三态 |
| 页面卡在"进行中" | app_task_run 的 running 行 + worker_id | owner_alive() 能否解析;进程已退出应立即收敛为 interrupted |
| 昨天某任务整日没跑且无告警 | sched.{key}.last_success_date / last_run_date | 是否被记成运行日导致当日时点被吞(见上文 danger 块) |
| 报 success 但表里 0 行 | app_task_step + fox_data_record | run_fuser 的空交付判定:候选非空却 0 行记 failed + 质量记录 |
| 改了任务代码行为没变 | worker 启动时刻 vs 模块 mtime | worker 不热重载,须重启 run.py;当日产出还要手动重跑覆盖 |
| 起了 worker 看不到"已拿到锁" | 级别 | 该日志在 DEBUG;GET_LOCK 被别的实例持有时本进程转 30s 竞选,属正常 |
相关文档
- 数据流转(定时任务时序)
- 配置指南(环境变量与 SysConfig)
- 数据管理架构
- 常见问题
docs/reports/data-coverage-backfill-2026-10-07.md(数据覆盖与回补报告) — 各任务产出表最新覆盖日期与缺口- 方案原文(工程文档,不进站点):
docs/plans/任务中心建设方案.md