fix(HEL-494): 数据中枢独占调度,主网站不再回退旧接口

主网站只向中枢要业务数据;来源选择、切源、补数全部在中枢内部完成,失败不再走东财/腾讯/Tushare 保底。

Co-authored-by: Cursor <cursoragent@cursor.com>
Co-authored-by: multica-agent <github@multica.ai>
This commit is contained in:
总工
2026-09-08 21:43:31 +08:00
co-authored by Cursor multica-agent
parent ef13d6feb5
commit 0b8419abca
23 changed files with 1159 additions and 368 deletions
+80 -47
View File
@@ -19,6 +19,7 @@ from backend.data.datahub.redact import redact_text, redact_value
from backend.data.datahub.route_state import LEDGER
from backend.data.datahub.settings import DatahubSettings
from backend.data.providers.tushare_client import TushareClient
from backend.data.providers.tushare_transport import TushareError
LOGGER = logging.getLogger("xiaobai.datahub")
ShadowSink = Callable[[dict[str, Any]], None]
@@ -171,6 +172,45 @@ class DatahubBridge:
self._log_failure("index_quotes", exc)
return None
def try_sector_quote(self, code: str, trade_date: str = "") -> dict[str, Any] | None:
flags = self.settings.flags("quotes")
if not flags.read:
return None
try:
response = self.client.sector_quote(code, trade_date)
data = response.data
if not isinstance(data, dict) or not data:
raise DatahubError("EMPTY", "datahub sector quote empty")
row = dict(data)
if (response.meta or {}).get("stale"):
row["delayed"] = True
row["delay_seconds"] = int((response.meta or {}).get("staleness_seconds") or 0)
row["delay_notice"] = str((response.meta or {}).get("delay_notice") or "")
self._record_route("quotes", "datahub", str((response.meta or {}).get("source") or "datahub"))
return row
except Exception as exc:
self._log_failure("quotes", exc)
return None
def try_limit_pool(self, trade_date: str = "") -> list[dict[str, Any]] | None:
flags = self.settings.flags("limit_events")
if not flags.read:
return None
try:
response = self.client.limit_pool(trade_date)
rows = [dict(item) for item in (response.data or []) if isinstance(item, dict)]
if not rows:
raise DatahubError("EMPTY", "datahub limit pool empty")
self._record_route(
"limit_events",
"datahub",
str((response.meta or {}).get("source") or "datahub"),
)
return rows
except Exception as exc:
self._log_failure("limit_events", exc)
return None
def try_daily_chart(
self,
code: str,
@@ -191,6 +231,11 @@ class DatahubBridge:
self.client.index_bars,
{"code": code, "from": start, "to": compact_end},
)
elif dataset == "sector_daily":
response = self._paginate(
self.client.sectors,
{"code": code, "from": start, "to": compact_end},
)
else:
response = self._paginate(
self.client.daily_bars,
@@ -263,53 +308,33 @@ class DatahubBridge:
fields: str,
legacy_query: Callable[..., list[dict[str, Any]]],
) -> list[dict[str, Any]]:
del legacy_query # 主网站不再直连 Tushare;调度全部由数据中枢完成。
if api_name == "rt_sw_k":
raise TushareError("rt_sw_k is disabled; use published sw_daily or free Shenwan realtime")
dataset = API_TO_DATASET.get(api_name)
# 问天按实际数据依赖接入:已映射到 hub 的 API 跟随开关;未映射的继续旧链路。
if not dataset:
return legacy_query(api_name, params, fields)
flags = self.settings.flags(dataset)
if not flags.read and not flags.shadow:
return legacy_query(api_name, params, fields)
hub_rows: list[dict[str, Any]] | None = None
hub_meta: dict[str, Any] = {}
hub_error: str | None = None
hub_canonical: list[dict[str, Any]] = []
try:
response = self._fetch_dataset(dataset, params or {}, api_name=api_name)
hub_canonical = self._extract_rows(dataset, response, params or {})
hub_rows = to_native_rows(dataset, hub_canonical)
hub_meta = dict(response.meta)
self._validate_usable(dataset, hub_rows, response)
except Exception as exc:
hub_error = self._error_text(exc)
self._log_failure(dataset, exc)
if flags.shadow:
try:
legacy_rows = legacy_query(api_name, params, fields)
except Exception as exc:
if flags.read and hub_rows is not None and hub_error is None:
self._emit_shadow(
compare_rows(dataset, [], hub_canonical, hub_meta, self._error_text(exc), fields)
)
return project_fields(hub_rows, fields)
raise
self._emit_shadow(compare_rows(dataset, legacy_rows, hub_canonical, hub_meta, hub_error, fields))
if flags.read and hub_rows is not None and hub_error is None:
self._record_route(dataset, "datahub", str(hub_meta.get("source") or "datahub"))
return project_fields(hub_rows, fields)
if dataset:
flags = self.settings.flags(dataset)
if flags.read:
self._record_route(dataset, "legacy", "tushare", hub_error or "")
return legacy_rows
if flags.read and hub_rows is not None and hub_error is None:
self._record_route(dataset, "datahub", str(hub_meta.get("source") or "datahub"))
return project_fields(hub_rows, fields)
result = legacy_query(api_name, params, fields)
if flags.read:
self._record_route(dataset, "legacy", "tushare", hub_error or "")
return result
try:
response = self._fetch_dataset(dataset, params or {}, api_name=api_name)
hub_canonical = self._extract_rows(dataset, response, params or {})
hub_rows = to_native_rows(dataset, hub_canonical)
self._validate_usable(dataset, hub_rows, response)
self._record_route(dataset, "datahub", str(response.meta.get("source") or "datahub"))
return project_fields(hub_rows, fields)
except Exception as exc:
self._log_failure(dataset, exc)
try:
response = self.client.query_api(api_name, params or {}, fields)
rows = [dict(item) for item in (response.data or []) if isinstance(item, dict)]
if dataset:
self._record_route(dataset, "datahub", str((response.meta or {}).get("source") or "datahub"))
else:
self._record_route(api_name, "datahub", str((response.meta or {}).get("source") or "datahub"))
return rows if not fields else project_fields(rows, fields)
except Exception as exc:
self._log_failure(dataset or api_name, exc)
raise TushareError(self._error_text(exc)) from exc
def _fetch_dataset(self, dataset: str, params: dict[str, Any], api_name: str = "") -> DatahubResponse:
date = yyyymmdd(params.get("trade_date") or params.get("date"))
@@ -429,8 +454,8 @@ class DatahubBridge:
def _log_failure(self, dataset: str, exc: Exception) -> None:
error = redact_text(self._error_text(exc), self.settings.secrets())
LOGGER.warning("datahub fallback dataset=%s error=%s", dataset, error)
self._record_route(dataset, "legacy", "pending-legacy", error)
LOGGER.warning("datahub unavailable dataset=%s error=%s", dataset, error)
self._record_route(dataset, "datahub", "unavailable", error)
def _record_route(self, dataset: str, route: str, source: str = "", error: str = "") -> None:
LEDGER.record(dataset, route, source, redact_text(error, self.settings.secrets()))
@@ -529,6 +554,8 @@ class DatahubAwareTushareClient:
legacy.try_market_quotes = self.try_market_quotes
legacy.try_quotes = self.try_quotes
legacy.try_index_quotes = self.try_index_quotes
legacy.try_sector_quote = self.try_sector_quote
legacy.try_limit_pool = self.try_limit_pool
legacy.record_datahub_legacy = self.record_datahub_legacy
def query(
@@ -548,6 +575,12 @@ class DatahubAwareTushareClient:
def try_index_quotes(self) -> list[dict[str, Any]] | None:
return self._bridge.try_index_quotes()
def try_sector_quote(self, code: str, trade_date: str = "") -> dict[str, Any] | None:
return self._bridge.try_sector_quote(code, trade_date)
def try_limit_pool(self, trade_date: str = "") -> list[dict[str, Any]] | None:
return self._bridge.try_limit_pool(trade_date)
def record_datahub_legacy(self, dataset: str, source: str = "", error: str = "") -> None:
self._bridge.record_legacy(dataset, source, error)
+51 -7
View File
@@ -90,6 +90,24 @@ class DatahubClient:
params["dataset"] = dataset
return self.get("/v1/batches", params)
def query_api(self, api_name: str, params: dict[str, Any] | None = None, fields: str = "") -> DatahubResponse:
return self.post(
"/v1/query",
{"api_name": api_name, "params": params or {}, "fields": fields},
)
def sector_quote(self, code: str, date: str = "") -> DatahubResponse:
payload: dict[str, Any] = {"code": code}
if date:
payload["date"] = date
return self.get("/v1/sectors/quote", payload)
def limit_pool(self, trade_date: str = "") -> DatahubResponse:
params: dict[str, Any] = {}
if trade_date:
params["date"] = trade_date
return self.get("/v1/limit-pool", params)
def get(self, path: str, params: dict[str, Any] | None = None) -> DatahubResponse:
if not self.settings.token:
raise DatahubError("NOT_CONFIGURED", "DATAHUB_TOKEN is not configured")
@@ -118,15 +136,41 @@ class DatahubClient:
)
raise last_error or DatahubError("INTERNAL", "datahub request failed")
def _request(self, url: str) -> DatahubResponse:
def post(self, path: str, body: dict[str, Any] | None = None) -> DatahubResponse:
if not self.settings.token:
raise DatahubError("NOT_CONFIGURED", "DATAHUB_TOKEN is not configured")
url = self.settings.base_url + path
attempts = 1 + max(0, self.settings.retries)
last_error: DatahubError | None = None
payload = json.dumps(body or {}, ensure_ascii=False).encode("utf-8")
for attempt in range(attempts):
try:
return self._request(url, method="POST", data=payload)
except DatahubError as exc:
last_error = exc
if exc.code not in {"TIMEOUT", "UNAVAILABLE"} or attempt + 1 >= attempts:
raise
LOGGER.warning(
"datahub retry %s/%s %s",
attempt + 1,
attempts,
redact_text(str(exc), self.settings.secrets()),
)
raise last_error or DatahubError("INTERNAL", "datahub request failed")
def _request(self, url: str, method: str = "GET", data: bytes | None = None) -> DatahubResponse:
headers = {
"Accept": "application/json",
"X-Datahub-Token": self.settings.token,
"User-Agent": "XiaobaiReviewDatahub/1.0",
}
if data is not None:
headers["Content-Type"] = "application/json"
request = urllib.request.Request(
url,
headers={
"Accept": "application/json",
"X-Datahub-Token": self.settings.token,
"User-Agent": "XiaobaiReviewDatahub/1.0",
},
method="GET",
data=data,
headers=headers,
method=method,
)
try:
with self._urlopen(request, timeout=self.settings.timeout_seconds) as response:
+1 -1
View File
@@ -38,7 +38,7 @@ class DataGateway:
if dataset_id:
self.policy.assert_allowed(dataset_id, "tushare", usage)
legacy = self.tushare_provider.client()
legacy.realtime_aggregator = self.realtime_observer
legacy.realtime_aggregator = None
return DatahubAwareTushareClient(legacy, self.datahub)
def dataset_status(self, trade_date: str) -> list[dict[str, Any]] | None:
+2 -3
View File
@@ -186,12 +186,11 @@ class DailyMarketMixin:
return mapped
def _free_board_map(self, trade_date: str) -> dict[str, dict[str, Any]]:
aggregator = getattr(self, "realtime_aggregator", None)
loader = getattr(aggregator, "eastmoney_limit_pool", None) if aggregator else None
loader = getattr(self, "try_limit_pool", None)
if not callable(loader):
return {}
try:
rows = loader(trade_date)
rows = loader(trade_date) or []
except Exception:
return {}
return {
+16 -61
View File
@@ -251,27 +251,23 @@ class DashboardMixin:
quotes = hub(trade_date)
if quotes:
return list(quotes), "datahub"
rt_error = ""
named = getattr(self, "try_quotes", None)
code_list = [item for item in str(codes or "").split(",") if item]
if callable(named) and code_list:
collected: list[dict[str, Any]] = []
for index in range(0, len(code_list), 60):
collected.extend(named(code_list[index:index + 60]) or [])
if collected:
delayed = any(item.get("delayed") for item in collected)
return collected, "datahub_delayed" if delayed else "datahub"
try:
quotes = self.query("rt_k", {"ts_code": codes})
if quotes:
self._mark_quote_legacy("tushare_rt_k", rt_error)
return list(quotes), "tushare_rt_k"
rt_error = f"No realtime data returned for {trade_date}"
delayed = any(item.get("delayed") for item in quotes)
return list(quotes), "datahub_delayed" if delayed else "datahub"
except TushareError as exc:
rt_error = str(exc)
try:
quotes, quote_source = self._free_realtime_quotes(trade_date, codes)
except Exception as exc:
raise TushareError(
f"当天盘中实时行情不可用:rt_k={rt_error};免费源={exc}"
) from exc
if not quotes:
raise TushareError(
f"当天盘中实时行情不可用:rt_k={rt_error};免费源=empty"
)
self._mark_quote_legacy(quote_source, rt_error)
return quotes, quote_source
raise TushareError(f"当天盘中实时行情不可用:{exc}") from exc
raise TushareError("当天盘中实时行情不可用:数据中枢未返回可用行情")
def _mark_quote_legacy(self, source: str, error: str = "") -> None:
marker = getattr(self, "record_datahub_legacy", None)
@@ -283,27 +279,8 @@ class DashboardMixin:
trade_date: str,
codes: str = "",
) -> tuple[list[dict[str, Any]], str]:
aggregator = self._realtime_aggregator()
last_error = ""
try:
quotes = aggregator.eastmoney_market_quotes(expected_date=trade_date)
if quotes:
return quotes, "eastmoney_clist"
except Exception as exc:
last_error = str(exc)
code_list = [item for item in str(codes or "").split(",") if item]
try:
if code_list:
quotes = aggregator.tencent_stock_quotes(code_list, expected_date=trade_date)
else:
quotes = aggregator.tencent_market_quotes(code_list, expected_date=trade_date)
except Exception as exc:
raise TushareError(
f"eastmoney={last_error or 'empty'}tencent={exc}"
) from exc
if not quotes:
raise TushareError(f"eastmoney={last_error or 'empty'}tencent=empty")
return quotes, "tencent_qt"
del trade_date, codes
raise TushareError("主网站不再直连免费行情源,请走数据中枢")
def _free_realtime_indices(self) -> list[dict[str, Any]]:
hub = getattr(self, "try_index_quotes", None)
@@ -312,14 +289,7 @@ class DashboardMixin:
converted = [item for item in (_hub_index_quote(row) for row in rows or []) if item]
if converted:
return converted
try:
rows = self._realtime_aggregator().eastmoney_indices()
marker = getattr(self, "record_datahub_legacy", None)
if callable(marker):
marker("index_quotes", "eastmoney_push2")
return rows
except Exception:
return []
return []
def _load_realtime_reference(
self,
@@ -470,21 +440,6 @@ class DashboardMixin:
return dict(rows[0])
except TushareError:
pass
aggregator = getattr(self, "realtime_aggregator", None)
if aggregator is None:
return {}
for loader in (
getattr(aggregator, "eastmoney_stock_quote", None),
getattr(aggregator, "tencent_stock_quote", None),
):
if not callable(loader):
continue
try:
quote = loader(ts_code, expected_date=reference_date)
except Exception:
continue
if quote:
return dict(quote)
return {}
def _stock_activity_metrics(
+4 -63
View File
@@ -63,22 +63,8 @@ class IndexMixin:
if callable(hub):
rows = hub()
if rows:
try:
return self._hub_realtime_market_indices(requested_date, rows)
except TushareError:
pass
try:
payload = self._tushare_realtime_market_indices(requested_date)
marker = getattr(self, "record_datahub_legacy", None)
if callable(marker):
marker("index_quotes", "tushare_rt_idx_k")
return payload
except TushareError:
payload = self._free_realtime_market_indices(requested_date)
marker = getattr(self, "record_datahub_legacy", None)
if callable(marker):
marker("index_quotes", str(payload.get("source") or "eastmoney_push2"))
return payload
return self._hub_realtime_market_indices(requested_date, rows)
raise TushareError("Realtime index quotes are incomplete")
def _hub_realtime_market_indices(
self,
@@ -199,50 +185,5 @@ class IndexMixin:
}
def _free_realtime_market_indices(self, requested_date: str) -> dict[str, Any]:
trade_date, _ = self.resolve_trade_context(requested_date)
aggregator = getattr(self, "realtime_aggregator", None)
if aggregator is None:
raise TushareError("免费实时源未配置")
quotes = aggregator.eastmoney_indices()
index_names = {
"000001": ("000001.SH", "上证指数"),
"399001": ("399001.SZ", "深证成指"),
"399006": ("399006.SZ", "创业板指"),
}
indices = []
for quote in quotes:
mapped = index_names.get(str(quote.get("code") or ""))
if not mapped:
continue
ts_code, name = mapped
close = _number(quote.get("price"))
previous_close = _number(quote.get("previous_close"))
if close <= 0 or previous_close <= 0:
continue
indices.append(
{
"ts_code": ts_code,
"name": str(quote.get("name") or name).strip(),
"trade_date": trade_date,
"close": close,
"pct_chg": round(_number(quote.get("change")) or (close / previous_close - 1) * 100, 3),
"return_5d": 0,
"amount_billion": round(_number(quote.get("amount_billion")), 2),
"quote_time": quote.get("quote_time") or "",
"source": quote.get("source") or "eastmoney_push2",
}
)
if len(indices) != 3:
raise TushareError("Realtime index quotes are incomplete")
return {
"trade_date": trade_date,
"source": "eastmoney_push2",
"realtime": True,
"precise": True,
"indices": indices,
"aggregate": {
"average_pct_chg": round(sum(item["pct_chg"] for item in indices) / len(indices), 3),
"average_return_5d": 0,
"average_return_20d": 0,
},
}
del requested_date
raise TushareError("主网站不再直连免费行情源,请走数据中枢")
+5 -38
View File
@@ -618,27 +618,19 @@ class ShenwanIndustryMixin:
trade_date: str,
finalized: bool = False,
) -> tuple[dict[str, Any], str, str]:
aggregator = getattr(self, "realtime_aggregator", None)
loader = getattr(aggregator, "eastmoney_shenwan_quote", None) if aggregator else None
if callable(loader):
hub = getattr(self, "try_sector_quote", None)
if callable(hub):
try:
row = loader(sector_code, expected_date="" if finalized else trade_date)
row = hub(sector_code, "" if finalized else trade_date)
except Exception as exc:
message = str(exc)
if finalized:
return {}, "", f"申万行业 {sector_code} 盘后正式数据待入库"
return {}, "", f"免费申万实时暂不可用:{message[:180]}"
return {}, "", f"数据中枢申万实时暂不可用:{message[:180]}"
if row:
return dict(row), str(row.get("source") or "eastmoney_sw"), ""
return dict(row), str(row.get("source") or "datahub"), ""
if finalized:
return {}, "", f"申万行业 {sector_code} 当日盘后正式数据尚未入库"
if aggregator and sector_name:
try:
row = aggregator.eastmoney_sector(sector_name)
except Exception as exc:
return {}, "", f"免费行业实时暂不可用:{str(exc)[:180]}"
if row:
return dict(row), str(row.get("source") or "eastmoney_sector"), ""
return {}, "", f"申万行业 {sector_code} 当日外显待补充"
def _load_member_realtime_quotes(
@@ -677,31 +669,6 @@ class ShenwanIndustryMixin:
delayed = any(item.get("delayed") for item in filtered)
return filtered, "datahub_delayed" if delayed else "datahub"
aggregator = getattr(self, "realtime_aggregator", None)
eastmoney_loader = getattr(aggregator, "eastmoney_stock_quotes", None) if aggregator else None
if callable(eastmoney_loader):
try:
filtered = consider(eastmoney_loader(wanted, expected_date=trade_date) or [], "eastmoney_ulist")
if len(filtered) >= max(1, int(len(wanted) * 0.9)):
return filtered, "eastmoney_ulist"
except Exception:
pass
tencent_loader = getattr(aggregator, "tencent_stock_quotes", None) if aggregator else None
if callable(tencent_loader):
try:
filtered = consider(tencent_loader(wanted, expected_date=trade_date) or [], "tencent_qt")
if len(filtered) >= max(1, int(len(wanted) * 0.9)):
return filtered, "tencent_qt"
except Exception:
pass
try:
quotes, source = self._free_realtime_quotes(trade_date, ",".join(wanted))
consider(quotes, source)
except TushareError:
pass
if best_rows:
delayed = any(item.get("delayed") for item in best_rows)
if delayed and not str(best_source).endswith("_delayed"):