# -*- 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), "pageNumber": str(page), # 东财 datacenter 用 pageNumber(cninfo # 才是 pageNum; 写错被静默忽略=翻页 # 失效, 每页都回第一页, 冒烟 12 页同 500 行实锤) "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()