Files
sanguo_vnpy_v2/sanguo_factor/monthly_batch.py
T

275 lines
14 KiB
Python

# sanguo_factor/monthly_batch.py
"""先批后评·第一段(决议 J/§11.3): 月度批评日全量重评批——跑批从「人点」变「日历」.
薄 CLI: 注册表全量(排除 retired/graveyard)→组窗(12M,end 归一自然月末)→复用既有
run_batch_eval 真跑落 eval_db→第二段 monthly_review 读刚落库的批判定.
零评估逻辑(补中间不重建);window_months/排除态为起步默认,首年校准.
manifest 缺失/不可读→标记件 {report_dir}/data_gaps.json+告警照跑批(P2-11/12);
--manifest/env 建议绝对路径(调度器/env CWD 各异,相对路径易静默错位).
"""
from __future__ import annotations
import argparse
import calendar
import json
import os
import sqlite3
import sys
from typing import Any
import yaml
from sanguo_factor import data_gap_check as dgc
from sanguo_factor.batch_eval import run_batch_eval
from sanguo_factor.eval_store import default_eval_db_path
from sanguo_factor.registry import get_factor
from sanguo_factor.version_registry import load_registry
EXCLUDED_STATUSES = frozenset({"retired", "graveyard"}) # 死因子不占月度批算力
def month_end(as_of: str) -> str:
"""as_of 归一到当月最后一日(数据截止对齐自然月末,决议 H)."""
y, m = int(as_of[:4]), int(as_of[5:7])
return f"{y:04d}-{m:02d}-{calendar.monthrange(y, m)[1]:02d}"
def window_start(end: str, window_months: int = 12) -> str:
"""闭区间起点:含 end 当月共 window_months 个月的月初(起步默认 12M)."""
y, m = int(end[:4]), int(end[5:7])
total = y * 12 + (m - 1) - (window_months - 1)
return f"{total // 12:04d}-{total % 12 + 1:02d}-01"
def resolve_with_yaml(registry: dict):
"""D2 复合 resolver:内存静态库 miss→读 yaml 最新版 params.expression
动态注册.
分解产物无静态库文件,默认 get_factor 会把 yaml 因子丢进
skipped_unregistered(月度批评考不到,调研报告 §3 必修缺陷)——本 resolver
让 yaml 因子与手写因子同注册表同求值器(零特权).
"""
from .factor_guard import category_for_source, derive_source, precheck_expression
from .registry import get_factor, register_factor
def resolve(name: str):
found = get_factor(name)
if found is not None:
return found
entry = (registry.get("factors") or {}).get(name)
if entry is None:
return None
versions = entry.get("versions") or []
params = (versions[-1].get("params") or {}) if versions else {}
expr = params.get("expression")
if not expr:
return None
source = str(params.get("source") or "bars_daily")
# P1-4 纵深防御: yaml 是运行期可变文件,唯一把关方不能只有写入方——
# 注册(=下次 eval 执行)前跑 precheck+source 一致性,违规记
# skipped_unregistered+stderr,不直通 eval 汇点(eval 即代码执行)
violations = precheck_expression(expr)
derived = derive_source(expr) if not violations else None
if violations or not isinstance(derived, str) or derived != source:
why = (";".join(v.detail for v in violations) if violations
else f"source 声明 {source} 与推导 {derived} 不符")
print(f"[monthly_batch] 跳过未注册 {name}: {why}",
file=sys.stderr)
return None
register_factor(name, expr, category_for_source(source))
return get_factor(name)
return resolve
def plan_batch(registry: dict, as_of: str, window_months: int = 12,
resolver=None, daily: bool = False) -> dict[str, Any]:
"""注册表→批参数(main 只做 IO 编排,可单测).
「纯函数」契约勘正(P3-12): resolver=get_factor(默认)时纯读;
resolver=resolve_with_yaml 时带进程内注册副作用——首次解析 miss 名
会动态注册进全局因子库(register_factor),对进程内其他消费者可见;
跨调用幂等(已注册直接命中),但测试须留意跨用例泄漏.
血统名漂移防线(09-27 审计 F-3): 库中不可解析名剔除进
skipped_unregistered(不炸批,main 侧显式告警)——yaml 与因子库
曾漂移 5 名(fa_gross_margin/s0* vs 库名),月度批评会空转.
"""
resolve = resolver or get_factor
# 日度化(决议 H 修订):daily 时 as_of 直用(T-1 夜班),不归一月末;
# 月度缺省行为不变(月末快照=判定层语义,月度链照旧归一).
end = as_of if daily else month_end(as_of)
names = sorted(n for n, e in registry["factors"].items()
if e["status"] not in EXCLUDED_STATUSES)
skipped = [n for n in names if resolve(n) is None]
return {"factor_names": [n for n in names if resolve(n) is not None],
"skipped_unregistered": skipped,
"start": window_start(end, window_months),
"end": end,
"label": (f"daily_{as_of}" if daily else f"monthly_{as_of[:7]}")}
def decomposer_new_names(registry: dict, factor_names: list[str]) -> list[str]:
"""IC 闸「新因子」判定(审计 P3-11):条目级 origin 优先,旧数据回落
version params 全版本扫描(非仅最新).
双向漂移根因=旧判定读「最新版 params.origin」:分解因子月月在闸视野
永不收敛、手工升版(新 params 不带 origin)即无声退出闸视野。改以稳定
身份判定:条目级 origin 存在则以条目级为准(人工可显式改判退出);
否则回落扫全部版本 params(decompose 的 origin 戳落在含 v1 的各版本
上),手工升版不再影响判定——同注册表状态跨月判定一致.
"""
picked = []
for name in factor_names:
entry = registry["factors"].get(name) or {}
origin = entry.get("origin")
if origin is not None:
if origin == "decomposer":
picked.append(name)
continue
versions = entry.get("versions") or []
if any((v.get("params") or {}).get("origin") == "decomposer"
for v in versions):
picked.append(name)
return picked
def _reusable_daily_run(db: str, plan: dict[str, Any]) -> dict | None:
"""月末批可复用的同 end 日度批(2026-10-11 双层分工补丁,决议H:
物理分开≠逻辑重叠——次月1日夜班已产出 end=月末 的日度批时,月度链
stage1 不重算 40min,直接复用)。全条件满足才复用:label/end 钉住同
as_of、因子集一致(漂移=回退自算,不许部分复用)、完整批。
"""
from sanguo_factor import eval_store
want = f"daily_{plan['end']}"
plan_names = set(plan["factor_names"])
try:
runs = eval_store.list_runs(db)
except sqlite3.OperationalError: # 空/无 schema 库=无可复用,回退自算
return None
for r in runs:
if (r.get("label") == want and r.get("end") == plan["end"]
and r.get("factors_total") == len(plan_names)
and r.get("factors_done") == len(plan_names)):
got = {row["factor"] for row in eval_store.get_rows(db, r["run_id"])}
if got == plan_names:
return r
return None
def main(argv: list[str] | None = None) -> int:
ap = argparse.ArgumentParser(description="先批后评·第一段:注册表全量月度重评批")
ap.add_argument("--as-of", required=True, help="数据截止日 YYYY-MM-DD(归一自然月末)")
ap.add_argument("--registry",
default=os.environ.get("SANGUO_FACTOR_REGISTRY",
"config/factor_registry.yaml"))
ap.add_argument("--eval-db", default=None, help="缺省 eval_store.default_eval_db_path()")
ap.add_argument("--window-months", type=int, default=12)
ap.add_argument("--report-dir", default="reports/factor_monthly",
help="values 导出与 IC 正交闸报告目录")
ap.add_argument("--manifest", default="config/data_manifest.yaml",
help="底座数据域清单(gap 比对供给侧)")
ap.add_argument("--gap-issue", action="store_true",
help="有数据缺口时尝试开 Gitea issue(决议 M;默认关)")
ap.add_argument("--host", choices=("nas", "vps"), default=None,
help="运行侧别声明(--gap-issue 开单权威单机制 P3-14: "
"vps 侧拒开,NAS=权威侧;缺省=人工运行不限)")
ap.add_argument("--vnpy-db", default=None,
help="行情库路径 override(双机 cfg 路径各异时直传,"
"缺省走 config 解析)")
ap.add_argument("--fund-data-dir", default=None,
help="财务静态域根目录 override(VPS 镜像路径与 NAS 默认各异时直传,"
"缺省走 cfg/默认解析)")
ap.add_argument("--daily", action="store_true",
help="日度模式(决议 H 修订):as_of 直用不归一月末,"
"label=daily_ 前缀;夜班批 as-of=T-1")
args = ap.parse_args(argv)
registry = load_registry(args.registry)
plan = plan_batch(registry, args.as_of, args.window_months,
resolver=resolve_with_yaml(registry), daily=args.daily)
if not plan["factor_names"]:
print("[monthly_batch] 注册表无可跑因子(全 retired/graveyard?)", file=sys.stderr)
return 1
db = args.eval_db or default_eval_db_path()
if plan["skipped_unregistered"]:
print(f"[monthly_batch] ⚠️ 剔除库中不可解析名(血统漂移,须对齐): "
f"{plan['skipped_unregistered']}", file=sys.stderr)
# —— manifest 四态 fail-soft(P2-11/P2-12): 缺路径/坏语法/坏形状三态统一
# 落 manifest_missing 标记件+stderr 告警(数据侧最需告警的场景不再零输出),
# 缺口按缺失处理批照跑、不开单;正常态照常比对。标记件与独立 CLI 同位
# (report_dir/data_gaps.json);CLI 侧同情形 fail-fast exit 1(独立调度无批
# 可跑,契约差异见 data_gap_check docstring)。
gaps: dict[str, list[str]] = {}
manifest_failed = not os.path.exists(args.manifest)
if manifest_failed:
print(f"[monthly_batch] ⚠️ manifest 不存在,缺口按缺失处理照跑批,"
f"标记件 data_gaps.json 已落(建议 --manifest/env 用绝对路径): "
f"{args.manifest}", file=sys.stderr)
else:
try:
gaps = dgc.check_gaps(dgc.collect_requirements(registry),
dgc.load_manifest(args.manifest))
except (OSError, yaml.YAMLError, AttributeError) as exc:
manifest_failed = True
gaps = {}
print(f"[monthly_batch] ⚠️ manifest 不可读(坏语法/坏形状),缺口按"
f"缺失处理照跑批,标记件 data_gaps.json 已落: "
f"{args.manifest} ({exc!r})", file=sys.stderr)
if manifest_failed:
os.makedirs(args.report_dir, exist_ok=True)
with open(os.path.join(args.report_dir, "data_gaps.json"), "w",
encoding="utf-8") as f:
json.dump({"manifest_missing": True}, f, ensure_ascii=False, indent=1)
if gaps:
print(f"[monthly_batch] ⚠️ 数据缺口(决议 M,批照跑): "
+ "; ".join(f"{f}→{','.join(s)}" for f, s in sorted(gaps.items())),
file=sys.stderr)
if args.gap_issue:
num = dgc.maybe_open_gap_issue(gaps, args.manifest, host=args.host)
if num:
print(f"[monthly_batch] 已开数据缺口 issue #{num}")
print(f"[monthly_batch] {len(plan['factor_names'])} 因子 "
f"{plan['start']}~{plan['end']} → {db} (label={plan['label']})")
# 2026-10-11 双层分工补丁 B2(决议H): 月度链①段复用同 as_of 日度批——
# 日度模式自身不查(同 label 重跑=幂等覆盖);月度模式检测到 end=月末、
# label=daily_月末、因子集一致的完整批→跳过 stage1 重算与 ic_gate
# (不写 values/ic_gate 产物,日度链跑时已出过同 as_of 的;stage2
# monthly_review 经 pick_eval_run 自然选中该批)。日度批缺席/漂移→
# 落到下方原重算路径,行为不变。
reuse = None if args.daily else _reusable_daily_run(db, plan)
if reuse is not None:
print(f"[monthly_batch] 复用日度批 run_id={reuse['run_id']} "
f"label={reuse['label']}(决议H:物理分开≠逻辑重叠,重算白烧),"
f"跳过 stage1 重算与 ic_gate")
return 0
values_dir = os.path.join(args.report_dir, "values", plan["label"])
summary = run_batch_eval(plan["factor_names"], plan["start"], plan["end"],
db, plan["label"], factor_values_out=values_dir,
vnpy_db_override=args.vnpy_db,
fund_data_dir=args.fund_data_dir)
print(f"[monthly_batch] 完成: {summary}")
# IC 正交闸(D4): 新因子=origin decomposer;报告件不自动转移
new_names = decomposer_new_names(registry, plan["factor_names"])
if new_names:
from sanguo_factor.ic_gate import orthogonality_report, write_report
rep = orthogonality_report(values_dir, new_names, plan["factor_names"])
write_report(os.path.join(args.report_dir,
f"ic_gate_{plan['label']}.json"), rep)
# P2-9: 「查过且干净」与「根本没查成」可区分——缺件/零比对时
# 显式告警(系统性 eval 失败不再呈现绿灯);告警属信号语义,
# 不 exit 非零
if rep["checked"] < len(new_names) or rep["compared_pairs"] == 0:
print(f"[monthly_batch] ⚠️ IC 正交闸未完整执行: 新因子 "
f"{len(new_names)} 实查 {rep['checked']}(缺件跳过 "
f"{rep['skipped_missing_values']}),比对对数 "
f"{rep['compared_pairs']}——疑系统性 eval 失败,勿当绿灯",
file=sys.stderr)
if rep["flagged"]:
print(f"[monthly_batch] ⚠️ IC 正交闸: "
+ ", ".join(f"{f['factor']}(vs {f['vs']} ic={f['ic']})"
for f in rep["flagged"]))
return 0
if __name__ == "__main__":
raise SystemExit(main())