feat(HEL-478): 估值发布后晚间复核并原子追补上游修订

盘后成功发布后继续轻量比对 daily_basic 网站字段,发现修订才走质量门与整组原子切换,避免 17:10 快照落后于晚间上游改写。

Co-authored-by: Cursor <cursoragent@cursor.com>
Co-authored-by: multica-agent <github@multica.ai>
This commit is contained in:
总工
2026-09-07 21:19:35 +08:00
co-authored by Cursor multica-agent
parent 16ba83ec01
commit d175bb65d4
11 changed files with 982 additions and 4 deletions
+2
View File
@@ -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 / 日志"},
+13
View File
@@ -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,
+240
View File
@@ -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, ...],
+111
View File
@@ -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]
+218 -1
View File
@@ -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)
+14
View File
@@ -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,