236bd30eb8
- 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缺口入档
698 lines
25 KiB
Python
698 lines
25 KiB
Python
# -*- 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()
|