feat(data): P4-3 接线——corpus_reader.get_events 读出口+manifest 两域+run_corpus 尾两段(funnel/gate rc 隔离+LLM env 条件传递); fusion spec §20.16+runbook 落档; Mac 冒烟绿(词表件落+候选入队+门如实报缺) [vps] [nas]
CI/CD / test (push) Successful in 46s
CI/CD / nas-deploy (push) Successful in 2s
CI/CD / nas-verify (push) Successful in 10s

This commit is contained in:
2026-10-02 10:02:51 +08:00
parent 1ea143ec6c
commit e47dddd336
6 changed files with 96 additions and 0 deletions
+3
View File
@@ -35,3 +35,6 @@ domains:
sw_industry: {desc: 申万一级行业成份逐日快照(~5,400 行/夜), 状态: 生产}
dragon_tiger_seats: {desc: 龙虎榜席位明细(增量,挂载日起), 状态: 生产}
nb_daily: {desc: 北向 HKEX 日统计(T-1~T 窗), 状态: 生产}
# ---- 事件快照两域(P4-3 漏斗刀, spec §20.16) ----
events_in: {desc: LLM 候选输入切片(防编数:喂模型确切文本先落盘可回放), 状态: 生产}
events_llm: {desc: 事件判读快照(词表+LLM 统一,model 列区分;append-only 永不重算), 状态: 生产}
@@ -227,6 +227,7 @@ ssh 49.232.102.198 "powershell -Command \"Get-Content C:\sanguo_bigqmt\qmt_relog
- **ak-weekly 三新域(融合 spec §20.13,2026-09-22 挂载)**:`ak_stock_wrapper.ps1` 追加两条独立命令(不带 --force,marker 幂等增量):`--types pledge_ratio,pledge_detail --start 20210108` + `--types research_report --start 20210701`。pledge_ratio 仅枚举周五(非周五 EM result=None 必失败),未发布周五 failed 不写 marker 次周自愈;pledge_detail=RPTA_APP_ACCUMDETAILS 254 页走 fetch_em_dc_pages 节流通道(1.2s/页≈6 分钟)月桶留存;research_report=reportapi 直连周窗 anchor=周五。
- **反向镜像链(融合 spec §20.11,2026-09-17 新增)**:`push_vps_mirror_daily.sh`+`vps_mirror_lib.py`(纯逻辑,宿主 python3.8 stdlib)——corpus **五子域**(09-22 起 +flash_meta 快讯域)+dmsk 五域 NAS→VPS **推送**(09-18 起挂 `sanguo-corpus-daily` DSM 任务尾部=`run_corpus_standalone.sh` EXIT trap 完赛即推,典型 ~10:1x/backfill 日 ~17:3x;rc 隔离=corpus rc∈0/1/3 才推[09-19 首验勘误:rc() 把「墙停+噪声 failed」报 1,旧 0/3 闸门 11 天只推 1 天,仅 rc=2 限流不推]、推送失败只记日志不进聚合 rc;**admin 账号**=避开 root 属主坑;原独立 16:00 DSM 方案作废),分区粒度 journal 指纹 diff 只增不删;部署位 `/volume1/stock/vps-mirror/`(宿主文件,改 repo 副本后 scp -O 同步,同 corpus/dmsk runner 纪律),journal/watermark_*.json/push.log 同目录;VPS 落 `data\corpus`+`data\fundamentals\dmsk`+族根 `_mirror_watermark.json`(消费方 `get_watermark()` 晨检对表);任一失败不 commit journal/不落水位线=次日 diff 自愈;`--dry-run` 打印将推送清单=验收收敛证明。设计与消费契约=融合 spec §20.11。09-20 增**独立 static manifest lane**(spec §20.11.7):`data_backup/static_overlay/manifests/*.list`(补洗变更清单,每行 `<domain>/<file>`)文件级幂等推送、整单成功归档 pushed/、失败只 WARN 不影响 corpus/dmsk 聚合——三表 per-stock 文件由此 lane 推 VPS(不进 journal 指纹体系)。
- **flash 快讯 lane+xcheck 对账 lane(融合 spec §20.12.3/§20.12.4,09-22 上线)**:①`run_corpus_standalone.sh` 第三 lane `--lane flash`(金十快照+东财 kuaixun 翻页 T-1 边界,秒级完赛,rc 隔离不进聚合、独立跑不 gate 于 daily/backfill)→ corpus/flash_meta/ 落盘,镜像车五子域自动携带;②`run_corpus_standalone.sh` 第四 lane `xcheck.py --lane corpus_side`(ann_em 拉取+公告披露日对照+腾讯 qt 全市场估值快照对照,rc 隔离);③`run_dmsk.sh` 第四 lane `xcheck.py --lane fund_side`(读 wash manifests 当日命中股→新浪三表近 12 期+两源 NOTICE_DATE 对照)。xcheck 域落 `corpus/xcheck/` **不推 VPS**(NAS 侧质检设施);findings 报告+trace 追溯=融合 spec §20.12.4;容器名 sanguo-xcheck(杀=`docker stop sanguo-xcheck`)。首跑锚点:corpus cron.log `flash lane exit=0(隔离)` 行 + `xcheck corpus_side exit=0` 行 + dmsk cron.log `xcheck fund_side exit=0` 行 + 次日 push.log 出现 flash_meta 分区计数。
- **funnel 漏斗 lane+integrity gate lane(融合 spec §20.16/流水线 spec §4.8 P4-3,10-02 上线)**:`run_corpus_standalone.sh` 第五/六段——⑤`corpus_funnel.py --until 23:00 --llm-cap 200`(第一层词表八类事件零成本全量→events_llm 直落;第二层 LLM 只碰词表 miss:候选入 `state/funnel_pending.jsonl` 队列、串行慢爬、token 预算断路 rc=3 次夜续;水位 `state/funnel_watermark.json` 只前向,历史回填 `--start` 按预算分批)⑥`integrity_gate.py --profile default`(三层门 crit exit1/warn exit2)。两段 rc 隔离不进聚合。**LLM env 三件与 P4-1 向导同款(§运维表 288 行)**:DSM 任务环境或 run_corpus.sh 头 export 后由 wrapper 条件传递进容器;缺 key 且当晚有候选 → funnel rc=1(fail loud,非故障面=未配 env)。冒烟:corpus cron.log `funnel exit=0` 行+`integrity gate exit=0` 行+`corpus/events_llm/dt=*/` 分区出现;`--no-llm` 可只验第一层。容器名 sanguo-funnel/sanguo-igate(杀=docker stop)。
- **因子月度批评·先批后评四段链(流水线 spec §4.2 决议 H+§12.4 决议 J,09-23 上线;09-27 加 stage4 归因批 10-04 首班起)**:双机同评——NAS+VPS 每月第一个周日 09:30 各跑同一批评链(①②③段),JSON 对拍件(`{host}_{asof}.json`)落各自报告目录,人工/VPS 侧 scp 汇集到 NAS 后 `python -m sanguo_factor.review_compare nas_*.json vps_*.json` 对拍(一致=落档进告警流程;分歧 exit 2=只开排查 issue 不出降权告警)。ⓐ**NAS 侧**:`run_factor_monthly.sh`(源码=repo `scripts/nas_sync/run_factor_monthly_standalone.sh`,改后 scp -O 同步生效位 `/volume1/stock/sanguo_vnpy_v2/`)——一次性容器四段串行 `monthly_batch`(注册表全量 12M 窗重评批,重活)→`monthly_review --host nas`(判定报告,JSON append-only)→`data_gap_check`→`attribution_batch`(归因批: 15 截面增量[12 源+size_log+在位者+挑战者,daily_section 幂等 merge,as_of 撞周末模块内归一最近 bar 日]→归因双跑[#69 影子 A/B 归因侧供给,挑战者落 challenger/ 子目录]→cp 双投递生效位[在位者件→admin 家 `data/attribution/` 根 glob 取最新;挑战者件带因子名;生效位目录容器挂 rw `/effective-data`,代码卷恒 :ro];工作区 `/volume1/stock/factor_cross_section/attribution/` 留全历史);as_of=上月末脚本头自算;rc 语义=①非零停链、②③ exit 2 是有效信号按 0 聚合、④非零=真失败(截面断供/归因挂/投递缺)进 final;容器金丝雀坑=**必须 `-e HOME=/tmp -e MPLCONFIGDIR=/tmp/mpl`**(--user 1024 下 HOME=/ → vnpy import 建不了 `/.vntrader`)。报告目录 `/volume1/stock/sanguo_vnpy_v2/reports/factor_monthly/`;日志=factor_monthly_cron.log。DSM 任务 `sanguo-factor-monthly`=**每周日 09:30 admin(UI 手建)+脚本头日号闸**——DSM 计划任务月度档只有「每月几号」没有「第一个周日」原生选项,调度只能建每周日,「是不是第一个周日」由脚本头 `dom>7 静默退出` 兜住。ⓑ**VPS 侧**:wrapper=`scripts/factor_research/factor_monthly_wrapper.ps1`(三段同链批评部分,`--host vps`,eval_db=`C:\sanguo_vnpy_v2\data\factor_eval.db`,报告/日志=`C:\sanguo_vnpy_v2\reports\factor_monthly\`;ASCII-only 纪律;**stage4 归因批不镜像**——归因双跑落 NAS 工作区+生效位,VPS 侧归因件走 infra 偏差日报链另行供给);schtask 注册=`scripts/factor_research/register_factor_monthly_task.ps1`(**XML ScheduleByMonthDayOfWeek 原生「每月第一个周日」**,无需日号闸——zh-CN schtasks 不吃星期名组合走 ak-monthly 同款 XML 路线;SYSTEM/不限时,monthly_batch 重活首班实测耗时后再定是否拆任务错时)——vps-deploy 把两文件推上 VPS 后 ssh 跑一次注册脚本即成。首班=部署后第一个周日人工盯全程。
- **static 三表事件驱动补洗(融合 spec §20.11.7,2026-09-20 上线)**:`run_dmsk.sh` 第三 lane `static_overlay_wash.py`(dmsk daily/backfill 后 ~07:2x 起跑,--until 15:00 墙停,rc 隔离)——dmsk NOTICE_DATE 雷达命中股逐股 F10(emweb) 近 12 期 upsert 重写 `static/{income,balance,cashflow}/` per-stock 文件;pending 重试队列(F10 落库滞后/网络失败次日自愈,30 天老化);周日 sync_parquet 全量回冲=雷达窗内自愈、更深归月度全量洗兜底;变更清单→晚间镜像车 manifest lane 推 VPS。
- 数据归属口径四原则与 14 条管线全景=融合 spec §20.10(数据面唯一权威查询入口)。
@@ -1824,6 +1824,21 @@ schema 特殊低估交集教训);限速纪律沿 §19.9 采集纪律表;VP
①alerts 存储+API(sanguo_api 路由)→ ②VPS 监控脚本+告警接入(注册表从本节表导出+`akshare_static_download.py` TYPES 交叉核对防漏登记)→ ③前端告警页 → ④NAS 监控 v2(同 schema 推 VPS;含 NAS 侧特有项:corpus 双 lane/vps-mirror 水位线/03:00 sync 链/xcheck findings)。数据面监控+事件模型=data session 域;策略/因子告警生产方=各自 session 按 schema 接入。
### 20.16 增补 2026-10-02:漏斗分诊+事件快照+完整性门(P4-3 收官刀,流水线 spec §4.8 F1-F7 的数据线细节版)
**产物=corpus 族两新城**(append-only dt 分区=采集日,永不重算;`cd.ID_KEYS` 两行):
- `events_in`:LLM 候选输入切片(content_hash 内容寻址+src 域/id+event_date+title+text_snippet≤500 字)——防编数快照:喂模型的确切文本先落盘,可回放对拍;词表件不入(确定性层零成本可重放)。
- `events_llm`:事件判读统一域(event_id=hash(content_hash+type+code+direction)/content_hash 外键/event_type 八类枚举/direction/confidence/model=lexicon 或模型名)。
**采集件**(`scripts/data_platform/`):
- `funnel_lexicon.py`:第一层词表(八类+否定前缀守卫,sentiment_lexicon 范式;分类学扩充=词表加行)。
- `corpus_funnel.py`:夜批(自有 flock `state/funnel.lock`);水位 `state/funnel_watermark.json`(news_meta/ann_meta 各自只前向);LLM 候选持久队列 `state/funnel_pending.jsonl`(预算停摆残余次夜续,marker=content_hash);LLM 契约=stock_code/event_date 系统侧注入+确定性后检(单项违规丢弃计 parse_fail);Ctx 复用(token 预算=rc3 checkpoint;`--limit`=LLM 调用数;`--llm-cap` 每夜新鲜候选 200);卡片需求面(`data_check` 态 dataNeeds)记录进 stats 不驱动;stats 行 `[funnel] 完成 统计: {...}` 兼容监控过程层解析。
- `integrity_gate.py`:三层门(L1 profile 域存在性——新鲜度归监控不归门;L2 字段四元组 12 项 6 critical 行级有效率≥0.95;L3 充实度空率≤50%);profile-aware 分母(default=corpus 五域/events=两新城);exit 1(crit)/2(warn)/0;`SANGUO_SKIP_INTEGRITY=1` 唯一后门;单查隔离。
**挂点**:`run_corpus_standalone.sh` xcheck 段后两段(funnel→gate,rc 隔离不进聚合;LLM env 三件+SANGUO_PIPELINE_DB 条件传递,缺 key 且有候选→rc=1 fail loud)。**读出口**:`corpus_reader.get_events(start,end,models)`(照 get_flash 容错契约);manifest 已登记两域。
**已知语义**:词表层 v1 接受残余噪声(否定前缀守卫拦大头)换确定性;LLM abstain=空数组合法;pending 积压=预算停摆的显式信号(stats pending_after 可观测)。
## 参考(调查来源)
- xtdata 官方:https://dict.thinktrader.net/nativeApi/xtdata.html
+17
View File
@@ -66,6 +66,23 @@ def get_flash(start=None, end=None, sources=None):
return df
EVENTS_COLUMNS = ["event_id", "content_hash", "src_domain", "src_id",
"stock_code", "event_date", "event_type", "direction",
"confidence", "model", "extracted_at"]
def get_events(start=None, end=None, models=None):
"""事件判读快照(spec §20.16 P4-3,dt 分区直读)。models 过滤
({"lexicon"}=只看词表层,或模型名);空树返带 schema 的空 df
(容错分级契约同 get_flash)。"""
df = pr.read_partitions(pr.corpus_root(), "events_llm", start, end)
if df.empty:
return df.reindex(columns=EVENTS_COLUMNS)
if models:
df = df[df["model"].isin(set(models))]
return df
def get_watermark():
"""镜像水位线(None=尚无成功推送记录)。"""
return pr.read_watermark(pr.corpus_root())
+29
View File
@@ -93,6 +93,35 @@ GID_ADMINS=101 # administrators(Synology ACL 授权组,金丝雀实测必带)
/app/scripts/data_platform/xcheck.py --lane corpus_side
rc_xc=$?
echo "=== $(date '+%F %T') xcheck corpus_side exit=$rc_xc(隔离,不进聚合) ==="
# ===== funnel lane(spec §4.8 P4-3, 2026-10-02: 两层漏斗+事件快照) =====
# 第一层词表零成本全量; 第二层 LLM 只碰词表 miss(串行慢爬+token 预算断路
# rc=3 checkpoint 次夜续)。rc 隔离不进聚合; LLM env 三件条件传递(缺 key 且
# 有候选 → rc=1 fail loud, env 配方=runbook P4-1 节)。
_LLM_ENV=()
[ -n "${SANGUO_LLM_API_KEY:-}" ] && _LLM_ENV+=(-e SANGUO_LLM_API_KEY="$SANGUO_LLM_API_KEY")
[ -n "${SANGUO_LLM_BASE_URL:-}" ] && _LLM_ENV+=(-e SANGUO_LLM_BASE_URL="$SANGUO_LLM_BASE_URL")
[ -n "${SANGUO_LLM_MODEL:-}" ] && _LLM_ENV+=(-e SANGUO_LLM_MODEL="$SANGUO_LLM_MODEL")
[ -n "${SANGUO_PIPELINE_DB:-}" ] && _LLM_ENV+=(-e SANGUO_PIPELINE_DB="$SANGUO_PIPELINE_DB")
"$DOCKER" run --rm --name sanguo-funnel --user "${UID_ADMIN}:${GID_ADMIN}" \
--group-add "${GID_ADMINS}" --no-healthcheck --entrypoint python \
"${_LLM_ENV[@]+"${_LLM_ENV[@]}"}" \
-v /volume1/stock:/volume1/stock \
-v "$APP":/app:ro \
sanguo_vnpy_v2:lock-aligned \
/app/scripts/data_platform/corpus_funnel.py --until 23:00 --llm-cap 200
rc_funnel=$?
echo "=== $(date '+%F %T') funnel exit=$rc_funnel(隔离,不进聚合) ==="
# ===== integrity gate(P4-3 前置纪律②, spec §20.16; rc 隔离) =====
# crit>0 exit 1 / warn exit 2 / 绿 exit 0; 报告 stdout 进 cron.log。
"$DOCKER" run --rm --name sanguo-igate --user "${UID_ADMIN}:${GID_ADMIN}" \
--group-add "${GID_ADMINS}" --no-healthcheck --entrypoint python \
-v /volume1/stock:/volume1/stock \
-v "$APP":/app:ro \
sanguo_vnpy_v2:lock-aligned \
/app/scripts/data_platform/integrity_gate.py --root /volume1/stock/corpus \
--profile default
rc_gate=$?
echo "=== $(date '+%F %T') integrity gate exit=$rc_gate(隔离,不进聚合) ==="
} >> "$LOG" 2>&1
# ===== hot_rank NAS 采集+推 VPS(spec §20.14, 09-25: emappdata 云 IP 断连迁移) =====
# 独立 runner(采集容器+scp 推送+补推未推日期), 自身恒 exit 0 隔离, 失败只记
+31
View File
@@ -17,6 +17,15 @@ from scripts.data_platform import monitor_checks
from sanguo_data import partition_reader as pr
@pytest.fixture(autouse=True)
def _clear_day_buffers():
"""cd._DAY_BUFFERS 是模块级全局——未 commit 的测试残余会被下一测试的
commit() 冲刷进新 tmp 根(跨测试泄漏实证)。逐测试清空。"""
cd._DAY_BUFFERS.clear()
yield
cd._DAY_BUFFERS.clear()
@pytest.fixture
def root(tmp_path, monkeypatch):
monkeypatch.setattr(cd, "CORPUS_ROOT", tmp_path)
@@ -237,3 +246,25 @@ def test_llm_lane_skips_done_markers(root):
fake = FakeClient({"events": []})
n = asyncio.run(cf._llm_lane(ctx, fake, "m", cands, led, rec))
assert n == 0 and fake.calls == []
# ---------- 读出口 get_events(F7, 照 get_flash 范式) ----------
def test_get_events_reader(root, monkeypatch):
day = _day(1)
_seed_day(root, "news_meta", day, _news_rows(day, [
("A1", "600519", "控股股东拟增持公司股份公告", None)]))
cf.main(["--no-llm"])
monkeypatch.setenv("SANGUO_CORPUS_ROOT", str(root)) # reader 走 env 根
from sanguo_data import corpus_reader as cr
df = cr.get_events()
assert not df.empty and {"event_id", "event_type", "model"} <= set(df.columns)
assert cr.get_events(models={"lexicon"}).shape[0] == df.shape[0]
assert cr.get_events(models={"nope"}).empty
def test_get_events_empty_tree_schema(tmp_path, monkeypatch):
monkeypatch.setattr(cd, "CORPUS_ROOT", tmp_path)
from sanguo_data import corpus_reader as cr
df = cr.get_events()
assert df.empty and list(df.columns) == cr.EVENTS_COLUMNS