From d175bb65d4b771cea116a8bbfb0bd646744cba8c Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E6=80=BB=E5=B7=A5?= Date: Mon, 7 Sep 2026 21:19:35 +0800 Subject: [PATCH] =?UTF-8?q?feat(HEL-478):=20=E4=BC=B0=E5=80=BC=E5=8F=91?= =?UTF-8?q?=E5=B8=83=E5=90=8E=E6=99=9A=E9=97=B4=E5=A4=8D=E6=A0=B8=E5=B9=B6?= =?UTF-8?q?=E5=8E=9F=E5=AD=90=E8=BF=BD=E8=A1=A5=E4=B8=8A=E6=B8=B8=E4=BF=AE?= =?UTF-8?q?=E8=AE=A2?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 盘后成功发布后继续轻量比对 daily_basic 网站字段,发现修订才走质量门与整组原子切换,避免 17:10 快照落后于晚间上游改写。 Co-authored-by: Cursor Co-authored-by: multica-agent --- xiaobai-datahub/README.md | 12 + xiaobai-datahub/admin/app.js | 13 + .../config/hub-quality.config.json | 4 + xiaobai-datahub/datahub/admin_api.py | 2 + xiaobai-datahub/datahub/db.py | 13 + xiaobai-datahub/datahub/pipeline.py | 240 +++++++++++++ xiaobai-datahub/datahub/revision.py | 111 ++++++ xiaobai-datahub/datahub/scheduler.py | 219 +++++++++++- xiaobai-datahub/datahub/settings.py | 14 + xiaobai-datahub/tests/test_eod_retry.py | 25 +- xiaobai-datahub/tests/test_revision_review.py | 333 ++++++++++++++++++ 11 files changed, 982 insertions(+), 4 deletions(-) create mode 100644 xiaobai-datahub/datahub/revision.py create mode 100644 xiaobai-datahub/tests/test_revision_review.py diff --git a/xiaobai-datahub/README.md b/xiaobai-datahub/README.md index 607e299..d1f55de 100644 --- a/xiaobai-datahub/README.md +++ b/xiaobai-datahub/README.md @@ -123,6 +123,18 @@ python -m datahub eod-refresh --trade-date 20260904 --force --dataset valuation 管理后台「补数」对盘后正式数据集同样走 `force_republish_boundary`,不会绕过 A/B 整批边界。 +## 估值发布后复核与自动追补 + +Tushare `daily_basic` 会在盘后继续改当日字段。HEL-423 在 2026-09-07 观察到:中枢 17:10 发布 `003021.SZ turnover_rate=1.3565`,21:05 上游/旧链路已是 `1.3572`;其余 7 类观察对象当日一致。日 K、资金流、竞价、指数没有同类晚间修订证据,股票主档已有 20:00/23:10 刷新,因此默认只复核估值,不盲目全量重拉。 + +窗口(可配):交易日 **20:00–23:20**,每 30 分钟一次轻量比对(对齐网站 21:00 / 23:30 观察)。只拉取 `daily_basic`,按网站真实请求字段精确比较,无误差豁免。 + +- 无变化:不产生新批次,状态「已追平」。 +- 发现修订:重新走字段质量门、覆盖检查和 A 组整批原子发布;读者全程只能看到上一完整版本或新完整版本。 +- 上游空 / 接口失败 / 不完整 / 质量门拒绝:保留上一完整版本,状态「复核失败」。 +- 23:20 截止后停止当晚复核;下一自然日盘前对上一交易日再做一次安全追赶。 +- 与 `eod_a` / `eod_retry` 共用互斥锁;容器重启会在窗口内立即补一次。 + ## 备份 每日 00:40 任务把 `datahub.db` 备份到 `data/backups/`(保留 14 份)。也可手动: diff --git a/xiaobai-datahub/admin/app.js b/xiaobai-datahub/admin/app.js index 5e19f6e..afaf599 100644 --- a/xiaobai-datahub/admin/app.js +++ b/xiaobai-datahub/admin/app.js @@ -105,6 +105,7 @@ async function render() { const data = await api("/admin/api/overview"); $("phase").textContent = data.session_phase; const eod = data.eod_status || {}; + const rev = data.revision_status || {}; const eodLabels = { pending_first_attempt: "等待首次尝试", waiting_upstream: "等待上游", @@ -112,6 +113,14 @@ async function render() { cutoff_failed: "已截止失败", closed_day: "休市", }; + const revLabels = { + waiting_review: "等待复核", + review_failed: "复核失败", + aligned: "已追平", + cutoff: "已截止", + pending_publish: "待发布", + closed_day: "休市", + }; const eodExtra = []; if (eod.state === "waiting_upstream") { eodExtra.push(`已试 ${eod.attempts} 次`); @@ -121,12 +130,16 @@ async function render() { if (eod.state === "cutoff_failed" && eod.missing_datasets) { eodExtra.push(`缺 ${esc(eod.missing_datasets.join(","))}`); } + const revExtra = []; + if (rev.detail) revExtra.push(esc(String(rev.detail))); + if (rev.window) revExtra.push(esc(String(rev.window))); page.innerHTML = `
交易日
${esc(data.trade_date)}
阶段
${esc(data.session_phase)}
今日发布
${data.publications.length}
盘后补跑
${esc(eodLabels[eod.state] || eod.state || "-")}
${eodExtra.join(" · ")}
+
估值复核
${esc(revLabels[rev.state] || rev.state || "-")}
${revExtra.join(" · ")}
异常批次
${data.anomalies.length}

最近调用

diff --git a/xiaobai-datahub/config/hub-quality.config.json b/xiaobai-datahub/config/hub-quality.config.json index 93c4514..578dea8 100644 --- a/xiaobai-datahub/config/hub-quality.config.json +++ b/xiaobai-datahub/config/hub-quality.config.json @@ -17,6 +17,10 @@ "eod_retry_start": "15:15", "eod_retry_interval_minutes": 30, "eod_retry_cutoff": "23:30", + "revision_review_datasets": ["valuation"], + "revision_review_start": "20:00", + "revision_review_interval_minutes": 30, + "revision_review_cutoff": "23:20", "moneyflow_history_trading_days": 60, "stocks_refresh_times": [ "20:00", diff --git a/xiaobai-datahub/datahub/admin_api.py b/xiaobai-datahub/datahub/admin_api.py index 41f162b..f6c86a7 100644 --- a/xiaobai-datahub/datahub/admin_api.py +++ b/xiaobai-datahub/datahub/admin_api.py @@ -39,6 +39,7 @@ class AdminAPI: "session_phase": session_phase(now_shanghai(), is_open), "is_open_day": is_open, "eod_status": self.scheduler.eod_status(today), + "revision_status": self.scheduler.revision_status(today), "publications": pubs, "anomalies": failed, "recent_calls": _public_calls(calls), @@ -91,6 +92,7 @@ class AdminAPI: {"id": "eod_a", "at": "15:05", "title": "盘后批 A daily/valuation/moneyflow/auction"}, {"id": "eod_b", "at": "15:10", "title": "盘后批 B index_daily"}, {"id": "eod_retry", "at": "15:15-23:30", "title": "盘后未出数自动重试(每 30 分钟,成功即停)"}, + {"id": "eod_revise", "at": "20:00-23:20", "title": "估值发布后复核(轻量比对,有修订才整组原子追补)"}, {"id": "stocks_refresh", "at": stocks_times, "title": "股票主档刷新与正式发布(新上市/更名,无变化跳过)"}, {"id": "history_backfill", "at": "manual", "title": "回补历史日历与指数日 K"}, {"id": "cleanup", "at": "00:30", "title": "清理 staging / 日志"}, diff --git a/xiaobai-datahub/datahub/db.py b/xiaobai-datahub/datahub/db.py index bcea373..4139a35 100644 --- a/xiaobai-datahub/datahub/db.py +++ b/xiaobai-datahub/datahub/db.py @@ -238,6 +238,19 @@ CREATE TABLE IF NOT EXISTS eod_progress ( updated_at TEXT NOT NULL ); +CREATE TABLE IF NOT EXISTS revision_progress ( + trade_date TEXT PRIMARY KEY, + state TEXT NOT NULL, + attempts INTEGER NOT NULL DEFAULT 0, + last_attempt_at TEXT, + next_retry_at TEXT, + finished_at TEXT, + catchup_done INTEGER NOT NULL DEFAULT 0, + last_diff TEXT, + detail TEXT, + updated_at TEXT NOT NULL +); + CREATE TABLE IF NOT EXISTS audit_log ( id INTEGER PRIMARY KEY AUTOINCREMENT, actor TEXT NOT NULL, diff --git a/xiaobai-datahub/datahub/pipeline.py b/xiaobai-datahub/datahub/pipeline.py index c6c5d8c..0a499de 100644 --- a/xiaobai-datahub/datahub/pipeline.py +++ b/xiaobai-datahub/datahub/pipeline.py @@ -15,6 +15,12 @@ from datahub.governance.ratelimit import TokenBucket from datahub.governance.retry import RetryError, retry_call from datahub.logutil import get_logger from datahub.normalize import finite_number, normalize_daily +from datahub.revision import ( + compare_fields, + diff_published_vs_upstream, + official_table, + revision_datasets, +) from datahub.settings import Settings from datahub.timeutil import add_days, isoformat, now_shanghai, yyyymmdd @@ -687,6 +693,240 @@ class Pipeline: return self.run_eod_batch_b(trade_date, force=True) raise ValueError(f"dataset is not part of an EOD release boundary: {dataset}") + def published_official_rows(self, dataset: str, trade_date: str) -> list[dict[str, Any]]: + day = yyyymmdd(trade_date) + batch_id = self.active_batch(dataset, day) + if not batch_id: + return [] + fields = compare_fields(dataset) + table = official_table(dataset) + if fields: + columns = ",".join(fields) + return self.db.fetchall( + f"SELECT {columns} FROM {table} WHERE batch_id = ?", + (batch_id,), + ) + return self.db.fetchall(f"SELECT * FROM {table} WHERE batch_id = ?", (batch_id,)) + + def compare_revision(self, dataset: str, trade_date: str) -> dict[str, Any]: + """Light fetch of one revision-risk dataset vs the published official rows.""" + day = yyyymmdd(trade_date) + published = self.published_official_rows(dataset, day) + if not published: + return { + "dataset": dataset, + "trade_date": day, + "changed": False, + "state": "skipped", + "reason": "not_published", + } + try: + upstream = self._fetch_dataset(dataset, day) + except Exception as exc: + return { + "dataset": dataset, + "trade_date": day, + "changed": False, + "state": "failed", + "reason": "upstream_error", + "error": str(exc), + } + if not upstream: + return { + "dataset": dataset, + "trade_date": day, + "changed": False, + "state": "failed", + "reason": "upstream_empty", + "error": "revision review upstream empty", + "published_rows": len(published), + "upstream_rows": 0, + } + listed = self.db.fetchone( + "SELECT COUNT(*) AS n FROM stock_master WHERE list_status = 'L'", + ) + listed_n = int((listed or {}).get("n") or 0) + floor = float(self.settings.quality.get("daily_row_ratio") or 0.98) + if listed_n and len(upstream) / listed_n < floor: + return { + "dataset": dataset, + "trade_date": day, + "changed": False, + "state": "failed", + "reason": "incomplete", + "error": ( + f"revision review incomplete: upstream {len(upstream)} " + f"/ listed {listed_n} < {floor}" + ), + "published_rows": len(published), + "upstream_rows": len(upstream), + } + if len(upstream) < len(published) * floor: + return { + "dataset": dataset, + "trade_date": day, + "changed": False, + "state": "failed", + "reason": "incomplete", + "error": ( + f"revision review incomplete: upstream {len(upstream)} " + f"< published {len(published)} * {floor}" + ), + "published_rows": len(published), + "upstream_rows": len(upstream), + } + compared = diff_published_vs_upstream(dataset, published, upstream) + compared["trade_date"] = day + compared["state"] = "changed" if compared["changed"] else "unchanged" + compared["reason"] = "revised" if compared["changed"] else "unchanged" + return compared + + def review_published_revisions(self, trade_date: str) -> dict[str, Any]: + """Evening/morning catch-up: compare website fields, republish only on change. + + Unchanged → no new batch. Changed → full A/B boundary quality gate + + atomic switch (HEL-459/460/461). Empty/failed/incomplete upstream keeps + the previous complete official version. + """ + day = yyyymmdd(trade_date) + results: dict[str, Any] = {} + for dataset in revision_datasets(self.settings.quality): + compared = self.compare_revision(dataset, day) + if compared.get("state") == "skipped": + results[dataset] = compared + continue + if compared.get("state") == "failed": + results[dataset] = compared + LOGGER.warning( + "revision review kept previous official version", + extra={ + "hub": { + "dataset": dataset, + "trade_date": day, + "reason": compared.get("reason"), + "event": "revision_review_failed", + } + }, + ) + self.audit( + "pipeline", "revision-review", f"{dataset}:{day}", + json.dumps( + { + "state": "failed", + "reason": compared.get("reason"), + "error": compared.get("error"), + }, + ensure_ascii=False, + ), + ) + continue + if not compared.get("changed"): + results[dataset] = { + "dataset": dataset, + "trade_date": day, + "state": "aligned", + "reason": "unchanged", + "published_rows": compared.get("published_rows"), + "upstream_rows": compared.get("upstream_rows"), + } + self.audit( + "pipeline", "revision-review", f"{dataset}:{day}", + json.dumps({"state": "aligned", "reason": "unchanged"}, ensure_ascii=False), + ) + continue + LOGGER.info( + "revision review detected upstream rewrite, republishing boundary", + extra={ + "hub": { + "dataset": dataset, + "trade_date": day, + "diffs": compared.get("diffs"), + "event": "revision_review_changed", + } + }, + ) + try: + published = self.force_republish_boundary(dataset, day) + except Exception as exc: + results[dataset] = { + "dataset": dataset, + "trade_date": day, + "state": "failed", + "reason": "republish_error", + "error": str(exc), + "diffs": compared.get("diffs"), + } + LOGGER.warning( + "revision republish failed, previous official version keeps serving", + extra={ + "hub": { + "dataset": dataset, + "trade_date": day, + "reason": str(exc), + "event": "revision_review_failed", + } + }, + ) + self.audit( + "pipeline", "revision-review", f"{dataset}:{day}", + json.dumps( + {"state": "failed", "reason": "republish_error", "error": str(exc)}, + ensure_ascii=False, + ), + ) + continue + failures = self.eod_failures(published) + if failures: + results.update(published) + results[dataset] = { + **(published.get(dataset) or {}), + "dataset": dataset, + "trade_date": day, + "state": "failed", + "reason": "quality_gate", + "error": "; ".join(failures), + "diffs": compared.get("diffs"), + } + self.audit( + "pipeline", "revision-review", f"{dataset}:{day}", + json.dumps( + { + "state": "failed", + "reason": "quality_gate", + "error": "; ".join(failures), + "diffs": compared.get("diffs"), + }, + ensure_ascii=False, + ), + ) + continue + results.update(published) + results["review"] = { + "dataset": dataset, + "trade_date": day, + "state": "aligned", + "reason": "revised", + "diffs": compared.get("diffs"), + "missing_codes": compared.get("missing_codes"), + "extra_codes": compared.get("extra_codes"), + } + self.audit( + "pipeline", "revision-review", f"{dataset}:{day}", + json.dumps( + { + "state": "aligned", + "reason": "revised", + "diffs": compared.get("diffs"), + "switched": sorted( + name for name, item in published.items() + if isinstance(item, dict) and item.get("state") == "published" + ), + }, + ensure_ascii=False, + ), + ) + return results + def run_release_group( self, datasets: tuple[str, ...], diff --git a/xiaobai-datahub/datahub/revision.py b/xiaobai-datahub/datahub/revision.py new file mode 100644 index 0000000..7c3c304 --- /dev/null +++ b/xiaobai-datahub/datahub/revision.py @@ -0,0 +1,111 @@ +"""Post-publish revision review for datasets whose upstream may rewrite T-day fields. + +HEL-423 field evidence, not a whitelist of tolerated diffs: + +- 2026-09-07 valuation/daily_basic: hub published 003021.SZ turnover_rate=1.3565 + at 17:10; website legacy and a direct Tushare read at 21:05 both showed 1.3572. + The other seven observed objects (daily, moneyflow, auction, stocks, status, + index_daily, calendar) matched. Hub had already stopped the day after the + first successful publish, so the revision never self-healed. +- 2026-09-02: same dataset, opposite direction (hub already held the later + value). Confirms daily_basic is rewritten after the first complete dump. + +Daily bars, moneyflow, auction and index_daily have no same-evening field +revision evidence. Stocks already refreshes at 20:00/23:10. Review therefore +fetches only configured revision-risk datasets (default: valuation) and +compares the website-requested field set. No numeric tolerance. +""" + +from __future__ import annotations + +from typing import Any + +from datahub.db import DATASET_TABLES +from datahub.normalize import VALUATION_FIELDS +from datahub.numbers import finite_number, round4 + +# Datasets with proven same-evening upstream rewrites. Config may replace this +# list; it must not silently expand to a full EOD re-pull. +DEFAULT_REVISION_DATASETS = ("valuation",) + +# Website daily_basic request (HEL-423): ts_code/trade_date plus the eight +# value fields used by the old link and field_gates. +WEBSITE_COMPARE_FIELDS: dict[str, tuple[str, ...]] = { + "valuation": VALUATION_FIELDS, +} + +REVISION_STATES = ("waiting_review", "review_failed", "aligned", "cutoff") + + +def revision_datasets(quality: dict[str, Any] | None) -> tuple[str, ...]: + raw = (quality or {}).get("revision_review_datasets") + if isinstance(raw, (list, tuple)) and raw: + names = tuple(str(item) for item in raw if str(item)) + if names: + return names + return DEFAULT_REVISION_DATASETS + + +def compare_fields(dataset: str) -> tuple[str, ...]: + fields = WEBSITE_COMPARE_FIELDS.get(dataset) + if fields: + return fields + gate = {} + return tuple(str(item) for item in (gate.get("fields") or []) if str(item)) + + +def _norm_value(field: str, value: Any) -> Any: + if field in {"ts_code", "trade_date"}: + return str(value or "") + number = round4(finite_number(value)) + return number + + +def row_signature(row: dict[str, Any], fields: tuple[str, ...]) -> tuple[Any, ...]: + return tuple(_norm_value(field, row.get(field)) for field in fields) + + +def diff_published_vs_upstream( + dataset: str, + published: list[dict[str, Any]], + upstream: list[dict[str, Any]], + *, + max_diffs: int = 20, +) -> dict[str, Any]: + """Exact compare on website-requested fields. No tolerance / exemption.""" + fields = compare_fields(dataset) + if not fields: + fields = tuple(sorted({key for row in published + upstream for key in row if key != "batch_id"})) + pub_map = {str(row.get("ts_code") or "").upper(): row for row in published} + up_map = {str(row.get("ts_code") or "").upper(): row for row in upstream} + missing = sorted(code for code in pub_map if code not in up_map) + extra = sorted(code for code in up_map if code not in pub_map) + diffs: list[dict[str, Any]] = [] + for code in sorted(set(pub_map) & set(up_map)): + left = row_signature(pub_map[code], fields) + right = row_signature(up_map[code], fields) + if left == right: + continue + for field, old, new in zip(fields, left, right): + if old == new: + continue + diffs.append({"ts_code": code, "field": field, "published": old, "upstream": new}) + if len(diffs) >= max_diffs: + break + if len(diffs) >= max_diffs: + break + changed = bool(diffs or missing or extra) + return { + "changed": changed, + "dataset": dataset, + "fields": list(fields), + "published_rows": len(published), + "upstream_rows": len(upstream), + "missing_codes": missing[:max_diffs], + "extra_codes": extra[:max_diffs], + "diffs": diffs, + } + + +def official_table(dataset: str) -> str: + return DATASET_TABLES[dataset][0] diff --git a/xiaobai-datahub/datahub/scheduler.py b/xiaobai-datahub/datahub/scheduler.py index 51e2cb3..46bd3e3 100644 --- a/xiaobai-datahub/datahub/scheduler.py +++ b/xiaobai-datahub/datahub/scheduler.py @@ -1,5 +1,6 @@ from __future__ import annotations +import json import threading from collections.abc import Callable from datetime import datetime, time, timedelta @@ -8,13 +9,14 @@ from typing import Any from datahub.db import HubDB from datahub.logutil import get_logger from datahub.pipeline import Pipeline +from datahub.revision import revision_datasets from datahub.timeutil import isoformat, now_shanghai, yyyymmdd LOGGER = get_logger() JobFn = Callable[[str], Any] -EOD_JOB_IDS = {"eod_a", "eod_b", "eod_retry"} +EOD_JOB_IDS = {"eod_a", "eod_b", "eod_retry", "eod_revise"} def is_open_day(db: HubDB, day: str) -> bool: @@ -27,6 +29,20 @@ def is_open_day(db: HubDB, day: str) -> bool: return int(row["is_open"]) == 1 +def previous_open_day(db: HubDB, day: str) -> str | None: + row = db.fetchone( + """ + SELECT cal_date FROM trade_calendar + WHERE exchange = 'SSE' AND is_open = 1 AND cal_date < ? + ORDER BY cal_date DESC LIMIT 1 + """, + (day,), + ) + if row is None: + return None + return str(row["cal_date"]) + + def _hhmm(value: str) -> time: return datetime.strptime(value, "%H:%M").time() @@ -49,6 +65,7 @@ class Scheduler: "eod_a": self._eod_a, "eod_b": self._eod_b, "eod_retry": self._eod_retry, + "eod_revise": self._eod_revise, "stocks_refresh": self._stocks_refresh, "cleanup": self._cleanup, "backup": self._backup, @@ -118,6 +135,8 @@ class Scheduler: if job_id in {"eod_a", "eod_b"}: self._settle_eod(day) ran.extend(self._eod_retry_tick(now, day, open_day)) + ran.extend(self._revision_review_tick(now, day, open_day)) + ran.extend(self._revision_catchup_tick(now, day)) return ran # ------------------------------------------------------------------ @@ -219,6 +238,201 @@ class Scheduler: "detail": (row or {}).get("detail"), } + def revision_progress(self, day: str) -> dict[str, Any] | None: + return self.db.fetchone("SELECT * FROM revision_progress WHERE trade_date = ?", (day,)) + + def revision_status(self, trade_date: str | None = None, clock: datetime | None = None) -> dict[str, Any]: + """等待复核 / 复核失败 / 已追平 / 已截止.""" + day = yyyymmdd(trade_date or now_shanghai(clock)) + row = self.revision_progress(day) + open_day = is_open_day(self.db, day) + published = self._revision_ready(day) + if row and row["state"] in {"aligned", "review_failed", "cutoff", "waiting_review"}: + state = str(row["state"]) + elif not open_day: + state = "closed_day" + elif not published: + state = "pending_publish" + else: + state = "waiting_review" + return { + "trade_date": day, + "is_open_day": open_day, + "state": state, + "datasets": list(revision_datasets(self.pipeline.settings.quality)), + "attempts": int((row or {}).get("attempts") or 0), + "last_attempt_at": (row or {}).get("last_attempt_at"), + "next_retry_at": (row or {}).get("next_retry_at") if state in {"waiting_review", "review_failed"} else None, + "finished_at": (row or {}).get("finished_at"), + "catchup_done": bool(int((row or {}).get("catchup_done") or 0)), + "detail": (row or {}).get("detail"), + "window": f"{self.pipeline.settings.revision_review_start}-{self.pipeline.settings.revision_review_cutoff}", + } + + def _revision_ready(self, day: str) -> bool: + return all( + self.pipeline.active_batch(dataset, day) + for dataset in revision_datasets(self.pipeline.settings.quality) + ) + + def _revision_due(self, now: datetime, row: dict[str, Any] | None) -> bool: + if row is None or not row.get("last_attempt_at"): + return True + try: + last = datetime.fromisoformat(str(row["last_attempt_at"])) + except ValueError: + return True + interval = timedelta(minutes=self.pipeline.settings.revision_review_interval_minutes) + return now_shanghai(last).replace(tzinfo=None) + interval <= now.replace(tzinfo=None) + + def _revision_review_tick(self, now: datetime, day: str, open_day: bool) -> list[str]: + if not open_day or not self._revision_ready(day): + return [] + settings = self.pipeline.settings + current = now.time() + start = _hhmm(settings.revision_review_start) + cutoff = _hhmm(settings.revision_review_cutoff) + row = self.revision_progress(day) + if current < start: + if row is None: + self._save_revision_progress(day, state="waiting_review") + return [] + if current >= cutoff: + if row is None or row["state"] not in {"aligned", "cutoff"}: + detail = "复核窗口已截止" + self._save_revision_progress( + day, state="cutoff", finished_at=isoformat(now), detail=detail, + ) + with self.db.write() as connection: + connection.execute( + "INSERT INTO job_runs(job_id, state, started_at, finished_at, error, attempt, detail)" + " VALUES ('eod_revise','failed',?,?,?,?,?)", + ( + isoformat(now), isoformat(now), detail, + int((row or {}).get("attempts") or 0), "revision cutoff reached", + ), + ) + elif row["state"] == "aligned" and not row.get("finished_at"): + self._save_revision_progress(day, finished_at=isoformat(now)) + return [] + if not self._revision_due(now, row): + return [] + if "eod_revise" not in self.jobs: + return [] + return self._run_revision_job(day, now, catchup=False) + + def _revision_catchup_tick(self, now: datetime, day: str) -> list[str]: + prev = previous_open_day(self.db, day) + if prev is None or prev >= day: + return [] + if not self._revision_ready(prev): + return [] + row = self.revision_progress(prev) + if row and int(row.get("catchup_done") or 0): + return [] + if not self._revision_due(now, row): + return [] + if "eod_revise" not in self.jobs: + return [] + return self._run_revision_job(prev, now, catchup=True) + + def _run_revision_job(self, day: str, now: datetime, catchup: bool) -> list[str]: + attempts = int((self.revision_progress(day) or {}).get("attempts") or 0) + 1 + interval = self.pipeline.settings.revision_review_interval_minutes + self._save_revision_progress( + day, + state="waiting_review", + attempts=attempts, + last_attempt_at=isoformat(now), + next_retry_at=isoformat(now + timedelta(minutes=interval)), + ) + ran: list[str] = [] + try: + out = self.run_job("eod_revise", day) + except Exception as exc: + LOGGER.warning("revision review failed for %s: %s", day, exc) + self._save_revision_progress( + day, + state="review_failed", + detail="复核失败,保留上一完整版本", + ) + ran.append("eod_revise") + return ran + ran.append("eod_revise") + if out.get("state") == "skipped": + return ran + result = out.get("result") if isinstance(out.get("result"), dict) else {} + failed = [ + name for name, item in result.items() + if isinstance(item, dict) and item.get("state") == "failed" + ] + review = result.get("review") if isinstance(result.get("review"), dict) else None + watched = [ + result[name] + for name in revision_datasets(self.pipeline.settings.quality) + if isinstance(result.get(name), dict) + ] + diff_blob = None + if review and review.get("diffs"): + diff_blob = json.dumps(review.get("diffs"), ensure_ascii=False) + else: + for item in watched: + if item.get("diffs"): + diff_blob = json.dumps(item.get("diffs"), ensure_ascii=False) + break + revised = bool(review and review.get("reason") == "revised") + matched = any(item.get("reason") == "unchanged" or item.get("state") == "aligned" for item in watched) + if failed: + self._save_revision_progress( + day, + state="review_failed", + detail="复核失败,保留上一完整版本", + last_diff=diff_blob, + ) + elif revised or matched: + fields: dict[str, Any] = { + "state": "aligned", + "finished_at": isoformat(now), + "detail": "已追平" if revised else "已追平(无变化)", + "last_diff": diff_blob, + } + if catchup: + fields["catchup_done"] = 1 + self._save_revision_progress(day, **fields) + return ran + + def _save_revision_progress(self, day: str, **fields: Any) -> None: + columns = [ + "trade_date", "state", "attempts", "last_attempt_at", + "next_retry_at", "finished_at", "catchup_done", "last_diff", "detail", "updated_at", + ] + with self.db.write() as connection: + existing = connection.execute( + "SELECT trade_date FROM revision_progress WHERE trade_date = ?", + (day,), + ).fetchone() + if existing is None: + payload = {name: None for name in columns} + payload.update({ + "trade_date": day, + "state": "waiting_review", + "attempts": 0, + "catchup_done": 0, + }) + payload.update(fields) + payload["updated_at"] = isoformat() + placeholders = ",".join("?" for _ in columns) + connection.execute( + f"INSERT INTO revision_progress({','.join(columns)}) VALUES ({placeholders})", + tuple(payload[name] for name in columns), + ) + else: + assignments = ", ".join(f"{name} = ?" for name in fields) + connection.execute( + f"UPDATE revision_progress SET {assignments}, updated_at = ? WHERE trade_date = ?", + (*fields.values(), isoformat(), day), + ) + def _record_eod_attempt(self, day: str, now: datetime) -> None: row = self.eod_progress(day) attempts = int((row or {}).get("attempts") or 0) + 1 @@ -320,6 +534,9 @@ class Scheduler: def _eod_retry(self, trade_date: str) -> dict[str, Any]: return self.pipeline.run_eod_missing(trade_date) + def _eod_revise(self, trade_date: str) -> dict[str, Any]: + return self.pipeline.review_published_revisions(trade_date) + def _stocks_refresh(self, trade_date: str) -> dict[str, Any]: return self.pipeline.refresh_stocks(trade_date) diff --git a/xiaobai-datahub/datahub/settings.py b/xiaobai-datahub/datahub/settings.py index 3277899..6f27ef7 100644 --- a/xiaobai-datahub/datahub/settings.py +++ b/xiaobai-datahub/datahub/settings.py @@ -79,6 +79,20 @@ class Settings: def eod_retry_cutoff(self) -> str: return str(self.quality.get("eod_retry_cutoff") or "23:30") + @property + def revision_review_start(self) -> str: + # Before the 21:00 website shadow observation. + return str(self.quality.get("revision_review_start") or "20:00") + + @property + def revision_review_interval_minutes(self) -> int: + return int(self.quality.get("revision_review_interval_minutes") or 30) + + @property + def revision_review_cutoff(self) -> str: + # Last light review ~23:00; cutoff before the 23:30 observation. + return str(self.quality.get("revision_review_cutoff") or "23:20") + def load_settings( env: dict[str, str] | None = None, diff --git a/xiaobai-datahub/tests/test_eod_retry.py b/xiaobai-datahub/tests/test_eod_retry.py index 9b9f77f..b966625 100644 --- a/xiaobai-datahub/tests/test_eod_retry.py +++ b/xiaobai-datahub/tests/test_eod_retry.py @@ -111,15 +111,24 @@ class EodRetryTests(unittest.TestCase): self.assertEqual(progress["state"], "done") self.assertEqual(progress["attempts"], 4) # eod_a + eod_b + 2 retries - # success stops all further same-day requests + # success stops further eod_retry; revision window has not started yet batches_before = len(self._batches(db, day)) eod_calls_before = len(self._eod_calls(transport)) sched.tick(clock_at(day, 17, 0)) - sched.tick(clock_at(day, 23, 0)) self.assertEqual(len(self._job_runs(db, "eod_retry")), 2) + self.assertEqual(len(self._job_runs(db, "eod_revise")), 0) self.assertEqual(len(self._batches(db, day)), batches_before) self.assertEqual(len(self._eod_calls(transport)), eod_calls_before) + # 23:00 is inside the valuation review window: light daily_basic only, no new batch + sched.tick(clock_at(day, 23, 0)) + self.assertEqual(len(self._job_runs(db, "eod_retry")), 2) + self.assertEqual(len(self._job_runs(db, "eod_revise")), 1) + self.assertEqual(len(self._batches(db, day)), batches_before) + extra = [name for name in self._eod_calls(transport)[eod_calls_before:]] + self.assertTrue(extra) + self.assertTrue(all(name == "daily_basic" for name in extra)) + def test_never_ready_marks_cutoff_failed_and_stops(self) -> None: day = "20240902" db, transport, pipe, sched = self._make(set()) @@ -174,6 +183,7 @@ class EodRetryTests(unittest.TestCase): self.assertIn("eod_a", ran) self.assertIn("eod_b", ran) self.assertNotIn("eod_retry", ran) + self.assertIn("eod_revise", ran) self.assertEqual(self._published(db, day), OFFICIAL) after = db.fetchall("SELECT dataset, active_batch FROM publications WHERE trade_date = ?", (day,)) self.assertEqual( @@ -181,7 +191,9 @@ class EodRetryTests(unittest.TestCase): active_map, ) self.assertEqual(set(official_batches()), batches_before) # no duplicate batches - self.assertEqual(self._eod_calls(transport), calls_before) # no duplicate upstream EOD calls + extra = self._eod_calls(transport)[len(calls_before):] + self.assertTrue(extra) + self.assertTrue(all(name == "daily_basic" for name in extra)) self.assertEqual(sched2.eod_status(day, clock=clock_at(day, 21, 0))["state"], "done") def test_restart_with_partial_publish_only_fetches_missing(self) -> None: @@ -206,6 +218,7 @@ class EodRetryTests(unittest.TestCase): for hh, mm in ((15, 5), (15, 10), (15, 40), (16, 10), (20, 0), (23, 40)): ran = sched.tick(clock_at(day, hh, mm)) self.assertNotIn("eod_retry", ran) + self.assertNotIn("eod_revise", ran) eod_runs = db.fetchall("SELECT * FROM job_runs WHERE job_id LIKE 'eod%'") self.assertEqual(eod_runs, []) self.assertIsNone(db.fetchone("SELECT * FROM eod_progress WHERE trade_date = ?", (day,))) @@ -226,12 +239,18 @@ class EodRetryTests(unittest.TestCase): self.assertEqual(len(self._batches(db, day)), batches_before) self.assertEqual(len(transport.calls), calls_before) + revised = sched.run_job("eod_revise", day) + self.assertEqual(revised["state"], "ok") + self.assertEqual(len(self._batches(db, day)), batches_before) + sched._eod_lock.acquire() # simulate an in-flight EOD job try: busy = sched.run_job("eod_retry", day) self.assertEqual(busy["state"], "skipped") busy_a = sched.run_job("eod_a", day) self.assertEqual(busy_a["state"], "skipped") + busy_r = sched.run_job("eod_revise", day) + self.assertEqual(busy_r["state"], "skipped") finally: sched._eod_lock.release() self.assertEqual(len(self._batches(db, day)), batches_before) diff --git a/xiaobai-datahub/tests/test_revision_review.py b/xiaobai-datahub/tests/test_revision_review.py new file mode 100644 index 0000000..02eae2c --- /dev/null +++ b/xiaobai-datahub/tests/test_revision_review.py @@ -0,0 +1,333 @@ +from __future__ import annotations + +import copy +import unittest +from pathlib import Path +import tempfile + +from datahub.adapters.base import AdapterError +from datahub.adapters.tushare import TushareAdapter +from datahub.crypto import SecretVault +from datahub.db import HubDB +from datahub.pipeline import Pipeline +from datahub.scheduler import Scheduler +from datahub.serving import V1API +from datahub.settings import Settings +from datahub.timeutil import SHANGHAI +from tests.fixtures import RAW, TRADE_DATE, fake_transport +from tests.test_eod_retry import clock_at +from tests.test_quality_gates import FIELD_GATES + +SAMPLE_DAY = "20260907" +NEXT_DAY = "20260908" +SAMPLE_CODE = "003021.SZ" + + +def _dated(row: dict, day: str) -> dict: + item = dict(row) + if "trade_date" in item: + item["trade_date"] = day + return item + + +class RevisingTransport: + """Fixture transport that can rewrite daily_basic after the first publish.""" + + DATE_APIS = {"daily", "daily_basic", "adj_factor", "moneyflow", "stk_auction", "index_daily"} + + def __init__(self, extra_calendar: list[dict] | None = None) -> None: + self.calls: list[str] = [] + self.fail_daily_basic = False + self.empty_daily_basic = False + self.null_volume_ratio = False + self.turnover_by_code: dict[str, float] = {} + self.extra_calendar = extra_calendar or [] + + def __call__(self, api_name: str, params: dict, fields: str): + self.calls.append(api_name) + day = str(params.get("trade_date") or "") + if api_name == "trade_cal": + rows = fake_transport(api_name, params, fields) + extra = [ + row for row in self.extra_calendar + if str(params.get("start_date") or "") <= row["cal_date"] <= str(params.get("end_date") or "99999999") + ] + return rows + extra + if self.fail_daily_basic and api_name == "daily_basic": + raise AdapterError("tushare daily_basic unavailable") + if self.empty_daily_basic and api_name == "daily_basic": + return [] + if api_name == "index_daily": + code = params.get("ts_code") + rows = [row for row in RAW["index_daily"] if row["ts_code"] == code] + if day: + rows = [_dated(row, day) for row in rows] + return rows + rows = fake_transport(api_name, params, fields) + if api_name == "stock_basic": + rows = list(rows) + rows.append({ + "ts_code": SAMPLE_CODE, "symbol": "003021", "name": "兆威机电", + "area": "广东", "industry": "元器件", "market": "主板", + "list_status": "L", "list_date": "20201202", + }) + return rows + if api_name in self.DATE_APIS: + template = RAW.get(api_name) or [] + if not day: + return [_dated(row, TRADE_DATE) for row in template] + out = [_dated(row, day) for row in template] + extra = copy.deepcopy(template[0]) + extra["ts_code"] = SAMPLE_CODE + extra["trade_date"] = day + if api_name == "daily_basic": + extra["turnover_rate"] = self.turnover_by_code.get(SAMPLE_CODE, extra.get("turnover_rate")) + if self.null_volume_ratio: + extra["volume_ratio"] = None + for row in out: + row["volume_ratio"] = None + out.append(extra) + if api_name == "daily_basic": + for row in out: + code = str(row.get("ts_code") or "") + if code in self.turnover_by_code: + row["turnover_rate"] = self.turnover_by_code[code] + return out + return rows + + +def make_revision_env(quality_extra: dict | None = None, extra_calendar: list[dict] | None = None): + tmp = tempfile.TemporaryDirectory() + db = HubDB(Path(tmp.name) / "hub.db") + transport = RevisingTransport(extra_calendar=extra_calendar) + adapter = TushareAdapter("x", transport=transport) + quality = { + "daily_row_ratio": 0.5, + "null_rate_max": 0.5, + "max_publish_attempts": 2, + "publication_generations": 3, + "field_gates": FIELD_GATES, + "revision_review_start": "20:00", + "revision_review_interval_minutes": 30, + "revision_review_cutoff": "23:20", + "revision_review_datasets": ["valuation"], + } + if quality_extra: + quality.update(quality_extra) + settings = Settings( + encryption_key=SecretVault.generate_key(), + api_token="t" * 32, + db_path=db.path, + backup_dir=Path(tmp.name) / "backups", + quality=quality, + scheduler_enabled=False, + ) + pipe = Pipeline(db, adapter, settings) + sched = Scheduler(db, pipe) + return tmp, db, transport, pipe, sched + + +SAMPLE_CALENDAR = [ + {"exchange": "SSE", "cal_date": SAMPLE_DAY, "is_open": 1, "pretrade_date": "20260906"}, + {"exchange": "SSE", "cal_date": NEXT_DAY, "is_open": 1, "pretrade_date": SAMPLE_DAY}, +] + + +class RevisionReviewTests(unittest.TestCase): + def _publish(self, pipe: Pipeline, day: str) -> None: + pipe.ingest_reference(day) + pipe.run_eod_batch_a(day) + pipe.run_eod_batch_b(day) + + def _turnover(self, db: HubDB, day: str, code: str = SAMPLE_CODE) -> float | None: + pub = db.fetchone( + "SELECT active_batch FROM publications WHERE dataset='valuation' AND trade_date=?", + (day,), + ) + row = db.fetchone( + "SELECT turnover_rate FROM eod_valuation WHERE batch_id=? AND ts_code=?", + (pub["active_batch"], code), + ) + return None if row is None else row["turnover_rate"] + + def _batch_ids(self, db: HubDB, day: str) -> set[str]: + return {str(row["batch_id"]) for row in db.fetchall("SELECT batch_id FROM batches WHERE trade_date=?", (day,))} + + def test_no_change_does_not_create_a_new_batch(self) -> None: + tmp, db, transport, pipe, sched = make_revision_env() + self.addCleanup(tmp.cleanup) + self._publish(pipe, TRADE_DATE) + before = self._batch_ids(db, TRADE_DATE) + sched.tick(clock_at(TRADE_DATE, 20, 0)) + self.assertEqual(self._batch_ids(db, TRADE_DATE), before) + progress = db.fetchone("SELECT * FROM revision_progress WHERE trade_date=?", (TRADE_DATE,)) + self.assertEqual(progress["state"], "aligned") + self.assertIn("无变化", progress["detail"]) + status = sched.revision_status(TRADE_DATE, clock=clock_at(TRADE_DATE, 20, 0)) + self.assertEqual(status["state"], "aligned") + + def test_hel423_20260907_single_field_revision_is_caught_up(self) -> None: + tmp, db, transport, pipe, sched = make_revision_env(extra_calendar=SAMPLE_CALENDAR) + self.addCleanup(tmp.cleanup) + transport.turnover_by_code[SAMPLE_CODE] = 1.3565 + self._publish(pipe, SAMPLE_DAY) + self.assertEqual(self._turnover(db, SAMPLE_DAY), 1.3565) + first = db.fetchone( + "SELECT active_batch FROM publications WHERE dataset='valuation' AND trade_date=?", + (SAMPLE_DAY,), + )["active_batch"] + + transport.turnover_by_code[SAMPLE_CODE] = 1.3572 + seen: list[float | None] = [] + + def watch() -> None: + seen.append(self._turnover(db, SAMPLE_DAY)) + + pipe.before_commit = watch + ran = sched.tick(clock_at(SAMPLE_DAY, 20, 0)) + self.assertIn("eod_revise", ran) + self.assertEqual(seen, [1.3565]) # readers still see the previous complete version mid-switch + self.assertEqual(self._turnover(db, SAMPLE_DAY), 1.3572) + second = db.fetchone( + "SELECT active_batch FROM publications WHERE dataset='valuation' AND trade_date=?", + (SAMPLE_DAY,), + )["active_batch"] + self.assertNotEqual(second, first) + api = V1API(db, pipe, pipe.settings) + payload = api.valuation({"date": SAMPLE_DAY, "code": SAMPLE_CODE}) + row = next(item for item in payload["data"] if item["ts_code"] == SAMPLE_CODE) + self.assertEqual(row["turnover_rate"], 1.3572) + progress = db.fetchone("SELECT * FROM revision_progress WHERE trade_date=?", (SAMPLE_DAY,)) + self.assertEqual(progress["state"], "aligned") + self.assertEqual(progress["detail"], "已追平") + audit = db.fetchone( + "SELECT * FROM audit_log WHERE action='revision-review' ORDER BY id DESC" + ) + self.assertIn("1.3572", str(audit["detail"])) + self.assertIn(SAMPLE_CODE, str(audit["detail"])) + + def test_empty_or_failed_upstream_keeps_previous_version(self) -> None: + tmp, db, transport, pipe, sched = make_revision_env() + self.addCleanup(tmp.cleanup) + self._publish(pipe, TRADE_DATE) + active = db.fetchone( + "SELECT active_batch FROM publications WHERE dataset='valuation' AND trade_date=?", + (TRADE_DATE,), + )["active_batch"] + batches = self._batch_ids(db, TRADE_DATE) + + transport.empty_daily_basic = True + sched.tick(clock_at(TRADE_DATE, 20, 0)) + self.assertEqual( + db.fetchone( + "SELECT active_batch FROM publications WHERE dataset='valuation' AND trade_date=?", + (TRADE_DATE,), + )["active_batch"], + active, + ) + self.assertEqual( + db.fetchone("SELECT state FROM revision_progress WHERE trade_date=?", (TRADE_DATE,))["state"], + "review_failed", + ) + + transport.empty_daily_basic = False + transport.fail_daily_basic = True + sched.tick(clock_at(TRADE_DATE, 20, 30)) + self.assertEqual( + db.fetchone( + "SELECT active_batch FROM publications WHERE dataset='valuation' AND trade_date=?", + (TRADE_DATE,), + )["active_batch"], + active, + ) + self.assertEqual(self._batch_ids(db, TRADE_DATE), batches) + + def test_quality_gate_rejects_catchup_and_keeps_previous(self) -> None: + tmp, db, transport, pipe, sched = make_revision_env() + self.addCleanup(tmp.cleanup) + self._publish(pipe, TRADE_DATE) + active = db.fetchone( + "SELECT active_batch FROM publications WHERE dataset='valuation' AND trade_date=?", + (TRADE_DATE,), + )["active_batch"] + transport.turnover_by_code[SAMPLE_CODE] = 9.9999 + transport.null_volume_ratio = True + sched.tick(clock_at(TRADE_DATE, 20, 0)) + self.assertEqual( + db.fetchone( + "SELECT active_batch FROM publications WHERE dataset='valuation' AND trade_date=?", + (TRADE_DATE,), + )["active_batch"], + active, + ) + self.assertEqual( + db.fetchone("SELECT state FROM revision_progress WHERE trade_date=?", (TRADE_DATE,))["state"], + "review_failed", + ) + + def test_repeat_ticks_after_align_do_not_republish(self) -> None: + tmp, db, transport, pipe, sched = make_revision_env() + self.addCleanup(tmp.cleanup) + transport.turnover_by_code[SAMPLE_CODE] = 1.3565 + self._publish(pipe, TRADE_DATE) + transport.turnover_by_code[SAMPLE_CODE] = 1.3572 + sched.tick(clock_at(TRADE_DATE, 20, 0)) + after_fix = self._batch_ids(db, TRADE_DATE) + sched.tick(clock_at(TRADE_DATE, 20, 10)) # inside interval + self.assertEqual(len(db.fetchall("SELECT * FROM job_runs WHERE job_id='eod_revise'")), 1) + sched.tick(clock_at(TRADE_DATE, 20, 30)) # next light compare, no change + self.assertEqual(self._batch_ids(db, TRADE_DATE), after_fix) + self.assertEqual(self._turnover(db, TRADE_DATE), 1.3572) + + def test_restart_catches_up_inside_window(self) -> None: + tmp, db, transport, pipe, sched = make_revision_env() + self.addCleanup(tmp.cleanup) + transport.turnover_by_code[SAMPLE_CODE] = 1.3565 + self._publish(pipe, TRADE_DATE) + transport.turnover_by_code[SAMPLE_CODE] = 1.3572 + sched2 = Scheduler(db, pipe) + ran = sched2.tick(clock_at(TRADE_DATE, 21, 0)) + self.assertIn("eod_revise", ran) + self.assertEqual(self._turnover(db, TRADE_DATE), 1.3572) + + def test_cutoff_stops_evening_reviews_and_morning_catchup_runs(self) -> None: + tmp, db, transport, pipe, sched = make_revision_env(extra_calendar=SAMPLE_CALENDAR) + self.addCleanup(tmp.cleanup) + transport.turnover_by_code[SAMPLE_CODE] = 1.3565 + self._publish(pipe, SAMPLE_DAY) + sched.tick(clock_at(SAMPLE_DAY, 23, 25)) # past 23:20 cutoff, no review yet + cutoff = db.fetchone("SELECT * FROM revision_progress WHERE trade_date=?", (SAMPLE_DAY,)) + self.assertEqual(cutoff["state"], "cutoff") + self.assertEqual(self._turnover(db, SAMPLE_DAY), 1.3565) + + transport.turnover_by_code[SAMPLE_CODE] = 1.3572 + sched.tick(clock_at(SAMPLE_DAY, 23, 50)) # still same calendar day, no catch-up + self.assertEqual(self._turnover(db, SAMPLE_DAY), 1.3565) + + ran = sched.tick(clock_at(NEXT_DAY, 8, 45)) + self.assertIn("eod_revise", ran) + self.assertEqual(self._turnover(db, SAMPLE_DAY), 1.3572) + progress = db.fetchone("SELECT * FROM revision_progress WHERE trade_date=?", (SAMPLE_DAY,)) + self.assertEqual(progress["state"], "aligned") + self.assertEqual(int(progress["catchup_done"]), 1) + + batches = self._batch_ids(db, SAMPLE_DAY) + sched.tick(clock_at(NEXT_DAY, 8, 50)) + self.assertEqual(self._batch_ids(db, SAMPLE_DAY), batches) + + def test_only_valuation_is_light_fetched(self) -> None: + tmp, db, transport, pipe, sched = make_revision_env() + self.addCleanup(tmp.cleanup) + self._publish(pipe, TRADE_DATE) + before = [name for name in transport.calls] + sched.tick(clock_at(TRADE_DATE, 20, 0)) + extra = transport.calls[len(before):] + self.assertIn("daily_basic", extra) + self.assertNotIn("daily", extra) + self.assertNotIn("moneyflow", extra) + self.assertNotIn("stk_auction", extra) + self.assertNotIn("index_daily", extra) + + +if __name__ == "__main__": + unittest.main()