271 lines
10 KiB
Python
271 lines
10 KiB
Python
"""监控检查纯函数核(spec §20.15.2 四层证据链的结果层+过程层解析件).
|
|
|
|
约定: 交易日历=list['YYYY-MM-DD'] 升序; 期望值='YYYYMMDD'(与存盘 stem 对齐);
|
|
now 的 ready 语义=「数据当日就绪时刻」——now 早于 ready 则该日尚不该有数据。
|
|
"""
|
|
import datetime as _dt
|
|
import os
|
|
import re
|
|
import shutil
|
|
import sqlite3
|
|
|
|
|
|
# ---------- 交易日历(dbbardata 沪深300 日线当锚; 仓库无独立日历表) ----------
|
|
|
|
def trading_days(db_path, today, lookback=50):
|
|
"""近 lookback 日历日的交易日序列; db 不可达返 [](调用方降级周一~五启发式)."""
|
|
cutoff = (today - _dt.timedelta(days=lookback)).isoformat()
|
|
try:
|
|
conn = sqlite3.connect(f"file:{db_path}?mode=ro", uri=True, timeout=10)
|
|
try:
|
|
rows = conn.execute(
|
|
"SELECT DISTINCT datetime FROM dbbardata "
|
|
"WHERE symbol='000300' AND exchange='SSE' AND interval='d' "
|
|
"AND datetime>=? ORDER BY datetime", (cutoff,)).fetchall()
|
|
finally:
|
|
conn.close()
|
|
# D-3(10-05): 读侧 [:10] 归一+None 过滤, 对齐 strategy_checks 同名
|
|
# 版(双实现已分叉; dbbardata datetime 历史双格式, 写侧归一不能
|
|
# 担保读侧形状——脏值曾致 expected_daily 全 per_date 判定崩坏)
|
|
return sorted({str(r[0])[:10] for r in rows if r[0]})
|
|
except sqlite3.Error:
|
|
return []
|
|
|
|
|
|
def _hhmm(ready):
|
|
h, m = (int(x) for x in ready.split(":"))
|
|
return h * 60 + m
|
|
|
|
|
|
def expected_daily(days, now, ready="20:30", lag=0):
|
|
if not days:
|
|
return None
|
|
now_min = now.hour * 60 + now.minute
|
|
today = now.date().isoformat()
|
|
ok = [d for d in days if d < today
|
|
or (d == today and now_min >= _hhmm(ready))]
|
|
if len(ok) <= lag:
|
|
return None
|
|
return ok[-1 - lag].replace("-", "")
|
|
|
|
|
|
def expected_weekly(days, now, ready="06:00", confirmed=None):
|
|
"""最后周五锚 F: 当周六 ready 后 F 可验(ak-weekly 周六 03:00 产出).
|
|
|
|
confirmed(YYYYMMDD 集)非空时周五锚须为交易日——假定集假周五(2026-09-25
|
|
中秋/10-02 国庆)滤掉锚自动回上一真周五; confirmed 空集=dbbardata 不可达
|
|
→ 降级不过滤(旧行为), 防 exp=None 静默假绿(factor 交界提醒 0927)."""
|
|
if not days:
|
|
return None
|
|
deadline = now - _dt.timedelta(days=1)
|
|
floor = deadline.date() if now.hour * 60 + now.minute >= _hhmm(ready) \
|
|
else deadline.date() - _dt.timedelta(days=1)
|
|
fridays = [d for d in days
|
|
if _dt.date.fromisoformat(d).weekday() == 4
|
|
and d <= floor.isoformat()]
|
|
if confirmed:
|
|
fridays = [d for d in fridays if d.replace("-", "") in confirmed]
|
|
return fridays[-1].replace("-", "") if fridays else None
|
|
|
|
|
|
# ---------- 目录/文件 ----------
|
|
|
|
def latest_stem(directory, pattern):
|
|
if not os.path.isdir(directory):
|
|
return None
|
|
rx = re.compile(pattern)
|
|
best = None
|
|
for fn in os.listdir(directory):
|
|
stem = fn.rsplit(".", 1)[0] if "." in fn else fn
|
|
m = rx.match(stem)
|
|
if m and (best is None or stem > best[0]):
|
|
best = (stem, os.path.join(directory, fn))
|
|
return best
|
|
|
|
|
|
def parquet_rows(path):
|
|
try:
|
|
import pyarrow.parquet as pq
|
|
return pq.ParquetFile(path).metadata.num_rows
|
|
except Exception:
|
|
return None
|
|
|
|
|
|
# ---------- 滚动分位带 (P2-6, 2026-10-01 用户拍板「监控问题统一处理」) ----------
|
|
# 注册表 row_band 语义升格为 (floor, fallback_hi): floor=绝对下限永不下调
|
|
# (采集缓降不能靠带内绿掉); fallback_hi=冷启动/观测不足时的整带回退值, 不再
|
|
# 作常开上限——静态上限三次被真实增长涨破的手动校准史(0be832d/3f675ed/4755e8d)
|
|
# 即滚动化论证。带内判定=P1~P99(窗口=最近 ROLLING_WINDOW 个 prior 文件, 剔除
|
|
# 被测观测与 0 行真空文件); 上界失控由 trailing-P99 机理兜住(双写翻倍 Day1 即
|
|
# 黄; 公告驱动缓涨在窗口吸收期连黄数日~数周=非静默信号, 值班预期见 runbook)。
|
|
ROLLING_WINDOW = 30
|
|
ROLLING_MIN_OBS = 10
|
|
ROLLING_Q = (0.01, 0.99)
|
|
|
|
|
|
def rolling_rows_band(directory, pattern, exclude_stem, floor, fallback_hi,
|
|
window=None, min_obs=None, q=None):
|
|
"""滚动分位带计算: 返回 (lo_eff, hi_eff, n_obs, source).
|
|
|
|
source="rolling"=带内判据用 P1~P99(下界再被 floor 托底); "fallback"=观测
|
|
不足 min_obs 回退静态整带 (floor, fallback_hi)。exclude_stem=当前被检文件
|
|
(窗口必须是历史, 不含被测观测)。
|
|
"""
|
|
window = ROLLING_WINDOW if window is None else window
|
|
min_obs = ROLLING_MIN_OBS if min_obs is None else min_obs
|
|
q = ROLLING_Q if q is None else q
|
|
if not os.path.isdir(directory):
|
|
return float(floor), float(fallback_hi), 0, "fallback"
|
|
rx = re.compile(pattern)
|
|
stems = []
|
|
for fn in os.listdir(directory):
|
|
stem = fn.rsplit(".", 1)[0] if "." in fn else fn
|
|
if rx.match(stem) and stem != exclude_stem:
|
|
stems.append(stem)
|
|
stems.sort(reverse=True)
|
|
sample = []
|
|
for stem in stems[:window]:
|
|
rows = parquet_rows(os.path.join(directory, stem + ".parquet"))
|
|
if rows: # 0 行真空文件=「无数据」非「低数据」, 不进样本
|
|
sample.append(rows)
|
|
if len(sample) < min_obs:
|
|
return float(floor), float(fallback_hi), len(sample), "fallback"
|
|
import numpy as np
|
|
lo_q = float(np.quantile(sample, q[0]))
|
|
hi_q = float(np.quantile(sample, q[1]))
|
|
return max(lo_q, float(floor)), hi_q, len(sample), "rolling"
|
|
|
|
|
|
def dir_latest_mtime(directory):
|
|
latest = 0.0
|
|
for root, _dirs, files in os.walk(directory):
|
|
for fn in files:
|
|
try:
|
|
latest = max(latest, os.path.getmtime(os.path.join(root, fn)))
|
|
except OSError:
|
|
pass
|
|
return latest
|
|
|
|
|
|
def check_disk(path, warn_gb=20, crit_gb=8):
|
|
"""磁盘剩余水位; 阈值参数化(D-8 10-05: 注册表 warn_gb/crit_gb 消费——旧
|
|
硬编码 8/20 与注册表脱钩, crit_gb 曾全仓无消费=改表不改行为的校准陷阱)."""
|
|
target = path if os.path.exists(path) else os.path.dirname(os.path.abspath(path))
|
|
usage = shutil.disk_usage(target)
|
|
free_gb = usage.free / 1024 ** 3
|
|
if free_gb < crit_gb:
|
|
status = "red"
|
|
elif free_gb < warn_gb:
|
|
status = "yellow"
|
|
else:
|
|
status = "green"
|
|
return {"free_gb": round(free_gb, 1), "status": status}
|
|
|
|
|
|
# ---------- 过程层: 统计行+失败签名分诊(spec §20.15.2 ②) ----------
|
|
|
|
_STAT_RE = re.compile(r"\[(\w+)\] 完成.*?统计: (\{[^\n]*\})")
|
|
|
|
SIGNATURES = [
|
|
("断路器触发", "breaker", "red"),
|
|
("RemoteDisconnected", "disconnect", "red"),
|
|
("Expecting value", "source_changed", "red"),
|
|
("418 Client Error", "anti_crawl", "red"),
|
|
("'NoneType' object is not subscriptable", "source_empty", "yellow"),
|
|
("KeyError", "field_changed", "red"),
|
|
]
|
|
|
|
|
|
def parse_stat_lines(text):
|
|
import json
|
|
out = {}
|
|
for m in _STAT_RE.finditer(text):
|
|
try:
|
|
out[m.group(1)] = json.loads(m.group(2))
|
|
except ValueError:
|
|
continue
|
|
return out
|
|
|
|
|
|
def _error_lines(text):
|
|
"""只留 ERROR 行(重试 WARNING 良性) + [FATAL] 断路器行。"""
|
|
return [ln for ln in text.splitlines()
|
|
if " ERROR " in ln or "[FATAL]" in ln]
|
|
|
|
|
|
def scan_signatures(text):
|
|
"""10-05 三修: ①只扫 ERROR/[FATAL] 行——'Expecting value' 出现在重试
|
|
WARNING 行(financial_abstract/000963 实证), 良性重试不再误中
|
|
source_changed; ②'418' 针收紧为 '418 Client Error'——裸 418 子串误中
|
|
rows=41823/ok=19418 进度行(anti_crawl 假阳性实证)。
|
|
**注(10-05 补充审计)**: 生产路径已无调用方(签名分诊走 signatures_by_type,
|
|
data_monitor.log_layer_results/log_only_results)——本函数仅测试消费,
|
|
勿再接线; 与 signatures_by_type 的行为差异=不按 [type] 归因。"""
|
|
lines = _error_lines(text)
|
|
seen, out = set(), []
|
|
for needle, sig, severity in SIGNATURES:
|
|
if any(needle in ln for ln in lines) and sig not in seen:
|
|
seen.add(sig)
|
|
out.append({"sig": sig, "severity": severity})
|
|
return out
|
|
|
|
|
|
_TYPE_TAG = re.compile(r"\[(\w+)\]")
|
|
|
|
|
|
def signatures_by_type(text):
|
|
"""签名按日志行 [type] 归因(10-05): 'ERROR [share_capital] ...' 行的
|
|
签名只归 share_capital——防跨类型连坐(xueqiu 的 418 把 share_capital
|
|
的 failed 抬红实证)。断路器行归触发类型, data_monitor 侧全局取。
|
|
[FATAL] 是级别标记不是类型(P3 补充审计): 裸 FATAL 行
|
|
(akshare_static_download.py:2721 实格式)归 _global, [type] [FATAL]
|
|
双标记行(:2160 实格式)归 [type]。"""
|
|
out = {}
|
|
for ln in _error_lines(text):
|
|
m = _TYPE_TAG.search(ln)
|
|
while m and m.group(1) == "FATAL":
|
|
m = _TYPE_TAG.search(ln, m.end())
|
|
t = m.group(1) if m else "_global"
|
|
sigs = out.setdefault(t, [])
|
|
for needle, sig, severity in SIGNATURES:
|
|
if needle in ln and sig not in {s["sig"] for s in sigs}:
|
|
sigs.append({"sig": sig, "severity": severity})
|
|
return out
|
|
|
|
|
|
# ---------- 调度层: schtasks /query /fo LIST /v 解析(双语字段名) ----------
|
|
|
|
_FIELD_PATTERNS = {
|
|
"task_name": r"^(?:TaskName|任务名)\s*:\s*(.*)$",
|
|
"next_run": r"^(?:Next Run Time|下次运行时间)\s*:\s*(.*)$",
|
|
"status": r"^(?:Status|状态)\s*:\s*(.*)$",
|
|
"last_run": r"^(?:Last Run Time|上次运行时间)\s*:\s*(.*)$",
|
|
"last_result": r"^(?:Last Result|上次结果)\s*:\s*(.*)$",
|
|
}
|
|
_COMPILED = {k: re.compile(p) for k, p in _FIELD_PATTERNS.items()}
|
|
|
|
|
|
def parse_schtasks_list(text):
|
|
out = {}
|
|
cur = None
|
|
for line in text.splitlines():
|
|
stripped = line.strip()
|
|
if not stripped:
|
|
continue
|
|
matched = False
|
|
for key, rx in _COMPILED.items():
|
|
m = rx.match(stripped)
|
|
if not m:
|
|
continue
|
|
val = m.group(1).strip()
|
|
if key == "task_name":
|
|
cur = val.lstrip("\\")
|
|
out.setdefault(cur, {})
|
|
elif cur is not None:
|
|
out[cur][key] = val
|
|
matched = True
|
|
break
|
|
if not matched and cur is None:
|
|
continue # HostName 等无关行
|
|
return out
|