From 3eded57d7e9af911da1dd7a04683888d0f78728c Mon Sep 17 00:00:00 2001 From: zlt Date: Tue, 18 Aug 2026 17:31:37 +0800 Subject: [PATCH] =?UTF-8?q?=E5=8A=A0=E5=85=A5=E8=82=A1=E6=B1=87=E5=8F=82?= =?UTF-8?q?=E6=95=B0=E5=92=8C=E7=9B=B8=E5=BA=94=E9=80=BB=E8=BE=91?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- MACRO_TIMING_PLAN.md | 382 +++++++++++++++++++ scripts/calibrate_macro_signal.py | 593 ++++++++++++++++++++++++++++++ 2 files changed, 975 insertions(+) create mode 100644 MACRO_TIMING_PLAN.md create mode 100644 scripts/calibrate_macro_signal.py diff --git a/MACRO_TIMING_PLAN.md b/MACRO_TIMING_PLAN.md new file mode 100644 index 0000000..979d64b --- /dev/null +++ b/MACRO_TIMING_PLAN.md @@ -0,0 +1,382 @@ +# tradingSystem · 宏观择时接入方案 V2(股汇对冲指数 → 组合仓位 + 个股宏观闸) + +> 状态: **V2,按用户 2026-08-18 三条拍板意见修订,待标定后定稿** | 日期: 2026-08-18 +> **V2 修订(依用户三点反馈)**:①触发口径不预设默认值,**先跑标定脚本**(`scripts/calibrate_macro_signal.py`,已交付)用历史数据对比"进区即动 vs 极值回落再动"后再定;②仓位调整由固定步长改为**对数映射**——极值越深加码越多、边际递减,参数由标定给建议值;③新增**个股宏观闸**:偏热区确认后,自主增持类个股动作(新建仓/回踩补足/盈利加仓/补仓)一并暂停——这同时堵上 V1 的一个洞(宏观降仓释放的现金会被 `PMS_OPEN_AUTONOMY=full` 的自主新建仓在下一分钟买回去);④数据源按用户口径:`gp_fx_daily` / `gp_shibor` 在 **192.168.16.153(153 代理)**,标定脚本顺带完成列名探查与核实。 +> 写作纪律沿袭项目既有文档:术语先白话、阈值给硬数字且全部可配置、失败路径显式、未确认事项进清单。 +> 关联:`POSITION_MGMT_DESIGN.md`(总体设计)、`README.md`(模块地图与三条铁律)、`DEVLOG.md` 2026-08-06「动 tradingSystem 以外的系统之前先答三题」。 + +--- + +## 0. 一句话 + +给 PMS 加一个「宏观择时层」:每个交易日算一次股汇对冲指数,指数进入极热区就**像用户一样在命令台下降仓命令**(REDUCE_EXPOSURE,幅度按对数映射随极值深度加码),进入极冷区下升仓命令,中性不动;偏热确认期间同时**暂停自主增持类个股动作**(宏观闸,防止降仓的钱被自动建仓买回去)。它产生仓位变化的唯一出口是**命令**与**闸**;方案生成、规则闸、择时出手、下发、对账、进度回报,全部走现有链路,**一行不改**。框架按「信号注册表」做,未来再接别的宏观信号(bionic 宏观水温、择时层日 QRS 等)只加注册表项。 + +--- + +## 1. 现状盘点:目标需要什么 vs 系统已有什么 + +| 环节 | 这个目标需要 | 系统现状 | 差距 | +|---|---|---|---| +| 信号数据 | `zs_day_data`(上证)、`gp_fx_daily`(USD/CNH)、`gp_shibor`(SHIBOR 1M) | `zs_day_data` 在 199 库有实证(bionic market_regime 在读,列 `symbol/timestamp/close/amount`,代码点式 `000001.SH`);**用户口径:fx 与 shibor 两表在 153 代理**。两表在五个项目里零代码引用,列名未经实测 | **小**:连接现成(153=proxy 源、199=index 源都已配),列名由标定脚本探查后回填 §4.2 | +| 信号计算 | 20 日对数收益 → spread → SHIBOR 修正 → 40 日 Z-Score ×10 | 全生态无任何股汇对冲指数实现(`hedge` 关键词全项目零命中)。bionic 有 `market_regime.py`(宏观水温),但那是决策系统内部件,PMS 不读 | **缺**:需新写,约百行纯函数 | +| 组合级自动决策 | 极值 → 升降整体仓位 | **无**。设计 §5 原话:「总仓上限由用户命令控制(升降仓判断权在用户,**不自动联动**)」。本需求是对这条设计决策的一次**有意扩展**,必须带着同样的档位/确认哲学落地 | **缺**:决策引擎(区制/确认/冷却/对数映射/上下限) | +| 个股动作联动 | 偏热时管住增量(用户拍板新增) | 自主提议五类动作只看个股与组合约束,不看宏观;`组合刹车`有"回撤后停自主增持"的同型机制可借鉴 | **缺**:宏观闸(§5.2,实现手势照抄刹车) | +| 执行面 | 把「降 X%」翻成逐股卖出并执行 | ✅ **已完备**。`REDUCE_EXPOSURE` / `INCREASE_EXPOSURE` 任务命令全生命周期:`command_service.issue()` → planner 四档凑额(停新买→清弱票→收利润→等比微减)→ 方案落表 → executor 转指令 → 规则闸终检 → 择时出手 → 回报进度。`pms_command.issued_by` 字段与 `/api/commands` 的 `by` 参数都在 | **无差距,全量复用** | +| 冲突仲裁 | 自动命令不得顶撞用户在途命令 | ✅ `command_spec.detect_conflicts` 已覆盖 REDUCE↔INCREASE、LIQUIDATE、HALT 等对 | 复用,宏观侧遇冲突**跳过不 force** | +| 档位/确认哲学 | 先建议后自动,人可随时收权 | ✅ 有成熟模板(`PMS_AUTONOMY` / `PMS_OPEN_AUTONOMY` 三档)。但 `pms_proposal` 是**个股级**,组合级建议放不进去 | **小**:宏观建议自己承载(§7) | +| 参数 | 阈值全部页面可调 | ✅ `param_store` 登记即出现在设置页、5 秒生效 | 登记即可 | +| 调度 | 每交易日定时算一次 | ✅ celery beat 骨架 + 三条守卫(交易日/休假/吞异常守成)成熟 | 加一个调度位 | +| 页面 | 看得见指数、区制、建议、闸状态 | ✅ 单页已有多面板先例 | 加一块面板 | +| 留痕 | 为什么动/为什么没动,可回查 | ✅ `pms_action_ledger` + 命令表 | 复用 + 新增信号历史表 | +| 历史存储 | 指数序列可回看、可判分、可标定 | 无 | **缺**:新表 `pms_macro_signal` | + +结论:**执行面九成现成,缺的是"信号计算 + 组合级决策 + 宏观闸 + 一小块页面"这一层薄壳**。风险不在工作量,在接入姿势——所以本方案把大部分篇幅花在护栏上。 + +--- + +## 2. 目标翻译成系统语言:为什么走命令通道 + 扫描层闸 + +「极热 → 降低整体仓位」在本系统里最忠实的表达就是一条 `REDUCE_EXPOSURE {pct, window_tdays}`——语义完全一致(释放金额 = 规模 × pct,planner 按 停新买→清弱票→收利润→等比微减 四档凑额,这正是"pms 根据规则执行对个股的降仓位操作")。反向就是 `INCREASE_EXPOSURE`。存量靠命令,增量靠闸:偏热期间把自主增持的**扫描**停掉(不产生候选就不会有指令),两头合起来才是完整的"仓位受宏观控制"。 + +刻意**不做**的两种姿势: + +1. **不做"连续仓位控制器"**(宏观层直接持有目标仓位、每分钟收敛)——绕开命令系统另立产生买卖的路径,与命令冲突时没有仲裁者,正是禁区。 +2. **不把宏观观点塞进规则闸**——规则闸是合规终检,塞行情观点会让"没通过合规"和"宏观不看好"在评审账本里分不开;宏观闸放在**扫描入口**(proposal_service),与"组合刹车""休假模式"同一落点、同一手势。 + +走命令通道的三个直接好处:宏观动作在**命令台可见、可撤销、可改窗口**(命令至上铁律天然成立);与用户命令的冲突由现有 `detect_conflicts` 仲裁;执行成本统计、进度日结、判分全部免费获得。 + +--- + +## 3. 总体架构与六条设计原则 + +``` + 每交易日 09:35 (beat: macro_scan) 页面「宏观择时」面板 + │ │ 手动扫描 / 一键采纳 + ▼ ▼ + ┌──────────────────────────────────────────────────────────┐ + │ macro_service.scan() │ + │ 取数(macro_repo, 三张源表单表查) → 计算(macro_rules 纯函数) │ + │ → 落 pms_macro_signal (当日一行, 重扫就地更新) │ + │ → 决策(macro_rules.decide: 区制/确认/冷却/对数映射/上下限) │ + │ → 分流: │ + │ autonomy=full → command_service.issue( │ + │ REDUCE/INCREASE_EXPOSURE, │ + │ issued_by="macro") │ + │ autonomy=propose_only → 当日建议落表, 面板等用户点头 │ + │ autonomy=off / ENABLED=False / 信号UNAVAILABLE → 不动 │ + │ → 每一次动作与每一次"想动被拦"落 pms_action_ledger │ + └──────────────────────────────────────────────────────────┘ + │ (命令一旦下达, 与用户手下的命令零区别) │ 宏观闸状态 (偏热确认中?) + ▼ ▼ + command_poll → planner → executor proposal_service 扫描入口: + → 规则闸 → 择时 → 下发 → 对账 偏热确认 → 增持类扫描暂停(留痕) +``` + +六条设计原则(每条都能在代码评审时逐一核对): + +1. **默认关**:`PMS_MACRO_ENABLED=False`。部署完成后系统行为与今天完全一致,开关和档位都在页面上,随时收权。 +2. **唯一出口**:宏观层产生仓位变化的途径只有两条——`command_service.issue()` 下命令、扫描入口的宏观闸(只拦不发)。不直接写 `pms_plan` / `pms_instruction` / `pms_position`,不调 executor,不碰账本。 +3. **故障即守成**:任一数据源取不到、样本不足、日龄超限 → 信号置 `UNAVAILABLE`,**不动作、不落闸**(拿不到 ≠ 偏热,闸的安全方向是"不额外拦"——增持本身另有规则闸把关),面板亮黄条。挂在调度器 `guard` 下,异常吞掉只记 ERROR。 +4. **命令至上**:与用户在途命令冲突(issue 返回 CONFLICT)→ 本轮放弃并留痕,**绝不 force**;`HALT_ALL`(休假)由调度守卫直接跳过;`PMS_GLOBAL_BUY_HALT` 生效时不发升仓;组合刹车期间默认不发升仓(可参数放开)。**用户命令永不受宏观闸限制**,且随时可在命令台撤销宏观命令。 +5. **限频防抖**:进区确认天数 + 迟滞退出带 + 同方向冷却 + 每信号每交易日至多一次动作 + 对数映射的"目标追踪"(§5.1,同一深度不重复加码)+ 上一条宏观命令在途不叠加。指数在极值区横跳不会造成反复买卖。 +6. **全留痕**:信号值每日入 `pms_macro_signal`(含中间量,可复算);每次触发/被拦/落闸入 `pms_action_ledger`(action=MACRO,hard_numbers 带当时指数值、区制、超额深度、当前仓位、本次步长);命令表 `issued_by='macro'` 一眼可辨。 + +--- + +## 4. 信号层设计 + +### 4.1 计算口径(按定义逐条落地) + +数据统一为按国内交易日对齐的日频序列,全部用**收盘值**: + +``` +① stock_ret = ln(close_t / close_{t-20}) # zs_day_data, symbol='000001.SH' + fx_ret = ln(fx_t / fx_{t-20}) # gp_fx_daily USD/CNH, 正值=人民币贬值 + (fx 的日期先 +1 自然日, 再对齐到国内交易日: 当日无值取最近前值, 即 as-of 对齐) +② spread = stock_ret + fx_ret +③ spread_adj = spread − β × (SHIBOR1M_t − SHIBOR1M_{t-20}) # β 默认 0.02 +④ hedge_index = (spread_adj − mean40(spread_adj)) / std40(spread_adj) × 10 +``` + +硬性口径,写进单测:窗口内任一序列缺口用 as-of 前值补齐但**记录补齐天数**;有效样本 < 20+40+5 个交易日 → `UNAVAILABLE`;`std40 = 0` → `UNAVAILABLE`;三源中最新日期落后当前交易日超过 `PMS_MACRO_STALE_TDAYS`(默认 3)个交易日 → `UNAVAILABLE`。 + +区制(zone)判定带迟滞,防横跳(阈值初值 ±25/±15,**以标定结果为准回填**): + +| zone | 进入条件 | 退出条件 | +|---|---|---| +| HOT(偏热,即 positive 区) | hedge_index > `PMS_MACRO_HOT_TH` | 回落到 < `PMS_MACRO_EXIT_BAND` | +| COLD(偏冷,negative 区) | hedge_index < `PMS_MACRO_COLD_TH` | 回升到 > −`PMS_MACRO_EXIT_BAND` | +| NEUTRAL | 其余 | — | +| UNAVAILABLE | 数据守卫未过 | 数据恢复 | + +### 4.2 数据访问(严格单表,三条独立查询) + +新增 `app/repo/macro_repo.py`。**用户口径:`gp_fx_daily` / `gp_shibor` 在 153 代理** → 走 `db.session.fetch_all(..., source="proxy")`;`zs_day_data` 已证实在 199(`source="index"`),若标定脚本探出 153 也可达则统一走 proxy(少一个依赖库)。每张表一条单表 SQL(`WHERE ... ORDER BY 日期 DESC LIMIT 120` 量级),完全符合单表守卫。**列名以标定脚本探查结果为准回填此处,回填前不写 macro_repo**: + +| 表 | 源 | 日期列 | 取值列 | 过滤条件 | +|---|---|---|---|---| +| zs_day_data | index(153 待探) | timestamp | close | symbol='000001.SH' | +| gp_fx_daily | proxy(153) | 待探 | 待探 | 待探(USD/CNH 筛选) | +| gp_shibor | proxy(153) | 待探 | 待探(1M) | 待探 | + +`pms_macro_signal` 的读写走 `source="proxy"`(与其余 pms_* 表同侧)。 + +### 4.3 信号注册表(多信号扩展点) + +`app/services/macro_service.py` 内一张注册表,V1 只有一项: + +```python +SIGNALS = { + "hedge_fx": { # 股汇对冲指数 + "label": "股汇对冲指数", + "fetch": macro_repo.fetch_hedge_inputs, # -> 原始序列 + "compute": macro_rules.compute_hedge_index, # -> {value, zone, detail} + }, + # 未来: "bionic_regime" (决策系统宏观水温) / "index_qrs" (择时层日QRS) —— + # 只加表项与各自参数(命名空间 PMS_MACRO__*), 决策层/执行层/闸/页面零改动。 +} +``` + +启用哪些由 `PMS_MACRO_SIGNALS`(逗号分隔,默认 `hedge_fx`)控制。多信号合成 V1 定死**最保守者优先**:任一启用信号 HOT → 整体偏热;无 HOT 而有 COLD → 偏冷;否则中性。加权合成属二期,不做预设计。 + +### 4.4 新表 `pms_macro_signal`(第 18 张表) + +```sql +CREATE TABLE IF NOT EXISTS pms_macro_signal ( + id BIGINT PRIMARY KEY AUTO_INCREMENT, + signal_key VARCHAR(32) NOT NULL COMMENT '信号名, 如 hedge_fx', + trade_date INT NOT NULL COMMENT '北京交易日 YYYYMMDD', + value DECIMAL(12,4) COMMENT 'hedge_index 值', + zone VARCHAR(16) NOT NULL DEFAULT 'NEUTRAL' COMMENT 'HOT/COLD/NEUTRAL/UNAVAILABLE', + detail_json TEXT COMMENT '中间量: stock_ret/fx_ret/spread/spread_adj/超额深度/各源末日/补齐天数/周期已执行调整', + action VARCHAR(32) NOT NULL DEFAULT 'NONE' COMMENT 'NONE/ADVICE_REDUCE/ADVICE_INCREASE/CMD_ISSUED/BLOCKED/GATED', + ref_id VARCHAR(64) NOT NULL DEFAULT '' COMMENT '关联 command_id (下达或采纳后回填)', + note VARCHAR(500) NOT NULL DEFAULT '' COMMENT '被拦原因/建议文案/闸状态', + updated_at DATETIME COMMENT '写入或重扫更新时间', + UNIQUE KEY uk_sig_date (signal_key, trade_date) +) COMMENT='宏观信号日快照与动作记录 (每信号每交易日一行, 重扫就地更新)'; +``` + +一行四用:历史序列(面板画近 20 日)、当日建议的承载(propose 档)、冷却与"周期已执行调整"的事实源、宏观闸的状态源(扫描层 5 秒缓存读当日行)。 + +--- + +## 5. 决策层设计 + +### 5.1 组合仓位:对数映射(V2 改,替代固定步长) + +极值越深,调整越多,但**边际递减**——既回应"越极端越该动",又天然防止 Z-Score 尾部把仓位一次打爆。纯函数 `macro_rules.decide(...)`,规则按优先级: + +1. `PMS_MACRO_ENABLED=False` 或档位 off 或 zone ∈ {NEUTRAL, UNAVAILABLE} → NONE。 +2. 确认:连续处于同一极值区 ≥ `PMS_MACRO_CONFIRM_DAYS`;触发口径 `PMS_MACRO_TRIGGER_MODE`(zone_enter / zone_exit,**默认值由标定结果定**,两种都实现、页面可切)。 +3. **对数映射(核心)**: + - 超额深度 `e = |hedge_index| − 进入阈值`(e ≥ 0,单位=指数点); + - 本轮极值周期的**累计目标调整** `target(e) = min( S0 × ln(1 + e / k), PMS_MACRO_SHIFT_MAX )`(占规模的百分点); + - 本次步长 `step = target(e_now) − done_shift`(done_shift = 周期内已执行的宏观调整,记在当日信号行 detail 里);`step < PMS_MACRO_STEP_MIN`(默认 0.02)→ 不动; + - **周期**定义:进区 → 迟滞退出;退出后 done_shift 清零,下个周期重新累计。同一深度不重复加码,指数横盘在区内不会连环下命令;指数**更深**才有新步; + - `S0` / `k` / `SHIFT_MAX` 初值由标定脚本给建议(目标形状:e 中位数 → 约 5 个点,e 95 分位 → 约 12 个点,封顶 20 个点),页面可调。 +4. 冷却:同方向上一次**生效**动作距今 < `PMS_MACRO_COOLDOWN_TDAYS`(默认 3 个交易日)→ 本轮不动(下轮 target 追踪仍在,冷却只是把步子隔开)。「生效」= 所发命令未被立刻 CANCELLED(零方案的升仓命令不占冷却,次日重试)。 +5. 在途:上一条宏观命令仍在 EXECUTING/PARTIAL → 不叠加。 +6. 方向裁剪: + - HOT → 降仓,`step` 再受地板裁剪:不把总仓位降到 `PMS_MACRO_MIN_PCT`(默认 0.20)以下(用户命令不受此限); + - COLD → 升仓,`step` 受 `PMS_MACRO_MAX_PCT`(默认 0=由 `PMS_PORTFOLIO_CAP` 兜底)裁剪; + - 升仓独有护栏:`PMS_GLOBAL_BUY_HALT` 生效不发;组合刹车生效且 `PMS_MACRO_RESPECT_BRAKE=True`(默认,用户未反对,标注为**假设**)不发。降仓不受这两条限制(与"冻结与刹车只挡增持不挡减持"同口径)。 + +### 5.2 个股宏观闸(V2 新增,用户拍板"个股交易动作受宏观信号控制") + +**语义**:偏热(HOT)确认期间,PMS 自己发起的**增持类**个股动作暂停——自主新建仓(OPEN)、回踩补足(FILL)、盈利加仓(ADD)、补仓(DCA)的扫描整体跳过并留痕;**减持类动作(TRIM/风控卖出/清仓)与用户命令永不受限**。COLD 期间不加额外限制(自主动作照常,加量由升仓命令负责)。这与"组合刹车"的既有语义完全同型("自主增持停,命令类不受限"),实现手势照抄。 + +**为什么必须有**:不加闸的话存在一条自我打架的路径——宏观降仓卖出 → 现金变多 → `PMS_OPEN_AUTONOMY=full` 的自主新建仓下一分钟把钱又买回去。闸把增量管住,命令管存量,才闭环。 + +**实现落点(两处,均为纯增量,与既有 `exec_halt` / 刹车检查并排)**: +- `proposal_service.scan_and_route` 入口:宏观闸生效 → 四类已有持仓动作中**买入侧**候选不扫(TRIM 照常),skipped 留痕"宏观偏热闸"; +- `proposal_service._scan_open` 入口:宏观闸生效 → 本轮不建仓,skipped 留痕。 + +闸状态由 `macro_service.gate_state()` 提供(读当日 `pms_macro_signal`,进程内 5 秒缓存,**读不到 = 不落闸**——闸的安全方向是"不额外拦",增持自有规则闸把关,不能让宏观层故障把整个自主引擎摁死)。开关 `PMS_MACRO_STOCK_GATE`(默认 True,随 `PMS_MACRO_ENABLED` 总开关生效;即便档位是 propose_only,闸也生效——"停止买入"是保守方向动作,与"减持方向不设确认门槛"同哲学)。 + +**边界(待拍板 §12.3)**:个股交易方案(网格/做T/跟踪)的买入侧要不要一并受闸。V2 建议**先不动策略层**——策略是用户特意设的(signal_service 对挂策略票的处理哲学是"只提示不替你做"),宏观闸自动暂停它们容易造成困惑;若你要管,复用现成的 `strategy_service.pause_buy`(页面可恢复)一行接入即可。 + +--- + +## 6. 执行接入(兼容性核心节) + +### 6.1 下达即普通命令 + +```python +res = command_service.issue( + "REDUCE_EXPOSURE", # 或 INCREASE_EXPOSURE + {"pct": step, "window_tdays": param_store.get_int("PMS_MACRO_WINDOW_TDAYS", 3)}, + issued_by="macro", + note="宏观择时: hedge_fx=+31.2 偏热确认(阈值+25, 超额6.2), 周期目标降8.6点已降0点, 本次降8.6点") +``` + +- `issue()` 内部自动做冲突检测:返回 `ok=False` 且 errors 含 CONFLICT → 宏观层**留痕跳过**(ledger verdict=REJECT),当日不再重试。 +- 命令落表后与用户命令完全同路:`command_poll` 每分钟排方案、`intraday_exec` 择时出手、`refresh_progress` 日结进度。`shadow` 通道下只记账不真发(与现状一致,宏观层无需感知通道模式)。 + +### 6.2 调度时点:09:35,不是盘前 + +一个已核实的坑决定了时点:`positions_view` 在取不到实时价时用摊薄成本顶住价格并标 `price_ok=False`,而 `planner._usable` **排除** `price_ok=False` 的票(2026-07-31 静默失败专项第九条)。行情 db13 只有当日分钟线,盘前全部无价——盘前下达的 REDUCE 命令会在一分钟内被排成**零方案并置 CANCELLED**。所以 `macro_scan` 定在 **09:35**(开盘后行情已就位,也在 08:40 拉候选池之后,升仓有票可选)。信号用的是 T-1 日终数据,几点算值都一样。beat 表加一行: + +```python +"macro_scan": {"task": "pms.macro_scan", "schedule": crontab(hour=9, minute=35)}, +``` + +任务本体 `@guard(trade_day=True, respect_exec_halt=True)`——休假模式自动跳过,异常吞掉守成。宏观闸的状态在 09:35 扫描落表后当日恒定(信号是日频慢变量),盘中每一跳扫描层读的是同一行。另设手动端点(§7.3)供盘中任意时刻重扫/试算。 + +### 6.3 护栏对照表(每条对应代码里一处显式检查) + +| 场景 | 行为 | 靠什么保证 | +|---|---|---| +| 休假模式 HALT_ALL | 整跳跳过 | 调度 guard `respect_exec_halt`;手动端点内 `macro_service` 自查一遍 | +| 用户在途 REDUCE/INCREASE/LIQUIDATE 冲突 | 跳过 + 留痕,不 force | `issue()` 现有 `detect_conflicts`,宏观层永不传 `force_conflict=True` | +| 全局暂停买入 | 升仓不发 | 引擎显式读 `PMS_GLOBAL_BUY_HALT`(此开关在参数表而非命令表,冲突检测看不见它,必须自查) | +| 组合刹车 | 升仓默认不发 | 引擎读 `PMS_BRAKE_UNTIL` 比对当日 | +| 数据缺失/停更 | 不动作、不落闸 + 面板黄条 | 信号 UNAVAILABLE 分支 | +| 极值区横跳/横盘 | 不反复动 | 迟滞退出带 + 对数目标追踪(同深度不重复加码)+ 冷却 + 每日一次 | +| 上一条宏观命令未走完 | 不叠加 | 查 `issued_by='macro'` 的在途命令 | +| 参数表读不到 | 宏观层整体失效(安全方向) | `PMS_MACRO_ENABLED` 文件初值 False,表挂了退初值即关;闸读不到状态 = 不拦 | +| 仓位已到地板/天花板 | 不动 + 留痕 | decide() 映射裁剪 | +| 宏观闸误伤减持 | 不可能 | 闸只挡买入侧扫描,TRIM/EXIT/风控卖出路径不经过闸 | + +--- + +## 7. 档位、建议承载与页面 + +### 7.1 三档(沿用系统词汇) + +`PMS_MACRO_AUTONOMY`: `off` / `propose_only`(**默认**)/ `full`。 +- off:只算信号、只落历史、只展示,不出建议不下命令(**闸也不生效**——off 是"只看不动"档)。 +- propose_only:触发时把建议写进当日 `pms_macro_signal` 行(action=ADVICE_*,note 带完整理由与建议 pct),面板出卡片等一键采纳;**当日有效**,隔日自动作废。**宏观闸照常生效**(保守方向自动,见 §5.2)。 +- full:直接 `issue(issued_by="macro")`。三条铁律下用户仍可在命令台撤。 + +不复用 `pms_proposal`:它的采纳路径是"生成该股指令",是个股级语义;硬塞组合级建议要改共用 decide 端点。宏观建议放自己的表行 + 自己的采纳端点,互不沾染。 + +### 7.2 采纳路径 + +`POST /api/macro/adopt {signal_key}` → 校验建议仍是当日且未采纳 → 重跑一遍 decide 的护栏(冲突/开关/地板可能在建议挂出后变化)→ `command_service.issue(..., issued_by="user")`(用户点的头就是用户意志)→ 回填 `ref_id=command_id, action=CMD_ISSUED`。幂等:同一行只能采纳一次。 + +### 7.3 API 与页面 + +- `GET /api/macro/status`:各信号 {value, zone, 超额深度, 近 20 日序列, 数据日龄, 当日建议, 闸状态, 最近宏观命令及进度, 冷却剩余, 周期已执行调整}。 +- `POST /api/ops/macro-scan?dry_run=true|false`:手动扫描;dry_run 只算不落表不下达。 +- `POST /api/macro/adopt`:见上。 +- 页面 `index.html` 加「宏观择时」面板:指数当前值 + 区制色带(热红/冷蓝/中性灰/不可用黄)+ 近 20 日迷你走势 + 当日建议卡(一键采纳/忽略)+ **宏观闸状态条**("偏热闸生效中:自主增持已暂停")+ 最近宏观命令进度。参数不用做界面——`PMS_MACRO_*` 登记后自动出现在「参数设置」。 +- Makefile 顺手加 `t-macro`(status + 试算),非必需。 + +--- + +## 8. 参数清单(settings.py 初值;全部经 param_store 页面可调;标 ⚙ 的初值由标定回填) + +| 参数 | 默认 | 说明 | +|---|---|---| +| PMS_MACRO_ENABLED | False | 宏观择时总开关(关=调度位空转,不取数不计算,面板显示"未启用") | +| PMS_MACRO_AUTONOMY | propose_only | off / propose_only / full | +| PMS_MACRO_SIGNALS | hedge_fx | 启用的信号清单(逗号分隔) | +| PMS_MACRO_STOCK_GATE | True | 个股宏观闸:偏热确认期间暂停自主增持类扫描(off 档整体不生效) | +| PMS_MACRO_HOT_TH ⚙ | 25.0 | 偏热进入阈值 | +| PMS_MACRO_COLD_TH ⚙ | -25.0 | 偏冷进入阈值 | +| PMS_MACRO_EXIT_BAND ⚙ | 15.0 | 迟滞退出带(\|值\| 回落到此内算离区) | +| PMS_MACRO_CONFIRM_DAYS ⚙ | 1 | 进区连续 N 个交易日才触发 | +| PMS_MACRO_TRIGGER_MODE ⚙ | zone_enter | zone_enter=进区即动 / zone_exit=极值回落再动(**默认值待标定对比后定**) | +| PMS_MACRO_LOG_S0 ⚙ | 0.12 | 对数映射系数:target = S0·ln(1+e/k) | +| PMS_MACRO_LOG_K ⚙ | 10.0 | 对数映射尺度 k | +| PMS_MACRO_SHIFT_MAX ⚙ | 0.20 | 单个极值周期的累计调整封顶(占规模的百分点) | +| PMS_MACRO_STEP_MIN | 0.02 | 裁剪后步长小于此不动 | +| PMS_MACRO_COOLDOWN_TDAYS | 3 | 同方向两次动作最小间隔(交易日) | +| PMS_MACRO_MIN_PCT | 0.20 | 宏观降仓地板(自动降仓不把总仓位降到此下;用户命令不受限) | +| PMS_MACRO_MAX_PCT | 0.0 | 宏观升仓天花板;0=不单设,由 PMS_PORTFOLIO_CAP 兜底 | +| PMS_MACRO_WINDOW_TDAYS | 3 | 宏观命令的执行窗口 | +| PMS_MACRO_RESPECT_BRAKE | True | 组合刹车期间不自动升仓(假设项,见 §12.4) | +| PMS_MACRO_STALE_TDAYS | 3 | 任一数据源日龄超此(交易日)→ UNAVAILABLE | +| PMS_MACRO_RET_WIN / Z_WIN | 20 / 40 | 收益窗 / Z-Score 窗 | +| PMS_MACRO_SHIBOR_BETA | 0.02 | SHIBOR 修正系数 β | + +`_RANGES` 同步登记,`DESC` 写中文说明。**不进 FAIL_CLOSED**:ENABLED 的文件初值就是 False,表读不到退初值即安全方向。 + +--- + +## 9. 触碰面清单与三条铁律自证 + +新增文件(5 个,互相独立,删净即回到今天——"三题"规矩第 3 题): + +| 文件 | 内容 | +|---|---| +| app/core/macro_rules.py | 纯逻辑:对齐与四步计算、zone 迟滞判定、对数映射 decide(无 IO,可单测) | +| app/repo/macro_repo.py | 三张源表单表取数 + pms_macro_signal 读写 | +| app/services/macro_service.py | 编排:守卫→取数→算→落表→决策→分流→留痕;status / adopt / gate_state | +| scripts/calibrate_macro_signal.py | 标定脚本(**已交付**,探查+分布+口径对比+对数参数,只读) | +| scripts/test_batch14_units.py | 单测约 25 例(§11) | + +触碰的既有文件(9 处,**全部纯增量**——只加行/加检查,不改既有行为,逐处可 diff 核对): + +| 文件 | 加什么 | 不动什么 | +|---|---|---| +| config/settings.py | PMS_MACRO_* 字段一段 | 既有字段零改动 | +| app/services/param_store.py | DESC / _RANGES / _range_check 各加宏观项 | 缓存、FAIL_CLOSED、既有键 | +| app/scheduler.py | `macro_scan` 任务 + beat 一行 | 既有九个调度位 | +| app/services/proposal_service.py | 两处入口各加一个宏观闸检查(与 exec_halt 检查并排,闸关/读不到=原行为) | 扫描、闸门、分流逻辑零改动 | +| app/web/main.py | 三个新端点 | 既有端点 | +| app/web/static/index.html | 一块面板 | 既有面板 | +| ddl_pms_v1.sql | 第 18 张表 CREATE IF NOT EXISTS | 既有表 | +| scripts/check_db.py | PMS_TABLES 加 pms_macro_signal | (顺带发现:pms_strategy / pms_op_log 也不在这份清单里,属既有欠账,按"不夹带"纪律单独提,不在本次一起改) | +| scripts/run_tests.py | SUITES 登记 batch14 | 既有批次 | + +三条铁律逐条自证:**命令至上**——宏观动作本身就是命令,冲突让位用户,可撤可改;宏观闸只拦 PMS 自主动作,用户命令与减持全不受限。**分工不越权**——宏观层只回答"整体该多重、增量该不该停",逐股怎么卖买仍由 planner/择时/规则闸决定,不碰决策系统(bionic/桥/择时层**零改动**,"三题"规矩天然满足)。**先记账后动作 + 故障即守成**——信号先落表再决策,命令走 issue 的先落表路径,任何取数/计算失败 = 不产生新指令、不落闸。 + +--- + +## 10. 前置标定(先于一切编码,脚本已交付) + +`scripts/calibrate_macro_signal.py`(**严格只读**,只发 SELECT/SHOW)在桥机跑,一次回答四件事: + +1. **列名探查**:对 153(proxy)与 199(index)两源探查三张表的位置、列名、样本行——结果回填 §4.2 后才写 macro_repo; +2. **分布体检**:hedge_index 历史分位数、±20/25/30 各阈值的触发天数与年频次、极值区段的持续天数与最大深度——回答"±25 合不合身"; +3. **触发口径对比**:zone_enter vs zone_exit 各自触发事件后 5/10/20 交易日上证的均值/中位数/胜率(HOT 后应偏负、COLD 后应偏正才算有效)——**注意样本量会很小,结论当参考不当真理**; +4. **对数映射标定**:极值周期超额深度 e 的分布 → 解出 S0/k 建议值(e 中位数→5 点、95 分位→12 点、封顶 20 点),并打 e→调整幅度对照表。 + +**运行方式(不需要重建镜像**——脚本从宿主工作树经 stdin 喂给容器 python,`python -` 的 sys.path[0] 是工作目录 /app,`config.settings` 照常可导入): + +```bash +cd ~/tradingSystem # git pull 之后 +docker compose run --rm -T pms-web python - < scripts/calibrate_macro_signal.py > /tmp/macro_calib_report.md +cat /tmp/macro_calib_report.md +# 需要完整序列时: +docker compose run --rm -T pms-web python - --dump-csv < scripts/calibrate_macro_signal.py > /tmp/hedge_series.csv +``` + +探查失败的退路:脚本会打出该表全部列名与样本行,按提示带 `--fx-table/--fx-date-col/--fx-price-col/--fx-where/--shibor-*` 覆盖重跑;表真不存在 → 数据先落地(数据管道补),或立项由桥/择时层供数(跨系统改动,按"三题"规矩单独议,本方案不带)。 + +--- + +## 11. 实施步骤与判收 + +| 步 | 内容 | 判收标准 | +|---|---|---| +| 1 | 桥机跑标定脚本(§10,可立即做) | 三表列名落实;报告出分布/口径对比/S0-k 建议;把报告拿回来定 §8 带 ⚙ 的初值 | +| 2 | 写码 + batch14 单测 | `make test` ALL SUITES PASS。用例:合成序列对照手算值验四步计算;fx +1 对齐与缺口 as-of;样本不足/停更→UNAVAILABLE;zone 迟滞进出与横跳;对数映射(target 单调、同深度不重复加码、周期清零、封顶);confirm/cooldown/每日一次;零方案命令不占冷却;冲突/HALT_BUY/刹车跳过(打桩);宏观闸(HOT 拦买入侧不拦 TRIM、读不到不拦、off 档不拦);建议采纳幂等与隔日作废 | +| 3 | ddl 建表 + `make deploy`(收盘后) | `make check` 见 18 表;页面参数区出现 PMS_MACRO_*,ENABLED=False,系统行为与部署前一致 | +| 4 | 手动 `POST /api/ops/macro-scan?dry_run=true` | 面板出指数值与区制;与标定脚本同日值一致(同一套纯函数,天然一致) | +| 5 | 开 ENABLED=True,档位 propose_only,观察 ≥1 周 | 每日一行信号历史;极值日出建议卡不出命令;宏观闸在 HOT 日真的把自主增持 skipped(留痕可查);采纳一次走通 建议→命令→方案→指令 全链 | +| 6 | 冲突演练:挂一条在途 INCREASE,手动扫描触发降仓 | 账本见 MACRO REJECT 留痕,无命令产生 | +| 7 | (拍板后)切 full | 首条自动命令端到端回执,日报关注区可见 | + +回滚:页面把 `PMS_MACRO_ENABLED` 置 False 即回到现状(秒级,命令、闸、建议全停);代码级回滚 = 还原上表 9 处触碰 + 删 5 个新文件;新表留着无任何读方,无副作用。 + +--- + +## 12. 待拍板清单(V2 更新) + +1. ~~触发口径默认值~~ → **改为由标定结果定**:跑完 §10 脚本,拿两种口径的事后对比数据来选(样本小的话建议 zone_enter + CONFIRM_DAYS 取 1~2 里标定更稳的那档)。 +2. ~~固定步长 vs 目标带~~ → **已定:对数映射**(用户拍板)。S0/k/SHIFT_MAX 初值等标定回填。 +3. **策略层(网格/做T)的买入侧要不要受宏观闸**:V2 建议先不动(理由 §5.2 末),要管的话复用 `strategy_service.pause_buy` 一行接入。 +4. **升仓是否受组合刹车约束**:上轮问题未获直接答复,V2 按**受约束(True)**落地并标注为假设;要放开改参数即可。 +5. **数据源**:已答——153。zs_day_data 是否也统一走 153 由探查定。 + +## 13. 风险与已知限制(说在前面) + +- **信号本身是慢变量**:20 日收益 + 40 日 Z,天然滞后,抓的是阶段冷热不是日内顶底;对数映射 + 冷却是为此设的缓冲。 +- **z-score 的尾部行为**:40 日窗里刚经历过极端行情时 std 变大,新极端可能"钝化";反之平静期小波动被放大。标定报告的分位数表就是用来看这个的。 +- **事后对比样本小**:±25 级别的极值一年可能只有几次,标定里的口径对比是参考不是显著性检验,别过度拟合。 +- **升仓依赖候选池**:极冷区触发升仓但当日无强传导候选 → 命令零方案 CANCELLED(留痕可见,不占冷却,次日重试)——"宁缺毋滥"纪律优先于宏观意愿,这是设计取向不是缺陷。 +- **宏观闸的机会成本**:偏热确认期间个股级的好机会也会被一并拦下(闸不看个股质地)——这是"个股动作受宏观控制"的题中之义;propose_only 观察期里可以从 skipped 留痕回看拦了什么,评估闸的松紧。 +- **UNAVAILABLE 长期化**:数据停更时宏观层安静失效(不动作、不落闸),只有面板黄条与日报提示——不会误动作,但也不会提醒你"该人工判断了",需要把面板纳入日常一眼。 diff --git a/scripts/calibrate_macro_signal.py b/scripts/calibrate_macro_signal.py new file mode 100644 index 0000000..c4350d4 --- /dev/null +++ b/scripts/calibrate_macro_signal.py @@ -0,0 +1,593 @@ +# -*- coding: utf-8 -*- +""" +宏观择时 · 股汇对冲指数 标定脚本 (只读, 一次性分析用) +===================================================== +配套 MACRO_TIMING_PLAN.md §10: 在编码接入之前, 先用真实数据回答四件事: + + 1. 三张源表 (zs_day_data / gp_fx_daily / gp_shibor) 的位置与列名 + —— 自动探查并打印, 供方案 §4.2 回填钉死; + 2. hedge_index 历史分布 —— ±25 阈值合不合身, 各阈值触发频率; + 3. 触发口径对比 —— 进区即动 (zone_enter, 含确认1/2日两档) vs 极值回落再动 + (zone_exit), 触发后 5/10/20 交易日上证走势孰优 (样本会很小, 当参考不当真理); + 4. 对数映射标定 —— target = S0·ln(1+e/k) 的 S0/k 建议值与 e→幅度 对照表。 + +运行 (桥机 factorevaluation, **不需要重建镜像** —— 从宿主工作树经 stdin 喂给容器 python; +`python -` 的 sys.path[0] 是工作目录 /app, config.settings 照常可导入): + + cd ~/tradingSystem # git pull 之后 + docker compose run --rm -T pms-web python - < scripts/calibrate_macro_signal.py \ + > /tmp/macro_calib_report.md + cat /tmp/macro_calib_report.md + + # 需要完整指数序列时 (CSV 到 stdout): + docker compose run --rm -T pms-web python - --dump-csv < scripts/calibrate_macro_signal.py \ + > /tmp/hedge_series.csv + +**严格只读**: 只发 SELECT / SHOW, 不写任何表、不建任何东西。 +列名自动探查失败时会打出该表全部列名与样本行, 按提示用 + --fx-table/--fx-date-col/--fx-price-col/--fx-pair-col/--fx-pair + --shibor-table/--shibor-date-col/--shibor-value-col/--shibor-term-col/--shibor-term + --zs-date-col/--zs-close-col/--zs-symbol-col/--index-code +覆盖后重跑。数据源优先级默认 proxy(153),index(199) —— 用户口径两张 gp 表在 153。 +""" +from __future__ import annotations + +import argparse +import json +import math +import os +import sys +from datetime import date, datetime, timedelta + +try: + from sqlalchemy import create_engine, text +except Exception as e: # pragma: no cover + print(f"FATAL: 需要 sqlalchemy (+pymysql), 请在 PMS 容器里跑。{e}", file=sys.stderr) + sys.exit(2) + + +# ================================================================ 连接 +def _dsns() -> dict: + """DSN 取值: 优先 config.settings (容器内), 退回环境变量 (裸跑)。""" + out = {} + try: + from config.settings import settings # noqa + out["proxy"] = settings.PROXY_DB_URL + out["index"] = settings.DB_MYSQL_URL + except Exception: + out["proxy"] = os.environ.get("PROXY_DB_URL", "") + out["index"] = os.environ.get("DB_MYSQL_URL", "") + return {k: v for k, v in out.items() if v} + + +_engines = {} + + +def _eng(name: str, dsn: str): + if name not in _engines: + _engines[name] = create_engine(dsn, pool_pre_ping=True, + connect_args={"connect_timeout": 5}, future=True) + return _engines[name] + + +def _rows(name: str, dsn: str, sql: str, params=None) -> list: + with _eng(name, dsn).connect() as c: + return [dict(r) for r in c.execute(text(sql), params or {}).mappings().fetchall()] + + +# ================================================================ 探查 +def discover(table: str, sources: list, dsns: dict) -> dict: + """在各源上找这张表: 返回 {source, columns[], sample[]}; 全部失败返回 {errors}。""" + errors = {} + for src in sources: + dsn = dsns.get(src) + if not dsn: + errors[src] = "DSN 未配置" + continue + cols = None + try: + cols = [str(r.get("Field") or r.get("field")) for r in + _rows(src, dsn, f"SHOW COLUMNS FROM `{table}`")] + except Exception as e1: + try: # 代理不支持 SHOW 时退化: 取一行读键名 + sample1 = _rows(src, dsn, f"SELECT * FROM `{table}` LIMIT 1") + cols = list(sample1[0].keys()) if sample1 else None + if cols is None: + errors[src] = f"表存在但为空? SHOW 失败: {e1}" + continue + except Exception as e2: + errors[src] = f"{type(e2).__name__}: {str(e2)[:160]}" + continue + try: + sample = _rows(src, dsn, f"SELECT * FROM `{table}` LIMIT 3") + except Exception: + sample = [] + return {"source": src, "columns": cols, "sample": sample} + return {"errors": errors} + + +def pick(cols: list, cands: list): + low = {c.lower(): c for c in cols} + for c in cands: + if c in low: + return low[c] + for c in cands: # 次选: 前缀/包含 + for lc, orig in low.items(): + if lc.startswith(c) or c in lc: + return orig + return None + + +DATE_CANDS = ["trade_date", "timestamp", "date", "ymd", "day", "quote_date", "data_date", + "trade_day", "dt"] + + +def norm_ymd(v): + """任意日期形态 → int YYYYMMDD; 解析不了返回 None。""" + if v is None: + return None + if isinstance(v, (datetime, date)): + return int(v.strftime("%Y%m%d")) + s = str(v).strip()[:10].replace("-", "").replace("/", "") + if len(s) >= 8 and s[:8].isdigit(): + return int(s[:8]) + return None + + +def ymd_plus_days(ymd: int, n: int) -> int: + d = datetime.strptime(str(ymd), "%Y%m%d").date() + timedelta(days=n) + return int(d.strftime("%Y%m%d")) + + +# ================================================================ 取数 +def fetch_zs(args, sources, dsns): + info = discover("zs_day_data", sources, dsns) + if "errors" in info: + return None, info + cols = info["columns"] + dcol = args.zs_date_col or pick(cols, DATE_CANDS) + ccol = args.zs_close_col or pick(cols, ["close", "close_price", "px_close", "price"]) + scol = args.zs_symbol_col or pick(cols, ["symbol", "ts_code", "code", "index_code"]) + info.update({"date_col": dcol, "close_col": ccol, "symbol_col": scol}) + if not (dcol and ccol and scol): + info["errors"] = {"guess": f"列名猜不全 date={dcol} close={ccol} symbol={scol}"} + return None, info + src, dsn = info["source"], dsns[info["source"]] + rows = [] + for code in (args.index_code, args.index_code.replace(".SH", ""), + "SH" + args.index_code.split(".")[0]): + try: + rows = _rows(src, dsn, + f"SELECT `{dcol}` AS d, `{ccol}` AS v FROM `zs_day_data` " + f"WHERE `{scol}` = :c ORDER BY `{dcol}` DESC LIMIT :n", + {"c": code, "n": args.days}) + except Exception as e: + info["errors"] = {"query": f"{type(e).__name__}: {str(e)[:160]}"} + return None, info + if rows: + info["symbol_used"] = code + break + series = sorted([(norm_ymd(r["d"]), float(r["v"])) for r in rows + if norm_ymd(r["d"]) and r["v"] not in (None, 0)], key=lambda x: x[0]) + return series, info + + +def fetch_fx(args, sources, dsns): + info = discover(args.fx_table, sources, dsns) + if "errors" in info: + return None, info + cols = info["columns"] + dcol = args.fx_date_col or pick(cols, DATE_CANDS) + vcol = args.fx_price_col or pick(cols, ["close", "price", "rate", "mid", "value", + "exchange_rate", "cnh", "px"]) + pcol = args.fx_pair_col or pick(cols, ["currency", "ccy_pair", "ccy", "pair", "symbol", + "code", "name", "curr", "currency_pair"]) + info.update({"date_col": dcol, "price_col": vcol, "pair_col": pcol}) + if not (dcol and vcol): + info["errors"] = {"guess": f"列名猜不全 date={dcol} price={vcol}"} + return None, info + src, dsn = info["source"], dsns[info["source"]] + where, params = "", {"n": args.days} + if args.fx_where: + where = f"WHERE {args.fx_where}" + elif pcol: + try: + vals = [str(list(r.values())[0]) for r in + _rows(src, dsn, f"SELECT DISTINCT `{pcol}` AS p FROM `{args.fx_table}` LIMIT 60")] + except Exception: + vals = [] + info["pair_values_seen"] = vals[:30] + want = args.fx_pair + if not want: + cands = [v for v in vals if "USD" in v.upper() and "CN" in v.upper()] + cnh = [v for v in cands if "CNH" in v.upper()] + want = (cnh or cands or [None])[0] + info["pair_used"] = want + if want: + where, params = f"WHERE `{pcol}` = :p", {"p": want, "n": args.days} + else: + info["note_pair"] = "没找到 USD/CN* 形态的品种值, 按整表取 (若整表就是 USDCNH 则正确)" + rows = _rows(src, dsn, + f"SELECT `{dcol}` AS d, `{vcol}` AS v FROM `{args.fx_table}` {where} " + f"ORDER BY `{dcol}` DESC LIMIT :n", params) + series = sorted([(norm_ymd(r["d"]), float(r["v"])) for r in rows + if norm_ymd(r["d"]) and r["v"] not in (None, 0)], key=lambda x: x[0]) + return series, info + + +def fetch_shibor(args, sources, dsns): + info = discover(args.shibor_table, sources, dsns) + if "errors" in info: + return None, info + cols = info["columns"] + dcol = args.shibor_date_col or pick(cols, DATE_CANDS) + wide = args.shibor_value_col or pick(cols, ["shibor_1m", "1m", "m1", "shibor1m", + "rate_1m", "one_month"]) + info.update({"date_col": dcol, "wide_1m_col": wide}) + if not dcol: + info["errors"] = {"guess": "找不到日期列"} + return None, info + src, dsn = info["source"], dsns[info["source"]] + if wide: # 宽表: 每期限一列 + rows = _rows(src, dsn, + f"SELECT `{dcol}` AS d, `{wide}` AS v FROM `{args.shibor_table}` " + f"ORDER BY `{dcol}` DESC LIMIT :n", {"n": args.days}) + else: # 长表: 期限一列 + 值一列 + tcol = args.shibor_term_col or pick(cols, ["term", "period", "tenor", "name", + "type", "item"]) + vcol = pick(cols, ["rate", "value", "shibor", "price", "close"]) + info.update({"term_col": tcol, "value_col": vcol}) + if not (tcol and vcol): + info["errors"] = {"guess": f"长表列名猜不全 term={tcol} value={vcol}"} + return None, info + terms = ([args.shibor_term] if args.shibor_term + else ["1M", "1m", "30", "30D", "1月", "1个月", "一个月"]) + rows = [] + for t in terms: + rows = _rows(src, dsn, + f"SELECT `{dcol}` AS d, `{vcol}` AS v FROM `{args.shibor_table}` " + f"WHERE `{tcol}` = :t ORDER BY `{dcol}` DESC LIMIT :n", + {"t": t, "n": args.days}) + if rows: + info["term_used"] = t + break + series = sorted([(norm_ymd(r["d"]), float(r["v"])) for r in rows + if norm_ymd(r["d"]) and r["v"] is not None], key=lambda x: x[0]) + return series, info + + +# ================================================================ 对齐与计算 +def asof_align(trade_days: list, series: list, shift_days: int = 0) -> tuple: + """把 (ymd,value) 序列 as-of 对齐到交易日历。shift_days: 先把数据日期 +N 自然日。 + 返回 (对齐后的值列表, 补齐天数)。找不到任何前值的交易日置 None。""" + if shift_days: + series = [(ymd_plus_days(d, shift_days), v) for d, v in series] + series.sort(key=lambda x: x[0]) + vals, filled, j, last = [], 0, 0, None + for t in trade_days: + while j < len(series) and series[j][0] <= t: + last = series[j][1] + j += 1 + exact = j > 0 and series[j - 1][0] == t + if last is not None and not exact: + filled += 1 + vals.append(last) + return vals, filled + + +def log_rets(vals: list, win: int) -> list: + out = [None] * len(vals) + for i in range(win, len(vals)): + a, b = vals[i], vals[i - win] + if a and b and a > 0 and b > 0: + out[i] = math.log(a / b) + return out + + +def diffs(vals: list, win: int) -> list: + out = [None] * len(vals) + for i in range(win, len(vals)): + if vals[i] is not None and vals[i - win] is not None: + out[i] = vals[i] - vals[i - win] + return out + + +def roll_z(vals: list, win: int) -> list: + out = [None] * len(vals) + for i in range(len(vals)): + w = [v for v in vals[max(0, i - win + 1): i + 1] if v is not None] + if len(w) < win: + continue + m = sum(w) / len(w) + var = sum((x - m) ** 2 for x in w) / (len(w) - 1) + sd = math.sqrt(var) + if sd > 1e-12 and vals[i] is not None: + out[i] = (vals[i] - m) / sd * 10.0 + return out + + +def pctl(sorted_vals: list, q: float): + if not sorted_vals: + return None + k = (len(sorted_vals) - 1) * q + lo, hi = int(math.floor(k)), int(math.ceil(k)) + if lo == hi: + return sorted_vals[lo] + return sorted_vals[lo] + (sorted_vals[hi] - sorted_vals[lo]) * (k - lo) + + +# ================================================================ 区段与事件 +def episodes(idx: list, th: float, exit_band: float, side: int) -> list: + """side=+1 找 HOT (v>th), side=-1 找 COLD (v<-th)。带迟滞: 退出条件 |v| th: + in_ep, st, mx = True, i, sv - th + elif in_ep: + mx = max(mx, sv - th) + if sv < exit_band: + out.append({"start": st, "exit_i": i, "max_depth": mx, "len": i - st}) + in_ep = False + if in_ep: + out.append({"start": st, "exit_i": None, "max_depth": mx, "len": len(idx) - st}) + return out + + +def fwd_ret(close: list, i: int, h: int): + if i + h < len(close) and close[i] and close[i + h]: + return close[i + h] / close[i] - 1.0 + return None + + +def ev_stats(close: list, events: list, horizons=(5, 10, 20)) -> dict: + out = {} + for h in horizons: + rs = [fwd_ret(close, i, h) for i in events] + rs = [r for r in rs if r is not None] + if not rs: + out[h] = {"n": 0} + continue + rs_sorted = sorted(rs) + out[h] = {"n": len(rs), "mean": sum(rs) / len(rs), + "median": pctl(rs_sorted, 0.5), + "win_pos": sum(1 for r in rs if r > 0) / len(rs)} + return out + + +def fmt_ev(st: dict) -> str: + ps = [] + for h in (5, 10, 20): + s = st.get(h) or {} + if not s.get("n"): + ps.append(f"{h}日:无样本") + else: + ps.append(f"{h}日: n={s['n']} 均值{s['mean']*100:+.2f}% " + f"中位{s['median']*100:+.2f}% 上涨占比{s['win_pos']*100:.0f}%") + return " · ".join(ps) + + +# ================================================================ 对数映射标定 +def calib_log(depths: list, mid_target=0.05, p95_target=0.12) -> dict: + """解 S0·ln(1+m/k)=mid_target 且 S0·ln(1+p/k)=p95_target。 + 比值方程对 k 二分; 无解 (p/m 太小) 则锚中位数, k 固定 10。""" + ds = sorted(depths) + if len(ds) < 3: + return {"ok": False, "why": f"极值区段太少 ({len(ds)} 段), 用默认 S0=0.12 k=10", + "S0": 0.12, "k": 10.0, "m": pctl(ds, 0.5) if ds else None} + m, p = max(pctl(ds, 0.5), 0.5), max(pctl(ds, 0.95), 1.0) + R = p95_target / mid_target + f = lambda k: math.log(1 + p / k) / math.log(1 + m / k) + lo, hi = 1e-3, 1e6 + if f(hi) < R: # k→∞ 比值→p/m 仍不够 → 无解 + S0 = mid_target / math.log(1 + m / 10.0) + return {"ok": False, "why": f"深度分布太窄 (中位{m:.1f} / 95分位{p:.1f}), " + f"锚中位数取 S0, k 固定 10", "S0": S0, "k": 10.0, + "m": m, "p": p} + for _ in range(200): + mid = math.sqrt(lo * hi) + # f(k) 随 k 单调递增 (k→0 时→1, k→∞ 时→p/m): 比值还不够大就要更大的 k + if f(mid) < R: + lo = mid + else: + hi = mid + k = math.sqrt(lo * hi) + S0 = mid_target / math.log(1 + m / k) + return {"ok": True, "S0": S0, "k": k, "m": m, "p": p} + + +# ================================================================ 主流程 +def main(): + ap = argparse.ArgumentParser(description="股汇对冲指数标定 (只读)") + ap.add_argument("--days", type=int, default=1600, help="回看条数 (交易日, 默认约6年)") + ap.add_argument("--index-code", default="000001.SH") + ap.add_argument("--beta", type=float, default=0.02) + ap.add_argument("--ret-win", type=int, default=20) + ap.add_argument("--z-win", type=int, default=40) + ap.add_argument("--th", type=float, default=25.0, help="主阈值") + ap.add_argument("--exit-band", type=float, default=15.0) + ap.add_argument("--thresholds", default="20,25,30", help="敏感性对比的阈值列表") + ap.add_argument("--source-priority", default="proxy,index", + help="按序尝试的数据源 (用户口径: 两张 gp 表在 153=proxy)") + ap.add_argument("--fx-table", default="gp_fx_daily") + ap.add_argument("--fx-date-col"), ap.add_argument("--fx-price-col") + ap.add_argument("--fx-pair-col"), ap.add_argument("--fx-pair") + ap.add_argument("--fx-where", help="整段 WHERE 逃生口, 如 \"ccy='USDCNH'\"") + ap.add_argument("--shibor-table", default="gp_shibor") + ap.add_argument("--shibor-date-col"), ap.add_argument("--shibor-value-col") + ap.add_argument("--shibor-term-col"), ap.add_argument("--shibor-term") + ap.add_argument("--zs-date-col"), ap.add_argument("--zs-close-col") + ap.add_argument("--zs-symbol-col") + ap.add_argument("--dump-csv", action="store_true", help="只输出完整序列 CSV") + args = ap.parse_args() + + dsns = _dsns() + if not dsns: + print("FATAL: PROXY_DB_URL / DB_MYSQL_URL 都拿不到 (容器内跑, 或导出环境变量)", + file=sys.stderr) + sys.exit(2) + sources = [s.strip() for s in args.source_priority.split(",") if s.strip() in dsns] + + # ---- 取数 ---- + zs, zi = fetch_zs(args, sources, dsns) + fx, fi = fetch_fx(args, sources, dsns) + sh, si = fetch_shibor(args, sources, dsns) + + def head(name, info, series): + lines = [f"### {name}"] + if info.get("source"): + lines.append(f"- 源: **{info['source']}** ({dsns[info['source']].split('@')[-1]})") + if info.get("columns"): + lines.append(f"- 全部列: `{', '.join(info['columns'])}`") + for k in ("date_col", "close_col", "symbol_col", "symbol_used", "price_col", + "pair_col", "pair_used", "wide_1m_col", "term_col", "value_col", + "term_used", "note_pair"): + if info.get(k): + lines.append(f"- {k}: `{info[k]}`") + if info.get("pair_values_seen"): + lines.append(f"- 品种值样本: {info['pair_values_seen']}") + if series: + lines.append(f"- 数据: {len(series)} 条, {series[0][0]} → {series[-1][0]}, " + f"末值 {series[-1][1]}") + if info.get("errors"): + lines.append(f"- **失败**: {info['errors']} —— 按文件头提示带覆盖参数重跑") + if info.get("sample"): + lines.append(f"- 样本行: `{json.dumps(info['sample'][:1], ensure_ascii=False, default=str)[:400]}`") + return "\n".join(lines) + + if not args.dump_csv: + print("# 股汇对冲指数 · 标定报告") + print(f"\n> 生成: {datetime.now().isoformat(timespec='seconds')} · " + f"参数: ret_win={args.ret_win} z_win={args.z_win} beta={args.beta} " + f"主阈值±{args.th} 退出带±{args.exit_band} 回看={args.days}\n") + print("## 一、表探查 (回填方案 §4.2 用)\n") + print(head("zs_day_data (上证)", zi, zs), "\n") + print(head(f"{args.fx_table} (USD/CNH)", fi, fx), "\n") + print(head(f"{args.shibor_table} (SHIBOR 1M)", si, sh), "\n") + + if not (zs and fx and sh): + if args.dump_csv: + print("FATAL: 取数不全, 先跑一遍报告模式看探查结果", file=sys.stderr) + else: + print("\n**取数不全, 后续标定跳过。** 按上面失败提示带覆盖参数重跑。") + sys.exit(1) + + # ---- 对齐 (交易日历 = zs 日期) ---- + tdays = [d for d, _ in zs] + close = [v for _, v in zs] + fx_al, fx_fill = asof_align(tdays, fx, shift_days=1) # 汇率 +1 自然日再 as-of + sh_al, sh_fill = asof_align(tdays, sh, shift_days=0) + + # ---- 计算 ---- + sr = log_rets(close, args.ret_win) + fr = log_rets(fx_al, args.ret_win) + sd = diffs(sh_al, args.ret_win) + spread = [None if (sr[i] is None or fr[i] is None) else sr[i] + fr[i] + for i in range(len(tdays))] + spread_adj = [None if (spread[i] is None or sd[i] is None) + else spread[i] - args.beta * sd[i] for i in range(len(tdays))] + idx = roll_z(spread_adj, args.z_win) + + if args.dump_csv: + print("ymd,sh_close,fx,shibor_1m,stock_ret20,fx_ret20,spread,spread_adj,hedge_index") + for i, d in enumerate(tdays): + row = [d, close[i], fx_al[i], sh_al[i], sr[i], fr[i], spread[i], + spread_adj[i], idx[i]] + print(",".join("" if v is None else (f"{v:.6f}" if isinstance(v, float) else str(v)) + for v in row)) + return + + valid = [v for v in idx if v is not None] + print("## 二、数据体检\n") + print(f"- 交易日历 (取自上证): {len(tdays)} 天, 最新 {tdays[-1]} (距今天 " + f"{(date.today() - datetime.strptime(str(tdays[-1]), '%Y%m%d').date()).days} 自然日)") + print(f"- 汇率末日 {fx[-1][0]} · as-of 补齐 {fx_fill} 天; " + f"SHIBOR 末日 {sh[-1][0]} · 补齐 {sh_fill} 天") + print(f"- hedge_index 有效样本 {len(valid)} 天 (预热损耗 {len(tdays) - len(valid)} 天)") + if valid: + print(f"- 当前值: **{valid[-1]:+.1f}**") + + if len(valid) < 120: + print("\n**有效样本不足 120 天, 分布与口径对比意义有限, 到此为止。**") + sys.exit(1) + + sv = sorted(valid) + print("\n分位数: " + " · ".join( + f"P{int(q*100)}={pctl(sv, q):+.1f}" for q in (0.01, 0.05, 0.25, 0.5, 0.75, 0.95, 0.99))) + + # ---- 阈值敏感性 ---- + print("\n## 三、阈值敏感性 (±TH 触发频率与事后走势)\n") + print("| TH | 超阈天数占比 | HOT段/年 | COLD段/年 | HOT进区后10日均值 | COLD进区后10日均值 |") + print("|---|---|---|---|---|---|") + years = max(len(valid) / 244.0, 0.1) + for th in [float(x) for x in args.thresholds.split(",")]: + eb = th * args.exit_band / args.th # 退出带按比例缩放 + hot = episodes(idx, th, eb, +1) + cold = episodes(idx, th, eb, -1) + beyond = sum(1 for v in valid if abs(v) > th) / len(valid) + h10 = ev_stats(close, [e["start"] for e in hot]).get(10, {}) + c10 = ev_stats(close, [e["start"] for e in cold]).get(10, {}) + f = lambda s: (f"{s['mean']*100:+.2f}% (n={s['n']})" if s.get("n") else "无样本") + print(f"| ±{th:.0f} | {beyond*100:.1f}% | {len(hot)/years:.1f} | " + f"{len(cold)/years:.1f} | {f(h10)} | {f(c10)} |") + + # ---- 触发口径对比 (主阈值) ---- + th, eb = args.th, args.exit_band + hot, cold = episodes(idx, th, eb, +1), episodes(idx, th, eb, -1) + print(f"\n## 四、触发口径对比 (主阈值±{th:.0f}, 退出带±{eb:.0f})\n") + print(f"- HOT 区段 {len(hot)} 段 (平均持续 " + f"{sum(e['len'] for e in hot)/max(len(hot),1):.1f} 天, " + f"最大深度 {max((e['max_depth'] for e in hot), default=0):.1f}); " + f"COLD 区段 {len(cold)} 段 (平均 " + f"{sum(e['len'] for e in cold)/max(len(cold),1):.1f} 天, " + f"最大深度 {max((e['max_depth'] for e in cold), default=0):.1f})\n") + + def second_day(eps): # 确认=2: 区段第 2 天 (不足 2 天的区段无事件) + return [e["start"] + 1 for e in eps if e["len"] >= 2] + + rows = [ + ("HOT · 进区即动(确认1)", [e["start"] for e in hot], "降仓视角: 越负越好"), + ("HOT · 进区确认2日", second_day(hot), "降仓视角: 越负越好"), + ("HOT · 极值回落再动", [e["exit_i"] for e in hot if e["exit_i"] is not None], + "降仓视角: 越负越好(但可能已回落)"), + ("COLD · 进区即动(确认1)", [e["start"] for e in cold], "升仓视角: 越正越好"), + ("COLD · 进区确认2日", second_day(cold), "升仓视角: 越正越好"), + ("COLD · 极值回落再动", [e["exit_i"] for e in cold if e["exit_i"] is not None], + "升仓视角: 越正越好"), + ] + for name, evs, note in rows: + print(f"**{name}** ({note})\n> {fmt_ev(ev_stats(close, evs))}\n") + print("> 提醒: 样本很小, 这是参考不是显著性检验; 两口径都已实现、页面可切, 这里只定默认档。\n") + + # ---- 对数映射标定 ---- + print(f"## 五、对数映射标定 target = S0 · ln(1 + e/k), e = |idx| − {th:.0f}\n") + depths = [e["max_depth"] for e in hot + cold] + cal = calib_log(depths) + if depths: + dsrt = sorted(depths) + print(f"- 区段最大深度 e 分布: 中位 {pctl(dsrt, 0.5):.1f} · " + f"P95 {pctl(dsrt, 0.95):.1f} · 最大 {dsrt[-1]:.1f} (共 {len(depths)} 段)") + tag = "标定成功" if cal.get("ok") else f"退化取值 ({cal.get('why')})" + print(f"- **建议初值: PMS_MACRO_LOG_S0 = {cal['S0']:.3f} · " + f"PMS_MACRO_LOG_K = {cal['k']:.1f}** —— {tag}") + print(f"- 目标形状: e 中位数 → 约 5 个点仓位, e 95分位 → 约 12 个点, " + f"SHIFT_MAX 封顶 20 个点\n") + print("| e (超额深度) | 2 | 5 | 8 | 10 | 15 | 20 | 30 | 40 |") + print("|---|" + "---|" * 8) + tgt = lambda e: min(cal["S0"] * math.log(1 + e / cal["k"]), 0.20) + print("| 累计调整幅度 | " + " | ".join(f"{tgt(e)*100:.1f}%" for e in + (2, 5, 8, 10, 15, 20, 30, 40)) + " |") + + print("\n## 六、把这些数拿回去干什么\n") + print("1. 表探查一节的 源/列名 → 回填 MACRO_TIMING_PLAN.md §4.2, 才能写 macro_repo;") + print("2. 第三节选 TH 与退出带 → PMS_MACRO_HOT_TH / COLD_TH / EXIT_BAND;") + print("3. 第四节选触发口径默认档 → PMS_MACRO_TRIGGER_MODE / CONFIRM_DAYS;") + print("4. 第五节的 S0/k → PMS_MACRO_LOG_S0 / LOG_K;") + print("5. 报告全文存档进仓库 (建议 docs/ 或贴回对话), 作为参数初值的依据留痕。") + + +if __name__ == "__main__": + main()