feat(data): 静态表调度挪周末+vintage盘前自检——③sanguo-ak-quarter月内1号02:00改财报月首个周六02:00(/d SAT,周六02:00→周一09:30有55h跑道;旧排程7.5h夜窗连单表8h都装不下,09-01注册后首爆16h+穿整个交易日,09:32决策窗撞balance(94%Q2'26)/income(0%)vintage撕裂致value ROE归零清仓;月内1-4号读旧但一致,策略配对修复侧兼容)④新增static_vintage_check.py每日09:00自检:三报表最大报告期+覆盖率盘点(REPORT_DATE列裁剪扫描),跨表撕裂/估值断更(>48h)告警exit 3,原子落vintage_status.json供策略盘前预判(策略session协作契约);register加sanguo-ak-vintage+新wrapper;+8测试 [vps]
CI/CD / test (push) Successful in 3s
CI/CD / nas-deploy (push) Successful in 8s
CI/CD / nas-verify (push) Successful in 14s

This commit is contained in:
2026-09-01 19:03:53 +08:00
parent aa7e7b8eba
commit 8dc3f494b7
5 changed files with 328 additions and 4 deletions
+4 -2
View File
@@ -1,5 +1,7 @@
# ak_quarter_wrapper.ps1 — sanguo-ak-quarter schtask wrapper (财报季 APR/MAY/SEP/NOV 02:00 akshare 三表+预告/快报 --force)
# balance+income+cashflow+forecast+express; per-stock 全量(5500×5), 财报季月夜间跑
# ak_quarter_wrapper.ps1 — sanguo-ak-quarter schtask wrapper (财报季月首个周六 02:00 akshare 三表+预告/快报 --force)
# balance+income+cashflow+forecast+express; per-stock 全量(5500×5表)单趟 20h+,
# 只能周末跑 (2026-09-01 事故: 月内1号02:00 起 7.5h 夜窗装不下, 洗库穿交易日,
# 09:32 决策窗撞 balance/income vintage 撕裂; 首个周六起跑周一 09:30 前有 55h 跑道)
$env:http_proxy = ''
$env:https_proxy = ''
$env:all_proxy = ''
@@ -0,0 +1,13 @@
# ak_vintage_wrapper.ps1 — sanguo-ak-vintage schtask wrapper (每日 09:00 静态表 vintage 一致性自检)
# 三报表(balance/income/cashflow)最大报告期不一致(洗库中段撕裂)或估值表断更 → ALERT + 退出码 3;
# 落 data/static/vintage_status.json 供策略盘前预判 (2026-09-01 洗库撕裂事故防线④)
$env:http_proxy = ''
$env:https_proxy = ''
$env:all_proxy = ''
Set-Location C:\sanguo_vnpy_v2
$ts = Get-Date -Format 'yyyyMMdd_HHmmss'
$logDir = 'C:\sanguo_vnpy_v2\data\migration_logs'
if (-not (Test-Path $logDir)) { New-Item -ItemType Directory -Path $logDir -Force | Out-Null }
$log = Join-Path $logDir "ak_vintage_$ts.txt"
C:\Python310\python.exe -X utf8 C:\sanguo_vnpy_v2\scripts\data_platform\static_vintage_check.py *>> $log
exit $LASTEXITCODE
@@ -9,10 +9,17 @@
schtasks --% /create /tn sanguo-ak-eod /tr "powershell -ExecutionPolicy Bypass -File C:\sanguo_vnpy_v2\scripts\data_platform\ak_eod_wrapper.ps1" /sc daily /st 19:00 /ru SYSTEM /rl HIGHEST /f
schtasks --% /create /tn sanguo-ak-events /tr "powershell -ExecutionPolicy Bypass -File C:\sanguo_vnpy_v2\scripts\data_platform\ak_events_wrapper.ps1" /sc daily /st 19:30 /ru SYSTEM /rl HIGHEST /f
schtasks --% /create /tn sanguo-ak-stock /tr "powershell -ExecutionPolicy Bypass -File C:\sanguo_vnpy_v2\scripts\data_platform\ak_stock_wrapper.ps1" /sc weekly /d SAT /st 03:00 /ru SYSTEM /rl HIGHEST /f
schtasks --% /create /tn sanguo-ak-quarter /tr "powershell -ExecutionPolicy Bypass -File C:\sanguo_vnpy_v2\scripts\data_platform\ak_quarter_wrapper.ps1" /sc monthly /m APR,MAY,SEP,NOV /st 02:00 /ru SYSTEM /rl HIGHEST /f
# ak-quarter: 财报月首个周六 02:00 (2026-09-01 事故修正: 旧 /d 缺省=每月1号02:00,
# 夜间窗口 7.5h 连一张表(~8h)都装不下, 09-01 首爆 balance 02:00-10:01 / income
# 10:01-18:17 穿整个交易日, 09:32 决策窗撞两表 vintage 撕裂 → value 清仓。
# 首个周六 02:00 → 周一 09:30 有 55h 跑道; 月内 1-4 号读旧但一致, 策略配对兼容)
schtasks --% /create /tn sanguo-ak-quarter /tr "powershell -ExecutionPolicy Bypass -File C:\sanguo_vnpy_v2\scripts\data_platform\ak_quarter_wrapper.ps1" /sc monthly /m APR,MAY,SEP,NOV /d SAT /st 02:00 /ru SYSTEM /rl HIGHEST /f
# ak-vintage: 每日 09:00 静态表 vintage 一致性自检 (防线④: 撕裂/断更告警退出码 3,
# 落 data/static/vintage_status.json 供策略盘前预判)
schtasks --% /create /tn sanguo-ak-vintage /tr "powershell -ExecutionPolicy Bypass -File C:\sanguo_vnpy_v2\scripts\data_platform\ak_vintage_wrapper.ps1" /sc daily /st 09:00 /ru SYSTEM /rl HIGHEST /f
schtasks --% /create /tn sanguo-index /tr "powershell -ExecutionPolicy Bypass -File C:\sanguo_vnpy_v2\scripts\data_platform\index_monthly_wrapper.ps1" /sc monthly /d 16 /st 19:50 /ru SYSTEM /rl HIGHEST /f
Write-Output "=== VERIFY Arguments 路径(应全路径, 非 \xxx.ps1)==="
foreach ($t in @('sanguo-ak-eod','sanguo-ak-events','sanguo-ak-stock','sanguo-ak-quarter','sanguo-index')) {
foreach ($t in @('sanguo-ak-eod','sanguo-ak-events','sanguo-ak-stock','sanguo-ak-quarter','sanguo-ak-vintage','sanguo-index')) {
$xml = schtasks /query /tn $t /xml 2>&1 | Out-String
if ($xml -match 'File\s+([^<"]+)') {
Write-Output (" {0,-20} File {1}" -f $t, $matches[1])
@@ -0,0 +1,181 @@
#!/usr/bin/env python3
# -*- coding: utf-8 -*-
"""静态表 vintage 一致性自检 (2026-09-01 洗库撕裂事故防线④)。
背景: sanguo-ak-quarter 财报季全量洗库逐股重写 balance→income→cashflow,
单表 ~8h、全程 20h+: 洗库中段两表 vintage 撕裂 (09-01 09:32 balance 已带
中报而 income 停在 Q1'26 → value ROE 配对归零清仓)。本脚本盘前盘点各报表
最新报告期分布, 跨表不一致/估值断更大声告警, 并落 vintage_status.json 供
策略侧盘前预判 (契约: tables.<type> = {max_period, coverage, n_at_max,
n_valid, n_empty, n_unreadable, newest_mtime}, checked_at 为扫描时刻)。
只读, 不依赖 akshare。用法 (VPS 09:00 schtask sanguo-ak-vintage):
python static_vintage_check.py # 默认扫 C:\\sanguo_vnpy_v2\\data\\static
python static_vintage_check.py --root <dir> --json-out <path>
退出码: 0=一致; 3=有告警 (schtask Last Result 可见)
"""
import argparse
import datetime
import json
import os
import sys
from pathlib import Path
from typing import Tuple
import pandas as pd
DEFAULT_ROOT = r"C:\sanguo_vnpy_v2\data\static"
STATEMENT_TABLES = ("balance", "income", "cashflow") # 有 REPORT_DATE 列
VALUATION_STALE_HOURS = 48.0 # 估值表(每日19:00 ak-eod)最新mtime超此数=断更
EMPTY_PARQUET_MIN_BYTES = 1024 # 与 akshare_static_download 同口径
def _norm_period(v) -> str:
"""REPORT_DATE 值 → 'YYYY-MM-DD' (兼容 str / datetime 两种 dtype)。"""
return str(v)[:10]
def scan_statement_table(table_dir: Path) -> dict:
"""扫一张报表: 每股最大 REPORT_DATE → 表级分布。
只读 REPORT_DATE 一列 (parquet 列裁剪, 5550 股 ~1min 级)。
空文件 (<1KB)/0 行/读失败分别计数, 不中断扫描。
"""
stats = {
"n_total": 0, "n_empty": 0, "n_unreadable": 0,
"n_valid": 0, "n_at_max": 0, "max_period": None,
"newest_mtime": None, "coverage": 0.0,
}
if not table_dir.is_dir():
stats["error"] = "目录不存在"
return stats
newest = 0.0
per_max = []
for f in sorted(table_dir.glob("*.parquet")):
stats["n_total"] += 1
try:
st = f.stat()
newest = max(newest, st.st_mtime)
if st.st_size < EMPTY_PARQUET_MIN_BYTES:
stats["n_empty"] += 1
continue
df = pd.read_parquet(f, columns=["REPORT_DATE"])
if df.empty:
stats["n_empty"] += 1
continue
per_max.append(_norm_period(df["REPORT_DATE"].max()))
stats["n_valid"] += 1
except Exception:
stats["n_unreadable"] += 1
if per_max:
mx = max(per_max)
stats["max_period"] = mx
stats["n_at_max"] = sum(1 for v in per_max if v == mx)
stats["coverage"] = round(stats["n_at_max"] / stats["n_valid"], 4)
if newest:
stats["newest_mtime"] = datetime.datetime.fromtimestamp(
newest).strftime("%Y-%m-%d %H:%M:%S")
return stats
def valuation_freshness(root: Path) -> dict:
"""估值表 (每日 ak-eod 19:00) 新鲜度: 最新 parquet mtime 距今几小时。"""
info = {"newest_mtime": None, "age_hours": None}
d = root / "valuation"
if not d.is_dir():
info["error"] = "目录不存在"
return info
newest = 0.0
for f in d.glob("*.parquet"):
try:
newest = max(newest, f.stat().st_mtime)
except OSError:
continue
if newest:
ts = datetime.datetime.fromtimestamp(newest)
info["newest_mtime"] = ts.strftime("%Y-%m-%d %H:%M:%S")
info["age_hours"] = round(
(datetime.datetime.now() - ts).total_seconds() / 3600, 1)
return info
def check_all(root: Path) -> Tuple[dict, list]:
"""盘点三报表 + 估值新鲜度。返 (status, alerts); alerts 非空 → 退出码 3。"""
status = {
"checked_at": datetime.datetime.now().strftime("%Y-%m-%d %H:%M:%S"),
"root": str(root),
"tables": {},
"valuation": {},
}
alerts = []
periods = {}
for t in STATEMENT_TABLES:
s = scan_statement_table(root / t)
status["tables"][t] = s
if s["max_period"]:
periods[t] = s["max_period"]
else:
alerts.append(
"[%s] 无有效数据 (n_valid=0): %s" % (
t, s.get("error", "全空/不可读")))
if len(periods) == len(STATEMENT_TABLES) and len(set(periods.values())) > 1:
detail = ", ".join(f"{t}={p}" for t, p in sorted(periods.items()))
alerts.append(
"报表 vintage 撕裂: %s (跨表配对将缺期, 疑似洗库进行中/中段, "
"策略盘前勿按新报告期决策)" % detail)
status["valuation"] = valuation_freshness(root)
v = status["valuation"]
if v.get("age_hours") is not None and v["age_hours"] > VALUATION_STALE_HOURS:
alerts.append(
"估值表断更: 最新 mtime %s 距今 %.1fh (>%.0fh, ak-eod 未跑成?)"
% (v["newest_mtime"], v["age_hours"], VALUATION_STALE_HOURS))
return status, alerts
def write_json_atomic(path: Path, payload: dict) -> None:
"""原子写 JSON (tmp → os.replace, 洗库同款语义)。"""
tmp = path.with_suffix(path.suffix + ".tmp")
tmp.parent.mkdir(parents=True, exist_ok=True)
tmp.write_text(
json.dumps(payload, ensure_ascii=False, indent=1), encoding="utf-8")
os.replace(tmp, path)
def main() -> int:
ap = argparse.ArgumentParser(
description="静态表 vintage 一致性自检 (三报表+估值)")
ap.add_argument(
"--root", default=os.environ.get("STATIC_VINTAGE_ROOT", DEFAULT_ROOT),
help="static 根目录 (默认 %s)" % DEFAULT_ROOT)
ap.add_argument(
"--json-out",
default=os.environ.get(
"STATIC_VINTAGE_JSON",
os.path.join(DEFAULT_ROOT, "vintage_status.json")),
help="vintage_status.json 输出路径")
args = ap.parse_args()
status, alerts = check_all(Path(args.root))
write_json_atomic(Path(args.json_out), status)
for t, s in status["tables"].items():
print("[vintage] %-11s max=%s coverage=%.1f%% (%d/%d) empty=%d "
"unreadable=%d newest=%s" % (
t, s["max_period"], 100 * s["coverage"],
s["n_at_max"], s["n_valid"], s["n_empty"],
s["n_unreadable"], s["newest_mtime"]))
v = status["valuation"]
print("[vintage] valuation newest=%s age=%sh" % (
v.get("newest_mtime"), v.get("age_hours")))
print("[vintage] json → %s" % args.json_out)
if alerts:
for a in alerts:
print("[vintage][ALERT] %s" % a)
return 3
print("[vintage] OK: 三报表 vintage 一致")
return 0
if __name__ == "__main__":
sys.exit(main())