Files
gao daef609e8b feat: QMT Bridge 设为主数据源,新增指数日K支持及数据源文档
- QmtBridgeSource: provides 增加 index_daily,新增 fetch_index_daily 方法
- task_kline_index: 优先级改为 qmt_bridge → mairui → sina 三级降级
- task_kline_5min: 优先用 qmt_bridge,不可用时降级 mairui
- task_kline_daily: 优先级加入 qmt_bridge(首位)
- config: 新增 qmt_bridge_url 配置项
- registry: qmt_bridge 注册信息同步更新
- docs: 新增 DATASETS_AND_SOURCES.md,完整说明数据集与数据源依赖关系
- AGENTS.md: 同步更新指数数据源描述
2026-07-23 23:39:59 +08:00

17 KiB
Raw Permalink Blame History

Agent 工作记忆

活跃监视

task_watch 在后台运行(PID 见 logs/task_watch.log),持续跟踪运行中的 task

  • 查询 dataset_registry 中 status='running' 的任务
  • 解析 config 表 lastMessage 中的耗时数据作为历史基准
  • 自适应检查间隔(剩余时间 × 0.15,clamp 15s-300s,无历史默认 60s

待推进的优化(按优先级)

P0 - 高优

  • task_watch 耗时记录入库:当前 watch 只读历史不做记录。任务完成后应将 dataset_registry.message 中的 elapsed_sec 写入 config 表(或新建 run_history 表),供下次估算 ETA 使用。
  • task_mairui_indicators 外层超时:与 kline_daily/kline_5min 对齐,给 as_completed 加上外层 timeout 兜底(app/tasks/task_mairui_indicators.py:227)。

P1 - 中优

  • 节假日自动更新trading_calendar_holidays 当前硬编码在 config 表。建议接入 https://github.com/shengshiyan/dataconf 或类似日历 API 自动拉取每年交易日历。
  • daily_sync_check 失败告警:昨日巡检因 FAILURE_WEBHOOK_URL/SECRET 未配置而 exit 1(期望的失败)。需确定是否要配 webhook,或改为非致命 warning。
  • kline_5min mairui 偶发全 worker 挂起:今天 recovery 4734 只失败发生在 mairui 单次 120s 全 hang。_fetch 已有 20s timeout + 重试,但仍抗不住 mairui 服务端全挂。建议监控 mairui 健康状态 + 失败股票重入队列。

P2 - 低优

  • MySQL 残留配置.env 仍包含 DB_BACKEND=mysqlMYSQL_* 变量,config.py 已只支持 PG。可清理。
  • stock_node 上次 warning:上次运行 7/10 返回 warning(41/1164 节点失败),需关注是否偶发。
  • kline_5min 大量 fail:今日 recovery 4734 只失败,需单独排查 mairui 数据源。

设计决策记录

  • systemd timer 为主,进程内 scheduler 为辅:宿主机直接部署用 systemd modeDocker 部署用 docker mode。二者通过 RUNTIME_MODE 隔离。
  • After= 只依赖 network-online.targetservice 间不设 After= 依赖链,避免一个任务阻塞后续所有任务。
  • fail-open 策略:节假日不过滤,让 task 在非交易日返回空数据报 warning 但不阻塞。周末由 systemd OnCalendar=Mon..Fri 保护。

数据集来源与依赖关系

数据来源与依赖关系是调度正确性的基础。原始源依赖外部 API;衍生源只读本地 PG,依赖上游原始源。

数据集清单

dataset_id 目标表 类型 来源 说明
stock_basic market_data.stocks 原始源 麦蕊 /hslt/list(主)+ Baostock query_stock_basic(备) 全市场 A 股代码/名称/交易所/上市状态;stock_basic 自身可附带雪球股本快照
kline_daily market_data.kline_stock 原始源 雪球 kline(主)→ 麦蕊 hsstock/history/*/d/n → 新浪 getKLineData(备) 全市场日 K OHLCV;支持 per-stock fallback
kline_index market_data.kline_index + indices 原始源 QMT Bridge(主)→ 麦蕊 hsindex/history → 新浪指数日 K(备) 上证/深证/创业板/沪深300/中证500/中证1000
kline_5min market_data.kline_5min 原始源 麦蕊 hsstock/history/*/5/n 全市场 5 分钟 K;历史深度 2023-06-14 起
tick_trade market_data.tick_trade 原始源 麦蕊 hsrl/zbjy 当天逐笔成交;每日 21:00 后发布
moneyflow market_data.moneyflow 原始源 麦蕊 hsstock/history/transaction 主力/大/中/小单净额;每日 21:30 后发布
longhubang market_data.longhubang_daily + longhubang_seat 原始源 akshare stock_lhb_detail_em / stock_lhb_stock_detail_em(东方财富封装) 龙虎榜聚合层 + 席位层;每日 22:00 触发
stock_node market_data.node_categories + nodes + stock_node_map 原始源 麦蕊 /hszg/{list,gg} 股票-指数/行业/概念映射;每周六 11:30
mairui_indicators market_data.kline_stock_macd_daily + kdj_daily + boll_daily 原始源 麦蕊 hsstock/history/{macd,kdj,boll} 日 K MACD/KDJ/BOLL;增量拉取
share_snapshot market_data.stocks + share 原始源 雪球 quote_detail 最新总股本/流通股本;需 XUEQIU_TOKEN
market_regime market_data.market_regime_daily 衍生源 本地计算(kline_stock 涨跌家数/ advance_ratio / 恐慌标记
mairui_ma_daily market_data.kline_stock_ma_daily 衍生源 本地计算(kline_stock.close MA5/10/20/60;未来可切麦蕊 /d/maN(需付费 licence

依赖关系

stock_basic ─┬─► kline_daily ─┬─► market_regime
             │                └─► mairui_ma_daily
             ├─► kline_5min
             ├─► tick_trade
             ├─► moneyflow
             ├─► longhubang
             ├─► stock_node
             └─► share_snapshot

kline_index (独立,无上游依赖)
  • 根任务stock_basic 是大多数任务的根,提供股票代码列表。kline_index 独立运行。
  • 关键链kline_daily 是 2 个衍生源(market_regimemairui_ma_daily)的实际数据基础。
  • 注意mairui_ma_dailyapp/core/sync/registry.py 中声明依赖 kline_daily,但只有调度顺序约束;若 kline_daily 运行后数据又被回填(如 2026-07-15 场景),mairui_ma_daily 不会自动重算,需手动触发或等下次调度。

默认调度时间(工作日)

时间 dataset_id 备注
09:00 stock_basic 开盘前刷新股票列表
15:30 kline_index 指数日 K
15:40 kline_daily 个股日 K
16:00 kline_5min 5 分钟 K
15:15 market_regime 市场情绪
16:30 mairui_ma_daily 日 K MA(本地计算)
16:40 mairui_indicators MACD/KDJ/BOLL
21:35 moneyflow 资金流(等 21:30 发布)
22:00 longhubang 龙虎榜(等 19:00-21:00 出齐)
周六 11:30 stock_node 节点映射
周六 11:30 share_snapshot 股本快照(雪球,周度)

systemd timer 与数据依赖对齐审核

总体:基本对齐,但存在 3 个需要注意的细节。

  1. mairui_indicators.serviceAfter= 多余依赖 已修复

    • 修复前:After=network-online.target market-sync.service market-sync-mairui-ma-daily.service
    • 修复后:After=network-online.target market-sync.service
    • 文件:deploy/systemd/units/market-sync-mairui-indicators.servicebin/systemd/market-sync-mairui-indicators.service
    • 待部署:/etc/systemd/system/ 仍保留旧版,需 sudo bash deploy/systemd/deploy.sh
  2. stock_node.service 未声明对 stock_basic 的依赖 已修复

    • 修复前:After=network-online.target
    • 修复后:After=network-online.target market-sync-morning.service
    • 文件:deploy/systemd/units/market-sync-stock-node.servicebin/systemd/market-sync-stock-node.service
    • 待部署:/etc/systemd/system/ 仍保留旧版,需 sudo bash deploy/systemd/deploy.sh
  3. After= 只保证启动顺序,不保证上游成功

    • systemd 的 After= 不会让下游等待上游成功完成,仅控制 timer 触发后的排队顺序。
    • 如果 market-sync.service 在 15:30 启动后因 kline_daily 卡住而 3h 超时失败,mairui_ma_daily 仍会在 16:30 启动(因为 timer 已触发)。
    • 结果:mairui_ma_daily 可能基于不完整的 kline_stock 计算,产生类似 2026-07-15 的 3,404 行缺失。

数据回填级联规则

核心原则:上游数据被回填或修复后,下游依赖它的衍生数据集必须手动/自动重跑,否则会出现“上游新、下游旧”的不一致。

上游任务/表 触发条件 必须重跑的下游任务
kline_daily / kline_stock 回填历史 K 线、修复错误日线、补充漏掉的股票/日期 market_regimemairui_ma_daily
stock_basic / stocks 新上市/退市股票、代码变更、上市状态修正 kline_dailykline_5mintick_trademoneyflowlonghubangstock_nodeshare_snapshot
kline_stock(任意修复) close/volume 等核心字段修正 market_regimemairui_ma_daily

操作建议:

  • 单次少量回填(如几只股票、几天):用 CLI 参数指定 codes / start / end,然后按上表手动触发下游。
  • 大量回填(如全市场、多月/多年历史):先跑上游,再按依赖链顺序跑下游;必要时禁用当日独立 timer,避免与 runall 并发。
  • 每日巡检(bin/daily_sync_check.py)已增加数据一致性检查:对比 kline_stockkline_stock_ma_dailymarket_regime_daily 的最新日期与缺失行数,发现缺口即告警。

历史问题与解决思路

记录真实故障案例、排查路径和修复方法,供后续复察或参考。

案例 12026-07-15 kline_stock_ma_daily 缺失 3,404 行

现象

  • kline_stock 11,728,408 行,kline_stock_ma_daily 只有 11,725,004 行,相差 3,404 行。
  • 差异集中在 2026-07-143,400 只股票)+ SH688287 的 4 天(2026-05-29/06-04/06-05/06-08)。
  • dataset_registrymairui_ma_daily 状态为 successlast_success_at=2026-07-14 17:11:48 UTC

根因

  • mairui_ma_daily 在 2026-07-14 17:11 完成时,kline_stock 的 2026-07-14 数据还不完整(少 3,400 只)。
  • kline_daily 在 2026-07-15 10:04 又跑了一次,补充了这 3,400 只 2026-07-14 的 K 线。
  • mairui_ma_daily 没有自动感知上游变化,导致 MA 表缺了这 3,404 行。

排查命令

# 1. 对比行数
SELECT COUNT(*) FROM market_data.kline_stock;
SELECT COUNT(*) FROM market_data.kline_stock_ma_daily;

# 2. 查缺失日期分布
SELECT trade_date, COUNT(*) FROM (
    SELECT stock_code, trade_date FROM market_data.kline_stock
    EXCEPT
    SELECT stock_code, trade_date FROM market_data.kline_stock_ma_daily
) t GROUP BY trade_date ORDER BY trade_date;

# 3. 查 dataset_registry 时间线
SELECT dataset_id, last_success_at, finished_at, message
FROM market_data.dataset_registry
WHERE dataset_id IN ('kline_daily', 'mairui_ma_daily');

修复方法

# 手动重跑 mairui_ma_daily(全量重算 ~22min
bin/market_sync_mairui_ma_daily_run.sh

预防/改进

  • bin/daily_sync_check.py 中增加数据一致性检查:每日 23:00 对比 kline_stockkline_stock_ma_daily 的覆盖差异。
  • 上游回填历史数据后,按「数据回填级联规则」手动触发下游重跑。

案例 2systemd After= 依赖设计缺陷

现象

  • market-sync-mairui-indicators.service 写了 After=market-sync-mairui-ma-daily.service
  • market-sync-stock-node.service 没有声明对 stock_basic 的依赖。

根因

  • 早期设计时对数据依赖关系梳理不够,把无依赖的任务串起来,或遗漏了必要的依赖声明。
  • systemdAfter= 只控制启动顺序,不检查上游任务是否成功完成

修复方法

  1. mairui_indicators.service:移除 market-sync-mairui-ma-daily.service,改为 After=network-online.target market-sync.service
  2. stock_node.service:增加 After=market-sync-morning.service,确保 stock_basic 已更新。
  3. 同步修改 deploy/systemd/units/bin/systemd/ 下的对应文件,并执行:
    sudo bash deploy/systemd/deploy.sh
    

设计原则

  • After= 只声明最小必要依赖:network-online.target + 直接上游 service。
  • 不要为了让下游等上游成功而过度串 service;成功/失败应通过 daily_sync_check 和回填规则兜底。

案例 3dataset_registry 状态卡在 running

现象

  • mairui_ma_dailysync_history 已记录成功,但 dataset_registry.status='running'finished_at=NULL
  • ps 中无对应进程。

根因

  • 状态更新与 sync_history 写入不在同一事务;极端情况下任务进程异常退出,导致 mark_sync_success 未执行。
  • 或并发/重入导致后一次 mark_sync_running 覆盖了前一次的成功状态。

排查命令

SELECT dataset_id, status, started_at, finished_at, last_success_at, message
FROM market_data.dataset_registry
WHERE dataset_id = 'mairui_ma_daily';

SELECT * FROM market_data.sync_history
WHERE dataset_id = 'mairui_ma_daily'
ORDER BY started_at DESC LIMIT 5;

修复方法

from app.core.db import ops as db_ops
from datetime import datetime

db_ops.update_dataset_registry_state(
    'mairui_ma_daily',
    status='success',
    finished_at='2026-07-15 12:02:21',
    last_success_at='2026-07-15 12:02:21',
    message='MA [5, 10, 20, 60] 共 11,728,408 行, 5511 只, 1349.2s',
    last_error=None,
    needs_resync=0,
    progress_current=0,
    progress_total=0,
    current_step='',
    updated_at=datetime.now().strftime('%Y-%m-%d %H:%M:%S'),
)

改进建议

  • 启动 runall_once.py 时自动调用 recover_interrupted_dataset_registry() 清理卡死状态(已实现)。
  • task_watch 可升级为:发现任务 running > 2× 历史平均耗时 且无进程时,自动标记为 failed。

Market Supervisor 工作流

本项目的 Agent 不仅是 code agent,也是 market data supervisor。修改代码/配置后必须完成部署与验证闭环,不能停留在 repo 层面。

Agent 职责

  1. 主动巡检:定期/按需运行 bin/daily_sync_check.py --report-only,确认 dataset 状态、systemd drift、数据一致性。
  2. 问题闭环:发现告警后,定位根因 → 修复代码/config → 部署到 /etc → 验证告警消除。
  3. 项目记忆同步:把已确定的问题、根因、修复方法、验证命令写回 AGENTS.md,供后续复察。
  4. 数据依赖守护:上游数据回填/修复后,按「数据回填级联规则」触发下游重跑,并验证一致性。

systemd unit 修改后的标准流程

任何修改了 deploy/systemd/units/bin/systemd/ 下 unit 文件的操作,必须:

  1. 保持两边同步:同时修改 deploy/systemd/units/bin/systemd/ 的对应文件(目前两者是镜像关系)。
  2. 执行 deploy
    cd /home/gao/Development/quant_home/market_sync
    sudo bash deploy/systemd/deploy.sh
    # 或跳过确认:sudo bash deploy/systemd/deploy.sh --yes
    
  3. 验证 drift
    .venv/bin/python bin/daily_sync_check.py --report-only
    
    确保 systemd_drifthas_drift=false
  4. 验证 unit 状态
    systemctl list-timers market-sync-*
    systemctl status market-sync-mairui-indicators.service
    systemctl status market-sync-stock-node.service
    

部署验证 Checklist

  • repo unit 与 /etc/systemd/system/ 内容一致(daily_sync_check 无 drift 告警)
  • systemctl daemon-reload 已执行
  • 相关 timer 仍 enabledsystemctl is-enabled <name>.timer
  • 修改后的 service 无语法错误(systemd-analyze verify /etc/systemd/system/<name>.service
  • 数据一致性检查通过(kline_stock 与下游表日期对齐)

已部署的优化(确认完成)

项目 日期 说明
RUNTIME_MODE=systemd 写入 .env 07-13 防止双调度器冲突
market-sync-worker.service mask 07-13 systemd 模式下禁用进程内 scheduler
mairui_ma_daily / mairui_indicators timer 部署 07-13 2 个缺失的 timer 补部署并启用
/etc market-sync.service 漂移修复 07-13 deploy.sh 同步 repo→etc,消除 drift
moneyflow/lhb After= 依赖链修复 07-13 去掉 After=market-sync-tick.service,任务不再串行阻塞
bin/task_watch.py 创建 07-13 自适应间隔的任务跟踪监视器(nohup 后台运行)
runall_once.py logger f-string 修复 07-13 超时日志改为 f-string
sync/base.py docstring 更新 07-13 _effective_sync_end 注释过期问题修复
call_with_timeout 超时后线程阻塞修复 07-13 改用 daemon 线程替代 ThreadPoolExecutor,避免超时后 shutdown(wait=True) 阻塞调用方
tick_trade 外层超时加固 07-13 加上 FETCH_HARD_TIMEOUT=120s + 手动 pool 管理,防止 mairui 挂起时永久卡住
kline_daily per-stock fallback 07-13 雪球失败自动降级到 mairui/新浪,fallback_ok 统计备胎救回数
kline_5min 超时后重入队列 07-13 mairui 全 hang 超时后冷却 30s 重试一次失败股票,救回临时故障
MCP Server SSE 模式 + systemd 服务 07-13 market-sync-mcp.service 开机自启,SSE 端点 http://127.0.0.1:8101/sse