From 236bd30eb8d5fd8f732a1a582e8e67487e507f4b Mon Sep 17 00:00:00 2001 From: claude_dev Date: Tue, 8 Sep 2026 00:29:23 +0800 Subject: [PATCH] =?UTF-8?q?feat(data):=20=E8=B4=A2=E5=8A=A1136=20W3'=20dms?= =?UTF-8?q?k=E5=8F=8Clane=E8=90=BD=E5=9C=B0+spec=C2=A719.8=E5=AE=9E?= =?UTF-8?q?=E6=96=BD=E8=AE=B0=E5=BD=95=20[nas]?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - fundamentals_dmsk_download.py 单脚本(2009→今backfill+近3期daily NOTIC_DATE>=T-7),corpus骨架六修正全带入 - marker修正: daily按日折叠(§19.4原文「一次永逸」=day-2同款坑,52756e3教训);backfill键=域+期(域缺席连带跳过,首测抓出) - 次新期金丝雀防「filter失效vs空期」不分导致的整段静默空跑 - run_dmsk_standalone.sh(DSM 07:10 sanguo-fundamentals-daily,一次性容器同corpus款)+verify_dmsk.py(键唯一硬门) - 21测全mock绿,全量295绿;W1'镜像对账全绿+vb2026 H1缺口入档 --- .../2026-07-21-data-source-fusion-design.md | 25 + .../fundamentals_dmsk_download.py | 697 ++++++++++++++++++ scripts/data_platform/verify_dmsk.py | 57 ++ scripts/nas_sync/run_dmsk_standalone.sh | 73 ++ .../test_fundamentals_dmsk_download.py | 326 ++++++++ 5 files changed, 1178 insertions(+) create mode 100644 scripts/data_platform/fundamentals_dmsk_download.py create mode 100644 scripts/data_platform/verify_dmsk.py create mode 100755 scripts/nas_sync/run_dmsk_standalone.sh create mode 100644 tests/data_platform/test_fundamentals_dmsk_download.py diff --git a/docs/superpowers/specs/2026-07-21-data-source-fusion-design.md b/docs/superpowers/specs/2026-07-21-data-source-fusion-design.md index e265b88..64967ce 100644 --- a/docs/superpowers/specs/2026-07-21-data-source-fusion-design.md +++ b/docs/superpowers/specs/2026-07-21-data-source-fusion-design.md @@ -818,6 +818,31 @@ valuation 今日刚同步/三表 09-06 全量),七域覆盖≡VPS(沪深 5215 缺口结论不变: 沪深退市 355+北交退市 258 三表缺(其日频估值已被 baostock 树覆盖)+ 北交现役 340 空+forecast/express 2020 前可选。 +### 19.8 实施记录(2026-09-08 凌晨) + +- **W1' 镜像对账全绿**: 七域 count(三表各 5555=5215 沪深有效+340 北交 920xxx + 零行空档,挂 .SZ 后缀=09-01 空 dump 残留;abstract/valuation 5557;forecast/ + express 26 期)+600519 四锚点容器内复验逐位一致+读取 0.99s。第 8 域 + valuation_baostock: 1990→2025 完整(2025 实测 1.25M 行/243 天/5211 股,含退市); + **2026 文件仅 08-13 起**(bs 日喂起点)→ 2026 H1 估值缺口由基线层 valuation + 覆盖,H1 回补列可选 backlog。 +- **W3' dmsk 双 lane 落地**: `fundamentals_dmsk_download.py`(单脚本红线,corpus + 骨架复用+六修正全带入)+21 测全 mock 绿+`run_dmsk_standalone.sh`(DSM + `sanguo-fundamentals-daily` 07:10,一次性容器 sanguo-fund-dmsk,corpus 同款 + --user 1024:100 --group-add 101)+`verify_dmsk.py`(键唯一硬门只读验收)。 + 落库=/volume1/stock/fundamentals/dmsk/{域}/dt=期/part-*.parquet,**绝不写镜像树**。 +- **§19.4 marker 语义修正**: 原文「daily 增量 unit=REPORT_DATE 一次永逸按期 + marker」有 day-2 同款坑(次日重扫被跳过漏新披露行,corpus 52756e3 教训)——实施 + 改为 daily marker 按日折叠(dt=子目录+日初清理);backfill 期 marker 才是一次 + 永逸(期内容历史不可变)。marker 键=域+期(域缺席会让后续域被首域 marker 连带 + 跳过,首测抓出)。 +- **次新期金丝雀**: 东财把「filter 契约失效」与「合法空期」都回 result=None + 无法区分→次新期(全市场已披露 3+ 月)空=契约漂移按失败处理不标 done,防整段 + 静默空跑标 done。 +- **profit/dupont 断更根因初判**: fundamentals_baostock/backfill_state.json + 一次性回补态、从未接日喂(与 vb 按年文件有 bs-daily 日喂形成对照);修法=季度 + 追赶一次性回补(~10 季×5200 股≈5.2 万发,分两晚),排 dmsk 之后,待用户点头。 + ## 参考(调查来源) - xtdata 官方:https://dict.thinktrader.net/nativeApi/xtdata.html diff --git a/scripts/data_platform/fundamentals_dmsk_download.py b/scripts/data_platform/fundamentals_dmsk_download.py new file mode 100644 index 0000000..f08f5ae --- /dev/null +++ b/scripts/data_platform/fundamentals_dmsk_download.py @@ -0,0 +1,697 @@ +# -*- coding: utf-8 -*- +"""fundamentals_dmsk_download.py — NAS 财务三表 dmsk 补缺层采集唯一入口(spec §19.4)。 + +运行形态: 一次性容器 docker run --rm(wrapper=/volume1/stock/fundamentals/run_dmsk.sh), +与 corpus/5m 同款;落 datacenter-web(东财数据中心),与 corpus 的 search-api/np-listapi +分域,各自限速。 + --lane backfill 2009Q1→最近已完期 全期回补(unit=(domain,REPORT_DATE) 一次永逸, + 期内容历史不可变;≈70期×3表×~12页≈2200发≈1h) + --lane daily 近 3 期重扫: NOTICE_DATE>=T-7 客户端过滤+账本去重 → 披露次日即可用 + (不等 VPS 月度洗库 5 周);marker 按日折叠(增量型 marker 必须带日期 + 维度——corpus day-2 坑 52756e3 教训,spec§19.4「一次永逸」原文据此修正) + --until HH:MM 墙钟自停: 完成当前 unit 后 checkpoint 退出(rc=3) + --limit N 冒烟: 每域最多 N unit, 绝不落 marker + +落库: {DMSK_ROOT}/dmsk/{balance,income,cashflow}/dt=YYYY-MM-DD/part-N.parquet + 57 列原样+SECURITY_CODE/REPORT_DATE 归一+fetch_date;行键 (SECURITY_CODE, + REPORT_DATE) 全局唯一(IdLedger);**绝不写 phase3 镜像树**(state/static 归 + phase3 04:00 链所有,写=被次日冲掉+破坏 ≡VPS 不变量)。 + +下游读法(双层口径, spec §19.7 写死): 基线层=镜像树全列(≤5周延迟)为主;本层=57列 +核心科目新鲜度补丁+退市/北交覆盖;键唯一由本层账本保证,读者可直 union。 + +断点续传: unit marker(state/markers/, tmp+rename 原子)。 +退出码: 0=完成 1=致命(unit 失败/域冷却) 2=限流让路(429) 3=墙钟/锁让路。 +单实例锁: fcntl.flock(state/dmsk.lock) — 跨容器内核级,进程死自动释放。 +""" +import argparse +import datetime as dt +import fcntl +import hashlib +import json +import logging +import os +import random +import shutil +import sys +import time +from pathlib import Path + +import numpy as np +import pandas as pd +import requests + +# ---------- 常量(测试可 monkeypatch) ---------- + +DMSK_ROOT = Path(os.environ.get("DMSK_ROOT", "/volume1/stock/fundamentals")) + +DC_URL = "https://datacenter-web.eastmoney.com/api/data/v1/get" + +DOMAINS = { # spec §19.2 实测契约: 全市场单期≈5223行/57列 + "balance": "RPT_DMSK_FN_BALANCE", + "income": "RPT_DMSK_FN_INCOME", + "cashflow": "RPT_DMSK_FN_CASHFLOW", +} +ID_KEYS = {d: ["SECURITY_CODE", "REPORT_DATE"] for d in DOMAINS} + +DC_PAGE_SIZE = 500 # 实测 5223 行≈12 页 +DC_MAX_PAGES = 500 # 防失控硬顶 +BACKFILL_FROM_YEAR = 2009 # 实测最早 2009 年报(akshare docstring「2010起」过时) +DAILY_RECENT_PERIODS = 3 # 近 3 期重扫窗 +DAILY_NOTICE_WINDOW_DAYS = 7 # NOTICE_DATE >= T-7 视为新披露 +BJ_PREFIXES = ("43", "83", "87", "88", "92") # 北交(920 新码+43/83/87/88 存量) + +COOLDOWN_ERRORS = 10 +FLUSH_EVERY = 10 # 账本/marker 批量提交周期(unit 计) +RATE = {"datacenter": 1.5} # 每域最小间隔(秒), 慢爬纪律 +JITTER = (0.0, 0.3) +UA = {"User-Agent": "Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36"} + +log = logging.getLogger("dmsk") + + +# ---------- 异常(契约同 corpus) ---------- + +class TransportError(Exception): + """传输类失败(超时/5xx/连接): unit 不标 done, 计入域名连续错。""" + + +class DomainCooldown(Exception): + """域名冷却: 连续错×10 或 429。rate_limited=True → 退出码走 2(限流让路)。""" + + def __init__(self, domain, rate_limited=False): + super().__init__(f"domain {domain} cooled (rate_limited={rate_limited})") + self.rate_limited = rate_limited + + +class HttpDeterministicError(Exception): + """确定性 4xx(非 429): 记日志跳过, 绝不烧重试。""" + + def __init__(self, status, url=""): + super().__init__(f"HTTP {status} {url}") + self.status = status + + +class WallClockStop(Exception): + """--until 墙钟到: 完成当前 unit 后抛出, 顶层收拾落盘并 rc=3。""" + + +# ---------- 单实例锁(flock, 跨容器) ---------- + +def acquire_lock(): + lock = DMSK_ROOT / "state" / "dmsk.lock" + lock.parent.mkdir(parents=True, exist_ok=True) + fd = os.open(str(lock), os.O_CREAT | os.O_RDWR, 0o666) + try: + fcntl.flock(fd, fcntl.LOCK_EX | fcntl.LOCK_NB) + except BlockingIOError: + os.close(fd) + return None + return fd + + +# ---------- 域名客户端(限速/连续错/429, 同 corpus) ---------- + +class DomainClient: + def __init__(self, name): + self.name = name + self.session = requests.Session() + self.session.trust_env = False # 直连不走代理 + self.consecutive_errors = 0 + self.rate_limited = False + self.cooled = False + self.n_requests = 0 + self._last = 0.0 + self._t_first = None + self._t_last = None + + def _pace(self): + interval = RATE.get(self.name, 1.0) + random.uniform(*JITTER) + wait = self._last + interval - time.monotonic() + if wait > 0: + time.sleep(wait) + self._last = time.monotonic() + if self._t_first is None: + self._t_first = self._last + self._t_last = self._last + self.n_requests += 1 + + def avg_interval(self): + if self.n_requests < 2 or self._t_first is None: + return None + return (self._t_last - self._t_first) / (self.n_requests - 1) + + def _transport_fail(self, msg): + self.consecutive_errors += 1 + if self.consecutive_errors >= COOLDOWN_ERRORS: + self.cooled = True + raise DomainCooldown(self.name) + raise TransportError(msg) + + def request(self, method, url, **kw): + if self.cooled: + raise DomainCooldown(self.name) + self._pace() + try: + r = self.session.request(method, url, timeout=25, **kw) + except requests.RequestException as e: + self._transport_fail(f"{type(e).__name__}: {e}") + if r.status_code == 429: + self.rate_limited = True + self.cooled = True + raise DomainCooldown(self.name, rate_limited=True) + if r.status_code >= 500: + self._transport_fail(f"HTTP {r.status_code}") + if r.status_code >= 400: + raise HttpDeterministicError(r.status_code, url) + self.consecutive_errors = 0 + return r + + +# ---------- 期枚举(纯函数) ---------- + +def quarter_end_dates(start_year, before_date, count=None): + """季度期末日序列, 新→旧(近年优先供 factor);只含 <= before_date 的已完期。 + 例: before=2026-09-08 → 2026-06-30, 2026-03-31, 2025-12-31, ...""" + out = [] + for year in range(before_date.year, start_year - 1, -1): + for md in ("12-31", "09-30", "06-30", "03-31"): + d = dt.date(year, *(int(x) for x in md.split("-"))) + if d <= before_date: + out.append(d) + if count is not None: + out = out[:count] + return out + + +def is_bj_code(code): + return bool(code) and str(code)[:2] in BJ_PREFIXES + + +# ---------- 归一化(纯函数) ---------- + +def norm_row(raw, fetch_date): + """行键列归一(SECURITY_CODE 补零串/REPORT_DATE 截日),其余 57 列原样透传。""" + row = dict(raw) + code = str(raw.get("SECURITY_CODE") or "").strip() + if code.isdigit() and len(code) < 6: + code = code.zfill(6) + row["SECURITY_CODE"] = code + row["REPORT_DATE"] = str(raw.get("REPORT_DATE") or "")[:10] + row["fetch_date"] = fetch_date + return row + + +# ---------- id 去重账本(同 corpus IdLedger) ---------- + +def _hash_id(s): + return int.from_bytes(hashlib.blake2b(s.encode("utf-8"), + digest_size=8).digest(), "big") & (2**63 - 1) + + +def _row_key(domain, row): + return "|".join(str(row.get(c) or "") for c in ID_KEYS[domain]) + + +class IdLedger: + """全局行键账本(state/ids_{domain}.parquet, int64 hash 有序数组)。""" + + def __init__(self, domain): + self.path = DMSK_ROOT / "state" / f"ids_{domain}.parquet" + if self.path.exists(): + self._arr = pd.read_parquet(self.path)["h"].to_numpy() + else: + self._arr = np.empty(0, dtype=np.int64) + self._pending = [] + + def add(self, ids): + self._pending.extend(_hash_id(i) for i in ids) + + def has_any(self, ids): + hs = np.array([_hash_id(i) for i in ids], dtype=np.int64) + out = [] + for h in hs: + i = np.searchsorted(self._arr, h) + out.append(bool(i < len(self._arr) and self._arr[i] == h) + or h in self._pending) + return out + + def flush(self): + if not self._pending: + return len(self._arr) + merged = np.unique(np.concatenate( + [self._arr, np.array(self._pending, dtype=np.int64)])) + tmp = self.path.with_suffix(".tmp") + pd.DataFrame({"h": merged}).to_parquet(tmp, index=False) + os.replace(tmp, self.path) + self._arr = merged + self._pending = [] + return len(self._arr) + + +def _iter_parts(domain): + base = DMSK_ROOT / "dmsk" / domain + if not base.exists(): + return [] + return sorted(p for p in base.glob("dt=*/part-*.parquet") if p.is_file()) + + +def reconcile_ledger_from_parts(domain, ledger): + """启动对账: 分区文件即真相——part 有而账本无的键补进账本(崩溃窗口自愈, + corpus _reconcile_ledgers 的全量版: dmsk 分区=期目录非日历日, 无近窗概念)。""" + keys = [] + for part in _iter_parts(domain): + df = pd.read_parquet(part, columns=ID_KEYS[domain]) + keys.extend("|".join(str(r[c]) for c in ID_KEYS[domain]) + for r in df.to_dict("records")) + if not keys: + return 0 + known = ledger.has_any(keys) + missing = [k for k, has in zip(keys, known) if not has] + if missing: + ledger.add(missing) + ledger.flush() + log.info("ledger reconcile %s: +%d (分区对账)", domain, len(missing)) + return len(missing) + + +# ---------- 分区追加(按期 part 文件, O(1)/unit) ---------- + +_PERIOD_BUFFERS = {} + + +def _period_dir(domain, report_date): + return DMSK_ROOT / "dmsk" / domain / f"dt={report_date}" + + +def append_rows(domain, report_date, rows): + """进 (domain, 期) 内存缓冲; flush_all() 时各写一个新 part-N(tmp+rename)。 + 旧 part 永不重写;行键唯一由调用方(账本+批内 seen)保证。""" + if not rows: + return 0 + key = (domain, report_date) + buf = _PERIOD_BUFFERS.get(key) + if buf is None: + buf = _PERIOD_BUFFERS[key] = PeriodBuffer(domain, report_date) + buf.rows.extend(rows) + return len(rows) + + +class PeriodBuffer: + """期目录追加缓冲 → flush 落一个新 part 文件(批量, 原子)。""" + + def __init__(self, domain, report_date): + self.domain = domain + self.report_date = report_date + self.rows = [] + d = _period_dir(domain, report_date) + idx = -1 + if d.exists(): + for p in d.glob("part-*.parquet"): + try: + idx = max(idx, int(p.stem.split("-")[1])) + except (ValueError, IndexError): + continue + self.next_idx = idx + 1 + + def flush(self): + if not self.rows: + return 0 + d = _period_dir(self.domain, self.report_date) + d.mkdir(parents=True, exist_ok=True) + n = len(self.rows) + tmp = d / f".part-{self.next_idx}.{os.getpid()}.tmp" + try: + pd.DataFrame(self.rows).to_parquet(tmp, index=False) + os.replace(tmp, d / f"part-{self.next_idx}.parquet") + except Exception: + tmp.unlink(missing_ok=True) + raise + self.next_idx += 1 + self.rows = [] + return n + + +def flush_all(): + for buf in _PERIOD_BUFFERS.values(): + buf.flush() + + +# ---------- unit marker ---------- + +def _marker_path(lane, stage, unit): + return DMSK_ROOT / "state" / "markers" / lane / stage / f"{unit}.done" + + +def is_done(lane, stage, unit): + return _marker_path(lane, stage, unit).exists() + + +def mark_done(lane, stage, unit): + p = _marker_path(lane, stage, unit) + p.parent.mkdir(parents=True, exist_ok=True) + tmp = p.with_suffix(".tmp") + tmp.write_text(dt.datetime.now().isoformat(), encoding="utf-8") + os.replace(tmp, p) + + +def _dunit(domain, report_date, day): + """daily lane marker 单元名: 按日折叠(dt=YYYY-MM-DD/ 子目录)。 + 增量窗每日一新——无日期维度会让首日 done 后次日整段跳过(corpus day-2 坑)。 + backfill 期 marker 保持平铺一次永逸(期内容历史不可变,语义不同)。""" + return f"dt={day.isoformat()}/{domain}_{report_date}" + + +def prune_daily_markers(): + """daily 日初清理: 只留今日 dt= 目录 + 清历史平铺 .done(防线同 corpus)。""" + today = dt.date.today().isoformat() + base = DMSK_ROOT / "state" / "markers" / "daily" + if not base.exists(): + return + for stage_dir in base.iterdir(): + if not stage_dir.is_dir(): + continue + for child in stage_dir.iterdir(): + if child.is_dir(): + if child.name != f"dt={today}": + shutil.rmtree(child, ignore_errors=True) + elif child.suffix == ".done": + child.unlink(missing_ok=True) + + +# ---------- fetch 层 ---------- + +def fetch_period_rows(client, domain, report_date): + """全市场单期全页。result=None/缺省=合法空期(2009Q1 等早期无数据), + 返回 []由调用方标 done;翻页按 result.pages。""" + out, page = [], 1 + while True: + params = { + "sortColumns": "SECURITY_CODE", "sortTypes": "1", + "pageSize": str(DC_PAGE_SIZE), "pageNum": str(page), + "reportName": DOMAINS[domain], "columns": "ALL", + "filter": f"(REPORT_DATE='{report_date}')", + } + r = client.request("GET", DC_URL, params=params, headers=UA) + try: + j = r.json() + except ValueError as e: + raise TransportError(f"json decode: {e}") from e + result = j.get("result") or {} + rows = result.get("data") or [] + out.extend(rows) + pages = int(result.get("pages") or 0) + if not rows or page >= pages or page >= DC_MAX_PAGES: + return out + page += 1 + + +# ---------- 运行上下文(同 corpus 纪律: 数据→账本→marker 次序) ---------- + +class Ctx: + def __init__(self, lane, until=None, limit=None): + self.lane = lane + self.limit = limit + self.deadline = None + if until: + hh, mm = until.split(":") + self.deadline = dt.datetime.combine( + dt.date.today(), dt.time(int(hh), int(mm))) + self.failed = 0 + self.rate_limited = False + self.hard_cool = False + self.wallclock = False + self.units = 0 + self.stage_units = 0 + self.clients = {} + self.ledgers = None + self.pending_marks = [] + + def set_stores(self, ledgers): + self.ledgers = ledgers + + def reset_stage(self): + """--limit 预算按域独立(冒烟覆盖每域)。""" + self.stage_units = 0 + + def budget_exhausted(self): + return self.limit is not None and self.stage_units >= self.limit + + def unit_done(self, stage, unit, mark=True): + self.units += 1 + self.stage_units += 1 + if mark and self.limit is None: + self.pending_marks.append((self.lane, stage, unit)) + if self.units % FLUSH_EVERY == 0: + self.commit() + + def commit(self): + """数据分区 → id 账本 → unit marker 依次落盘(崩溃宁可重拉不产生洞)。""" + flush_all() + if self.ledgers: + for led in self.ledgers.values(): + led.flush() + for lane, stage, unit in self.pending_marks: + mark_done(lane, stage, unit) + self.pending_marks.clear() + + def stop_now(self): + if self.deadline and dt.datetime.now() >= self.deadline: + self.wallclock = True + raise WallClockStop() + + def rc(self): + if self.failed or self.hard_cool: + return 1 + if self.rate_limited: + return 2 + if self.wallclock: + return 3 + return 0 + + +def _filter_new(domain, norm_rows, ledger): + """账本已知键 + 批内 seen 双层截留(corpus 六修正之三: 批内去重)。""" + keys = [_row_key(domain, r) for r in norm_rows] + known = ledger.has_any(keys) + seen = set() + out = [] + for r, k in zip(norm_rows, known): + if k or _row_key(domain, r) in seen: + continue + seen.add(_row_key(domain, r)) + out.append(r) + return out + + +def _absorb_new(domain, report_date, norm_rows, ledgers): + new = _filter_new(domain, norm_rows, ledgers[domain]) + if new: + append_rows(domain, report_date, new) + ledgers[domain].add([_row_key(domain, r) for r in new]) + return new + + +# ---------- BJ 含性工件(spec §19.4 实施首日验证) ---------- + +def save_bj_inclusion(counts): + """{domain: {period: bj_rows}} 跨 run 累积, 供北交所含性验证。""" + p = DMSK_ROOT / "state" / "bj_inclusion.json" + merged = {} + if p.exists(): + try: + merged = json.loads(p.read_text(encoding="utf-8")) + except ValueError: + log.warning("bj_inclusion.json 损坏, 重写(幂等无损)") + for domain, per_period in counts.items(): + slot = merged.setdefault(domain, {}) + for period, n in per_period.items(): + slot[period] = max(slot.get(period, 0), n) + tmp = p.with_suffix(".tmp") + tmp.write_text(json.dumps(merged, ensure_ascii=False, indent=1, + sort_keys=True), encoding="utf-8") + os.replace(tmp, p) + + +# ---------- backfill lane ---------- + +def _canary_period(periods): + """次新期金丝雀: 东财把「过滤器失效/契约漂移」与「合法空期」都回 result=None + 无法区分——次新期全市场已披露 3+ 个月必非空,它空=契约坏了。最新期不作金丝雀 + (披露季初可合法 0 行)。""" + return periods[1] if len(periods) > 1 else periods[0] + + +def run_backfill(ctx, cl, ledgers, start_year=BACKFILL_FROM_YEAR): + today = dt.date.today() + fetch_date = today.isoformat() + periods = quarter_end_dates(start_year, today) + canary = _canary_period(periods) + bj_counts = {} + log.info("backfill periods: %d (%s → %s), canary=%s", + len(periods), periods[-1], periods[0], canary) + try: + for domain in DOMAINS: + ctx.reset_stage() + bj_counts[domain] = {} + for report_date in periods: + unit = str(report_date) + mkey = f"{domain}_{unit}" # marker 键=域+期(域缺席会让后续 + if ctx.budget_exhausted(): # 域被首域 marker 连带跳过) + break + if is_done(ctx.lane, "period", mkey): + continue + try: + rows = fetch_period_rows(cl, domain, unit) + except DomainCooldown as e: + log.warning("backfill %s 域冷却 @%s: %s", domain, unit, e) + ctx.rate_limited |= e.rate_limited + ctx.hard_cool |= not e.rate_limited + return + except (TransportError, HttpDeterministicError) as e: + ctx.failed += 1 + log.warning("backfill %s %s 失败(不标done): %s", + domain, unit, e) + continue + if not rows and report_date == canary: + ctx.failed += 1 + log.error("backfill %s 金丝雀期 %s 空(契约漂移疑云, " + "不标done不落空库): 检查 filter/URL 契约", + domain, unit) + continue + norm = [norm_row(r, fetch_date) for r in rows] + new = _absorb_new(domain, unit, norm, ledgers) + n_bj = sum(1 for r in new if is_bj_code(r["SECURITY_CODE"])) + bj_counts[domain][unit] = n_bj + log.info("backfill %s %s: +%d/%d (bj=%d)", + domain, unit, len(new), len(rows), n_bj) + ctx.unit_done("period", mkey, mark=ctx.limit is None) + ctx.stop_now() + finally: + if any(bj_counts.values()): + save_bj_inclusion(bj_counts) + + +# ---------- daily lane ---------- + +def run_daily(ctx, cl, ledgers): + today = dt.date.today() + fetch_date = today.isoformat() + cutoff = (today - dt.timedelta(days=DAILY_NOTICE_WINDOW_DAYS)).isoformat() + periods = quarter_end_dates(BACKFILL_FROM_YEAR, today, + count=DAILY_RECENT_PERIODS) + canary = _canary_period(periods) + prune_daily_markers() + log.info("daily periods: %s, notice_cutoff=%s, canary=%s", + periods, cutoff, canary) + for domain in DOMAINS: + ctx.reset_stage() + for report_date in periods: + unit = str(report_date) + du = _dunit(domain, unit, today) + if ctx.limit is None: + if is_done(ctx.lane, "period", du): + continue + elif ctx.budget_exhausted(): + break + try: + rows = fetch_period_rows(cl, domain, unit) + except DomainCooldown as e: + log.warning("daily %s 域冷却 @%s: %s", domain, unit, e) + ctx.rate_limited |= e.rate_limited + ctx.hard_cool |= not e.rate_limited + return + except (TransportError, HttpDeterministicError) as e: + ctx.failed += 1 + log.warning("daily %s %s 失败(不标done): %s", domain, unit, e) + continue + if not rows and report_date == canary: + ctx.failed += 1 + log.error("daily %s 金丝雀期 %s 空(契约漂移疑云, 不标done): " + "检查 filter/URL 契约", domain, unit) + continue + norm = [norm_row(r, fetch_date) for r in rows] + fresh = [r for r in norm + if str(r.get("NOTICE_DATE") or "")[:10] >= cutoff] + new = _absorb_new(domain, unit, fresh, ledgers) + log.info("daily %s %s: +%d 新披露/%d 行", + domain, unit, len(new), len(rows)) + ctx.unit_done("period", du, mark=ctx.limit is None) + ctx.stop_now() + + +# ---------- lane 入口 ---------- + +def _ensure_log(): + if log.handlers: + return + (DMSK_ROOT / "logs").mkdir(parents=True, exist_ok=True) + ts = dt.datetime.now().strftime("%Y%m%d_%H%M%S") + fh = logging.FileHandler(DMSK_ROOT / "logs" / f"dmsk_{ts}.log", + encoding="utf-8") + fh.setFormatter(logging.Formatter("%(asctime)s %(levelname)s %(message)s")) + log.addHandler(fh) + log.addHandler(logging.StreamHandler()) + log.setLevel(logging.INFO) + + +def run_lane(lane, until=None, limit=None, start_year=BACKFILL_FROM_YEAR): + _PERIOD_BUFFERS.clear() + _ensure_log() + ctx = Ctx(lane=lane, until=until, limit=limit) + ledgers = {d: IdLedger(d) for d in DOMAINS} + for d in DOMAINS: + reconcile_ledger_from_parts(d, ledgers[d]) + cl = DomainClient("datacenter") + ctx.clients = {"datacenter": cl} + ctx.set_stores(ledgers) + log.info("lane=%s start: until=%s, limit=%s, start_year=%s", + lane, until, limit, start_year) + t0 = time.monotonic() + try: + if lane == "daily": + run_daily(ctx, cl, ledgers) + elif lane == "backfill": + run_backfill(ctx, cl, ledgers, start_year=start_year) + else: + log.error("unknown lane: %s", lane) + ctx.failed += 1 + except WallClockStop: + log.info("墙钟到(--until %s): 完成当前 unit 后 checkpoint 退出", until) + finally: + ctx.commit() + for name, c in ctx.clients.items(): + if c.n_requests: + avg = c.avg_interval() + log.info("stats %s: n_requests=%d avg_interval=%s%s", + name, c.n_requests, + f"{avg:.2f}s" if avg is not None else "n/a", + " [429]" if c.rate_limited else + " [cooled]" if c.cooled else "") + log.info("lane=%s done in %.0fs: units=%d failed=%d rc=%d", + lane, time.monotonic() - t0, ctx.units, ctx.failed, ctx.rc()) + return ctx.rc() + + +def main(): + ap = argparse.ArgumentParser(description="NAS 财务三表 dmsk 补缺层(spec §19.4)") + ap.add_argument("--lane", required=True, choices=["daily", "backfill"]) + ap.add_argument("--until", default=None, help="HH:MM 墙钟自停") + ap.add_argument("--limit", type=int, default=None, + help="冒烟: 每域最多 N unit") + ap.add_argument("--start-year", type=int, default=BACKFILL_FROM_YEAR) + args = ap.parse_args() + fd = acquire_lock() + if fd is None: + print("另一实例持有锁, 让路退出", file=sys.stderr) + sys.exit(3) + try: + rc = run_lane(args.lane, until=args.until, limit=args.limit, + start_year=args.start_year) + sys.exit(rc) + finally: + os.close(fd) + + +if __name__ == "__main__": + main() diff --git a/scripts/data_platform/verify_dmsk.py b/scripts/data_platform/verify_dmsk.py new file mode 100644 index 0000000..c213ee7 --- /dev/null +++ b/scripts/data_platform/verify_dmsk.py @@ -0,0 +1,57 @@ +# -*- coding: utf-8 -*- +"""verify_dmsk.py — dmsk 补缺层只读验收器(spec §19.5 ①覆盖+键唯一硬门)。 + +用法: python3 verify_dmsk.py [--root /volume1/stock/fundamentals] +检查(退出码 0=PASS 1=FAIL): + 1. 每域期目录数/part 数/行数/北交行数/最新 NOTICE_DATE; + 2. 行键 (SECURITY_CODE,REPORT_DATE) 全局唯一(重键=FAIL, 同 corpus verify 纪律); + 3. 退市股可查性(000003 三表至少一域有行 → 否则 WARNING 不 FAIL, 取决于回补深度); + 4. 北交所含性(920/43/83/87/88 前缀行数>0 → 含; =0 → 单列说明不 FAIL)。 +""" +import argparse +import json +import sys +from pathlib import Path + +import pandas as pd + +BJ_PREFIXES = ("43", "83", "87", "88", "92") + + +def main(): + ap = argparse.ArgumentParser() + ap.add_argument("--root", default="/volume1/stock/fundamentals") + args = ap.parse_args() + root = Path(args.root) / "dmsk" + fails = [] + report = {} + for domain in ("balance", "income", "cashflow"): + base = root / domain + parts = sorted(base.glob("dt=*/part-*.parquet")) if base.exists() else [] + if not parts: + fails.append(f"{domain}: 零 part 文件") + continue + frames = [pd.read_parquet(p, columns=["SECURITY_CODE", "REPORT_DATE", + "NOTICE_DATE"]) for p in parts] + df = pd.concat(frames, ignore_index=True) + dups = int(df.duplicated(["SECURITY_CODE", "REPORT_DATE"]).sum()) + n_bj = int(df["SECURITY_CODE"].astype(str).str[:2] + .isin(BJ_PREFIXES).sum()) + delisted = int((df["SECURITY_CODE"].astype(str) == "000003").sum()) + latest_notice = str(df["NOTICE_DATE"].astype(str).max())[:10] + report[domain] = dict(periods=len({p.parent.name for p in parts}), + parts=len(parts), rows=len(df), dup_keys=dups, + bj_rows=n_bj, delisted_000003_rows=delisted, + latest_notice_date=latest_notice) + if dups: + fails.append(f"{domain}: {dups} 重键行(键唯一硬门)") + print(json.dumps(report, ensure_ascii=False, indent=1)) + if fails: + print("FAIL:", "; ".join(fails)) + sys.exit(1) + print("PASS: 键唯一全过, 三域在库") + sys.exit(0) + + +if __name__ == "__main__": + main() diff --git a/scripts/nas_sync/run_dmsk_standalone.sh b/scripts/nas_sync/run_dmsk_standalone.sh new file mode 100755 index 0000000..b7054c6 --- /dev/null +++ b/scripts/nas_sync/run_dmsk_standalone.sh @@ -0,0 +1,73 @@ +#!/bin/bash +# run_dmsk_standalone.sh — 财务三表 dmsk 补缺层 runner(2026-09-08 spec §19.4) +# +# 模式: 与 corpus/5m 完全同款的「DSM 任务(admin) → bash wrapper → 一次性容器」 +# (docker run --rm 跑 lock-aligned 镜像,主容器重启/重建永不波及本进程)。 +# +# 落库: /volume1/stock/fundamentals/dmsk/{balance,income,cashflow}/ +# dt=/part-*.parquet —— **绝不写 phase3 镜像树** +# (/volume1/stock/sanguo_vnpy_v2/data/static 归 04:00 同步链所有, 写=被冲掉 +# +破坏 ≡VPS 不变量; 双层读法见 spec §19.7)。 +# +# 两 lane 串行(脚本内 flock 单实例锁天然互斥): daily --until 08:00 +# (近3期×3域=9 unit≈110发≈4min; NOTICE_DATE>=T-7 新披露, 披露次日可用); +# rc=0|3 续走 backfill --until 17:30(2009→今≈2200发一次性回补, 完成后每日 no-op +# 全 done 秒退); rc=1|2 跳过 backfill 次日自愈。rc 聚合同 corpus: 真失败(1/2) +# 优先于让路(3)。 +# +# 错峰: 07:10 晚于 corpus 06:00——东财端点分域(datacenter-web vs +# search-api/np-listapi)+各自限速 1.5s/1.0s, 并发重叠安全依据同 corpus(spec §19.3)。 +# +# 部署位: NAS /volume1/stock/fundamentals/run_dsmk.sh(宿主文件,CI 不覆盖; +# 改本 repo 副本后 scp -O 手动同步, 同 corpus/5m 纪律): +# scp -O scripts/nas_sync/run_dmsk_standalone.sh \ +# sanguo-nas:/volume1/stock/fundamentals/run_dsmk.sh +# DSM 任务(sanguo-fundamentals-daily, 每日 07:10, 用户=admin)命令: +# bash /volume1/stock/fundamentals/run_dsmk.sh +# +# --no-healthcheck: 镜像 HEALTHCHECK=curl :8000/health 是 web 服务专用(同 corpus)。 +set -u +DOCKER=/var/packages/Docker/target/usr/bin/docker +ROOT=/volume1/stock/fundamentals +LOG=$ROOT/cron.log +APP=/volume1/homes/admin/.sanguo_projects/sanguo_vnpy_v2 +UID_ADMIN=1024 # id admin (2026-09-06 实测) +GID_ADMIN=100 # users +GID_ADMINS=101 # administrators(Synology ACL 授权组, corpus 金丝雀实测必带) + +mkdir -p "$ROOT" +{ + echo "=== $(date '+%F %T') standalone-container dmsk run start ===" + "$DOCKER" run --rm --name sanguo-fund-dmsk --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/fundamentals_dmsk_download.py --lane daily --until 08:00 + rc_daily=$? + echo "=== $(date '+%F %T') dmsk daily lane exit=$rc_daily ===" + rc_back=0 + if [ "$rc_daily" -eq 0 ] || [ "$rc_daily" -eq 3 ]; then + "$DOCKER" run --rm --name sanguo-fund-dmsk --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/fundamentals_dmsk_download.py --lane backfill --until 17:30 + rc_back=$? + echo "=== $(date '+%F %T') dmsk backfill lane exit=$rc_back ===" + else + echo "=== daily 非零($rc_daily), 跳过 backfill lane(次日自愈) ===" + fi +} >> "$LOG" 2>&1 +# rc 聚合: 真失败(1/2)优先于让路(3);全让路=3;全零=0 +if [ "${rc_back:-0}" -ne 0 ] && [ "${rc_back:-0}" -ne 3 ]; then + exit "$rc_back" +fi +if [ "${rc_daily:-1}" -ne 0 ] && [ "${rc_daily:-1}" -ne 3 ]; then + exit "$rc_daily" +fi +if [ "${rc_daily:-1}" -ne 0 ] || [ "${rc_back:-0}" -ne 0 ]; then + exit 3 +fi +exit 0 diff --git a/tests/data_platform/test_fundamentals_dmsk_download.py b/tests/data_platform/test_fundamentals_dmsk_download.py new file mode 100644 index 0000000..171c958 --- /dev/null +++ b/tests/data_platform/test_fundamentals_dmsk_download.py @@ -0,0 +1,326 @@ +# -*- coding: utf-8 -*- +"""TDD for fundamentals_dmsk_download.py — 财务三表 dmsk 补缺层(spec §19.4)。 + +全 mock 零网络;日期全部相对 today(禁写死);DMSK_ROOT 重定向 tmp_path。 +契约来源: spec §19.2 11 发实测(全市场单期≈5223行/57列/NOTICE_DATE/含退市)。 +""" +import datetime as dt +import json +import os +from unittest.mock import MagicMock + +import pandas as pd +import pytest + +from scripts.data_platform import fundamentals_dmsk_download as fdd + + +# ---------- Fixtures ---------- + +@pytest.fixture(autouse=True) +def _fast(monkeypatch): + """限速归零 + 抖动归零: 测试不 sleep。""" + monkeypatch.setattr(fdd, "RATE", {k: 0.0 for k in fdd.RATE}) + monkeypatch.setattr(fdd, "JITTER", (0.0, 0.0)) + + +@pytest.fixture +def root(tmp_path, monkeypatch): + monkeypatch.setattr(fdd, "DMSK_ROOT", tmp_path) + for sub in ("dmsk", "state", "logs"): + (tmp_path / sub).mkdir(parents=True, exist_ok=True) + return tmp_path + + +def _raw(code="600519", period="2024-12-31 00:00:00", notice="2025-04-03 00:00:00", + **over): + row = { + "SECURITY_CODE": code, "SECURITY_NAME_ABBR": "贵州茅台", + "REPORT_DATE": period, "NOTICE_DATE": notice, + "TOTAL_PARENT_EQUITY": 233105984399.47, + } + row.update(over) + return row + + +class FakeResp: + def __init__(self, payload): + self.status_code = 200 + self._payload = payload + + def json(self): + return self._payload + + +# ---------- 期枚举 ---------- + +def test_quarter_end_dates_desc_and_bounds(): + out = fdd.quarter_end_dates(2009, dt.date(2026, 9, 8)) + assert out[0] == dt.date(2026, 6, 30) # 未来期 09-30 不入列 + assert out[1] == dt.date(2026, 3, 31) + assert out[-1] == dt.date(2009, 3, 31) + assert out == sorted(out, reverse=True) + + +def test_quarter_end_dates_count(): + out = fdd.quarter_end_dates(2009, dt.date(2026, 9, 8), count=3) + assert out == [dt.date(2026, 6, 30), dt.date(2026, 3, 31), + dt.date(2025, 12, 31)] + + +# ---------- 归一化 ---------- + +def test_norm_row_truncates_and_zfills(): + r = fdd.norm_row(_raw(code="600519"), "2026-09-08") + assert r["REPORT_DATE"] == "2024-12-31" + assert r["SECURITY_CODE"] == "600519" + assert r["fetch_date"] == "2026-09-08" + assert r["TOTAL_PARENT_EQUITY"] == 233105984399.47 # 其余列原样透传 + + +def test_norm_row_zfills_short_numeric_code(): + assert fdd.norm_row(_raw(code="92001"), "d")["SECURITY_CODE"] == "092001" + + +def test_is_bj_code(): + assert fdd.is_bj_code("920001") + assert fdd.is_bj_code("833171") + assert fdd.is_bj_code("430047") + assert not fdd.is_bj_code("600519") + assert not fdd.is_bj_code("000001") + assert not fdd.is_bj_code("") + assert not fdd.is_bj_code(None) + + +# ---------- fetch 层 ---------- + +def test_fetch_period_rows_paginates_to_pages_end(root, monkeypatch): + calls = [] + + def fake_request(method, url, **kw): + calls.append(kw["params"]) + page = int(kw["params"]["pageNum"]) + return FakeResp({"result": {"pages": 3, "count": 6, "data": [ + _raw(code=f"60000{page}")] * 2}, "success": True}) + + cl = fdd.DomainClient("datacenter") + monkeypatch.setattr(cl, "request", fake_request) + rows = fdd.fetch_period_rows(cl, "balance", "2024-12-31") + assert len(rows) == 6 + assert len(calls) == 3 + assert calls[0]["reportName"] == "RPT_DMSK_FN_BALANCE" + assert calls[0]["filter"] == "(REPORT_DATE='2024-12-31')" + + +def test_fetch_period_rows_result_none_is_legal_empty(root): + cl = fdd.DomainClient("datacenter") + cl.session = MagicMock() + resp = FakeResp({"result": None, "success": False, "message": ""}) + cl.session.request.return_value = resp + cl._pace = lambda: None + assert fdd.fetch_period_rows(cl, "balance", "2009-03-31") == [] + + +# ---------- backfill lane ---------- + +def _patch_fetch(monkeypatch, responder): + calls = [] + + def fake(client, domain, report_date): + calls.append((domain, report_date)) + return responder(domain, report_date) + + monkeypatch.setattr(fdd, "fetch_period_rows", fake) + return calls + + +def test_backfill_writes_parts_marks_and_resumes(root, monkeypatch): + _patch_fetch(monkeypatch, lambda d, p: [_raw(period=p + " 00:00:00")]) + rc1 = fdd.run_lane("backfill", start_year=2024) + assert rc1 == 0 + part = root / "dmsk" / "balance" / "dt=2024-12-31" / "part-0.parquet" + assert part.exists() + df = pd.read_parquet(part) + assert df["SECURITY_CODE"].iloc[0] == "600519" + # 全期 marker 落平铺(一次永逸, 键=域+期) + m = (root / "state" / "markers" / "backfill" / "period" + / "balance_2024-12-31.done") + assert m.exists() + # 重跑: 全 done → 零 fetch + calls2 = _patch_fetch(monkeypatch, lambda d, p: [_raw()]) + rc2 = fdd.run_lane("backfill", start_year=2024) + assert rc2 == 0 + assert calls2 == [] + + +def test_backfill_resumes_from_oldest_gap(root, monkeypatch): + _patch_fetch(monkeypatch, lambda d, p: [_raw(period=p + " 00:00:00")]) + fdd.run_lane("backfill", start_year=2024) + (root / "state" / "markers" / "backfill" / "period" + / "balance_2025-12-31.done").unlink() # 挖一个洞 + calls = _patch_fetch(monkeypatch, lambda d, p: [_raw(period=p + " 00:00:00")]) + fdd.run_lane("backfill", start_year=2024) + assert calls == [("balance", "2025-12-31")] + + +def test_backfill_ledger_prevents_rekey_across_runs(root, monkeypatch): + _patch_fetch(monkeypatch, lambda d, p: [_raw(period=p + " 00:00:00")]) + fdd.run_lane("backfill", start_year=2024) + # 清光 marker 重跑: 账本挡住 → 0 新行(键唯一硬门) + for p in (root / "state" / "markers" / "backfill").rglob("*.done"): + p.unlink() + fdd.run_lane("backfill", start_year=2024) + for part in (root / "dmsk").rglob("part-*.parquet"): + df = pd.read_parquet(part) + assert df.duplicated(["SECURITY_CODE", "REPORT_DATE"]).sum() == 0 + assert len(df) == 1 # 只有首跑那一行 + + +def test_in_batch_dedup_same_key_twice_in_response(root, monkeypatch): + # 同 (code,period) 两条同响应 → 账本外批内 seen 截留(corpus 六修正之三) + _patch_fetch(monkeypatch, lambda d, p: [ + _raw(period=p + " 00:00:00"), _raw(period=p + " 00:00:00")]) + fdd.run_lane("backfill", start_year=2025) + df = pd.read_parquet(root / "dmsk" / "balance" + / "dt=2025-12-31" / "part-0.parquet") + assert len(df) == 1 + + +def test_backfill_canary_empty_fails_and_never_marks(root, monkeypatch): + # 次新期空=契约漂移疑云: failed+1, 不标 done 不落空库 + def responder(d, p): + return [] if p == "2026-03-31" else [_raw(period=p + " 00:00:00")] + _patch_fetch(monkeypatch, responder) + rc = fdd.run_lane("backfill", start_year=2026) + assert rc == 1 + assert not (root / "state" / "markers" / "backfill" / "period" + / "balance_2026-03-31.done").exists() + + +def test_backfill_empty_early_period_is_legal_done(root, monkeypatch): + # 2009 年内非年报期早于首份数据: 合法空期照样标 done(金丝雀期有数据) + def responder(d, p): + if p.startswith("2009") and not p.endswith("12-31"): + return [] + return [_raw(period=p + " 00:00:00")] + _patch_fetch(monkeypatch, responder) + rc = fdd.run_lane("backfill", start_year=2009) + assert rc == 0 + assert (root / "state" / "markers" / "backfill" / "period" + / "balance_2009-03-31.done").exists() + assert not (root / "dmsk" / "balance" / "dt=2009-03-31").exists() + + +def test_backfill_bj_inclusion_artifact(root, monkeypatch): + _patch_fetch(monkeypatch, lambda d, p: [ + _raw(code="920001", period=p + " 00:00:00")]) + fdd.run_lane("backfill", start_year=2025) + j = json.loads((root / "state" / "bj_inclusion.json").read_text("utf-8")) + assert j["balance"]["2025-12-31"] >= 1 + + +def test_limit_never_marks(root, monkeypatch): + _patch_fetch(monkeypatch, lambda d, p: [_raw(period=p + " 00:00:00")]) + fdd.run_lane("backfill", limit=1, start_year=2024) + assert list((root / "state" / "markers" / "backfill").rglob("*.done")) == [] + + +# ---------- daily lane ---------- + +def test_daily_notice_filter_and_period_dir(root, monkeypatch): + today = dt.date.today() + fresh = (today - dt.timedelta(days=1)).isoformat() + stale = (today - dt.timedelta(days=30)).isoformat() + + def responder(d, p): + return [_raw(code="600519", period=p + " 00:00:00", notice=fresh), + _raw(code="000001", period=p + " 00:00:00", notice=stale)] + _patch_fetch(monkeypatch, responder) + rc = fdd.run_lane("daily") + assert rc == 0 + # 落点=dt=/ 而非 dt=today; 近3期各得 1 行新披露 + pdirs = sorted(p.name for p in (root / "dmsk" / "balance").glob("dt=*")) + assert len(pdirs) == 3 + assert f"dt={today.isoformat()}" not in pdirs + df = pd.read_parquet(root / "dmsk" / "balance" / pdirs[0] / "part-0.parquet") + assert list(df["SECURITY_CODE"]) == ["600519"] # NOTICE_DATE