Compare commits
1
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
d175bb65d4 |
@@ -7,9 +7,6 @@ TUSHARE_TOKEN=your_tushare_token_here
|
||||
|
||||
# Optional xiaobai-datahub client. All DATAHUB_READ_* / DATAHUB_SHADOW_* flags
|
||||
# default off in config/datahub.config.json, so the website keeps using Tushare.
|
||||
# Extended datasets (HEL-463): LIMIT_EVENTS POPULARITY DRAGON_TIGER SECTOR_DAILY
|
||||
# QUOTES INDEX_QUOTES INTRADAY — plus first-batch CALENDAR STOCKS DAILY INDEX_DAILY
|
||||
# VALUATION MONEYFLOW AUCTION STATUS.
|
||||
DATAHUB_BASE_URL=http://127.0.0.1:8766
|
||||
DATAHUB_TOKEN=
|
||||
|
||||
|
||||
@@ -21,34 +21,11 @@ from backend.data.providers.tushare_client import TushareClient
|
||||
|
||||
LOGGER = logging.getLogger("xiaobai.datahub")
|
||||
ShadowSink = Callable[[dict[str, Any]], None]
|
||||
|
||||
|
||||
def _usable_intraday_points(rows: list[Any]) -> list[dict[str, Any]]:
|
||||
points: list[dict[str, Any]] = []
|
||||
for row in rows:
|
||||
if not isinstance(row, dict):
|
||||
continue
|
||||
try:
|
||||
close = float(row.get("close") or 0)
|
||||
except (TypeError, ValueError):
|
||||
close = 0.0
|
||||
if close <= 0:
|
||||
continue
|
||||
point = dict(row)
|
||||
if "average" not in point and point.get("avg_price") is not None:
|
||||
point["average"] = point.get("avg_price")
|
||||
points.append(point)
|
||||
return points
|
||||
|
||||
|
||||
EMPTY_FAIL_DATASETS = {
|
||||
"stocks", "daily", "index_daily", "valuation", "moneyflow", "auction",
|
||||
"limit_events", "sector_daily",
|
||||
}
|
||||
EMPTY_FAIL_DATASETS = {"stocks", "daily", "index_daily", "valuation", "moneyflow", "auction"}
|
||||
|
||||
|
||||
def looks_like_heaven(module_name: str, filename: str = "") -> bool:
|
||||
"""问天调用栈识别(诊断用)。问天按数据集依赖接入,不再整栈强制旧链路。"""
|
||||
"""问天调用栈识别。问天未永久冻结,只是本阶段仍走旧 Tushare 链路。"""
|
||||
path = filename.replace("\\", "/")
|
||||
return module_name.startswith("backend.features.heaven") or "/features/heaven/" in path
|
||||
|
||||
@@ -111,34 +88,6 @@ class DatahubBridge:
|
||||
self._log_failure("status", exc)
|
||||
return None
|
||||
|
||||
def try_intraday(self, code: str) -> dict[str, Any] | None:
|
||||
flags = self.settings.flags("intraday")
|
||||
if not flags.read:
|
||||
return None
|
||||
try:
|
||||
response = self.client.intraday_points(code=code)
|
||||
data = response.data
|
||||
if not isinstance(data, dict):
|
||||
raise DatahubError("EMPTY", "datahub intraday payload invalid")
|
||||
points = _usable_intraday_points(data.get("points") or [])
|
||||
if not points:
|
||||
raise DatahubError("EMPTY", "datahub intraday empty")
|
||||
if (response.meta or {}).get("stale"):
|
||||
raise DatahubError("STALE", "datahub intraday stale")
|
||||
return {
|
||||
"entity_type": str(data.get("entity_type") or "stock"),
|
||||
"identifier": str(data.get("identifier") or code),
|
||||
"name": str(data.get("name") or ""),
|
||||
"code": str(data.get("code") or code),
|
||||
"trade_date": str(data.get("trade_date") or points[-1].get("date") or ""),
|
||||
"previous_close": float(data.get("previous_close") or 0),
|
||||
"points": points,
|
||||
"source": "datahub",
|
||||
}
|
||||
except Exception as exc:
|
||||
self._log_failure("intraday", exc)
|
||||
return None
|
||||
|
||||
def query(
|
||||
self,
|
||||
api_name: str,
|
||||
@@ -147,8 +96,8 @@ class DatahubBridge:
|
||||
legacy_query: Callable[..., list[dict[str, Any]]],
|
||||
) -> list[dict[str, Any]]:
|
||||
dataset = API_TO_DATASET.get(api_name)
|
||||
# 问天按实际数据依赖接入:已映射到 hub 的 API 跟随开关;未映射的继续旧链路。
|
||||
if not dataset:
|
||||
# 问天允许后续纳入 datahub;首批只读接入仍保持旧链路,避免误切。
|
||||
if not dataset or self.heaven_guard():
|
||||
return legacy_query(api_name, params, fields)
|
||||
flags = self.settings.flags(dataset)
|
||||
if not flags.read and not flags.shadow:
|
||||
@@ -159,7 +108,7 @@ class DatahubBridge:
|
||||
hub_error: str | None = None
|
||||
hub_canonical: list[dict[str, Any]] = []
|
||||
try:
|
||||
response = self._fetch_dataset(dataset, params or {}, api_name=api_name)
|
||||
response = self._fetch_dataset(dataset, params or {})
|
||||
hub_canonical = self._extract_rows(dataset, response, params or {})
|
||||
hub_rows = to_native_rows(dataset, hub_canonical)
|
||||
hub_meta = dict(response.meta)
|
||||
@@ -187,7 +136,7 @@ class DatahubBridge:
|
||||
return project_fields(hub_rows, fields)
|
||||
return legacy_query(api_name, params, fields)
|
||||
|
||||
def _fetch_dataset(self, dataset: str, params: dict[str, Any], api_name: str = "") -> DatahubResponse:
|
||||
def _fetch_dataset(self, dataset: str, params: dict[str, Any]) -> DatahubResponse:
|
||||
date = yyyymmdd(params.get("trade_date") or params.get("date"))
|
||||
start = yyyymmdd(params.get("start_date") or params.get("from") or date)
|
||||
end = yyyymmdd(params.get("end_date") or params.get("to") or date)
|
||||
@@ -204,10 +153,6 @@ class DatahubBridge:
|
||||
"valuation": self.client.valuation,
|
||||
"moneyflow": self.client.moneyflow,
|
||||
"auction": self.client.auction,
|
||||
"limit_events": self.client.limit_events,
|
||||
"popularity": self.client.popularity,
|
||||
"dragon_tiger": self.client.dragon_tiger,
|
||||
"sector_daily": self.client.sectors,
|
||||
}
|
||||
fetcher = fetchers[dataset]
|
||||
query: dict[str, Any] = {}
|
||||
@@ -222,23 +167,6 @@ class DatahubBridge:
|
||||
query["to"] = end
|
||||
if dataset == "daily":
|
||||
query["adjust"] = "none"
|
||||
if dataset == "limit_events":
|
||||
limit_type = str(params.get("limit_type") or "").strip().upper()
|
||||
if limit_type:
|
||||
query["limit_type"] = limit_type
|
||||
if dataset == "popularity":
|
||||
if api_name == "ths_hot":
|
||||
query["source"] = "ths"
|
||||
elif api_name == "dc_hot":
|
||||
query["source"] = "dc"
|
||||
if dataset == "sector_daily":
|
||||
family = {
|
||||
"ths_daily": "ths",
|
||||
"dc_index": "dc",
|
||||
"sw_daily": "sw",
|
||||
}.get(api_name, "")
|
||||
if family:
|
||||
query["family"] = family
|
||||
return self._paginate(fetcher, query)
|
||||
|
||||
def _paginate(self, fetcher: Callable[..., DatahubResponse], params: dict[str, Any]) -> DatahubResponse:
|
||||
|
||||
@@ -60,27 +60,6 @@ class DatahubClient:
|
||||
def auction(self, **params: Any) -> DatahubResponse:
|
||||
return self.get("/v1/auction", params)
|
||||
|
||||
def limit_events(self, **params: Any) -> DatahubResponse:
|
||||
return self.get("/v1/limit-events", params)
|
||||
|
||||
def popularity(self, **params: Any) -> DatahubResponse:
|
||||
return self.get("/v1/popularity", params)
|
||||
|
||||
def dragon_tiger(self, **params: Any) -> DatahubResponse:
|
||||
return self.get("/v1/dragon-tiger", params)
|
||||
|
||||
def sectors(self, **params: Any) -> DatahubResponse:
|
||||
return self.get("/v1/sectors", params)
|
||||
|
||||
def quotes_latest(self, **params: Any) -> DatahubResponse:
|
||||
return self.get("/v1/quotes/latest", params)
|
||||
|
||||
def index_quotes(self, **params: Any) -> DatahubResponse:
|
||||
return self.get("/v1/indexes/quotes", params)
|
||||
|
||||
def intraday_points(self, **params: Any) -> DatahubResponse:
|
||||
return self.get("/v1/intraday/points", params)
|
||||
|
||||
def dataset_status(self, date: str) -> DatahubResponse:
|
||||
return self.get("/v1/datasets/status", {"date": date})
|
||||
|
||||
|
||||
@@ -17,13 +17,6 @@ API_TO_DATASET = {
|
||||
"index_daily": "index_daily",
|
||||
"moneyflow": "moneyflow",
|
||||
"stk_auction": "auction",
|
||||
"limit_list_d": "limit_events",
|
||||
"ths_hot": "popularity",
|
||||
"dc_hot": "popularity",
|
||||
"hm_detail": "dragon_tiger",
|
||||
"ths_daily": "sector_daily",
|
||||
"dc_index": "sector_daily",
|
||||
"sw_daily": "sector_daily",
|
||||
}
|
||||
|
||||
SCALE_FIELDS = {
|
||||
@@ -42,16 +35,6 @@ SCALE_FIELDS = {
|
||||
"net_mf_amount": AMOUNT_WAN_YUAN,
|
||||
},
|
||||
"auction": {"vol": VOLUME_LOT, "float_share": AMOUNT_WAN_YUAN},
|
||||
"limit_events": {
|
||||
"limit_amount": AMOUNT_WAN_YUAN,
|
||||
"float_mv": AMOUNT_WAN_YUAN,
|
||||
"total_mv": AMOUNT_WAN_YUAN,
|
||||
},
|
||||
"dragon_tiger": {
|
||||
"buy_amount": AMOUNT_WAN_YUAN,
|
||||
"sell_amount": AMOUNT_WAN_YUAN,
|
||||
"net_amount": AMOUNT_WAN_YUAN,
|
||||
},
|
||||
}
|
||||
|
||||
|
||||
@@ -84,16 +67,6 @@ def to_native_row(dataset: str, row: dict[str, Any]) -> dict[str, Any]:
|
||||
converted[field] = _unscale(converted.get(field), factor)
|
||||
if dataset == "stocks":
|
||||
converted.pop("updated_at", None)
|
||||
if dataset == "popularity":
|
||||
# keep hub source; callers filter ths/dc themselves when needed
|
||||
if converted.get("ts_name") and not converted.get("name"):
|
||||
converted["name"] = converted.get("ts_name")
|
||||
if dataset == "dragon_tiger":
|
||||
if converted.get("ts_name") and not converted.get("name"):
|
||||
converted["name"] = converted.get("ts_name")
|
||||
if dataset == "sector_daily":
|
||||
if converted.get("pct_change") is not None and converted.get("pct_chg") is None:
|
||||
converted["pct_chg"] = converted.get("pct_change")
|
||||
return converted
|
||||
|
||||
|
||||
@@ -123,30 +96,6 @@ def row_key(dataset: str, row: dict[str, Any]) -> tuple[str, ...]:
|
||||
return (str(row.get("ts_code") or "").upper(),)
|
||||
if dataset == "status":
|
||||
return (str(row.get("dataset") or ""), yyyymmdd(row.get("trade_date")))
|
||||
if dataset == "limit_events":
|
||||
return (
|
||||
str(row.get("ts_code") or "").upper(),
|
||||
yyyymmdd(row.get("trade_date")),
|
||||
str(row.get("limit_type") or ""),
|
||||
)
|
||||
if dataset == "popularity":
|
||||
return (
|
||||
str(row.get("ts_code") or "").upper(),
|
||||
yyyymmdd(row.get("trade_date")),
|
||||
str(row.get("source") or ""),
|
||||
)
|
||||
if dataset == "dragon_tiger":
|
||||
return (
|
||||
str(row.get("ts_code") or "").upper(),
|
||||
yyyymmdd(row.get("trade_date")),
|
||||
str(row.get("hm_name") or ""),
|
||||
)
|
||||
if dataset == "sector_daily":
|
||||
return (
|
||||
str(row.get("ts_code") or "").upper(),
|
||||
yyyymmdd(row.get("trade_date")),
|
||||
str(row.get("family") or ""),
|
||||
)
|
||||
return (str(row.get("ts_code") or "").upper(), yyyymmdd(row.get("trade_date")))
|
||||
|
||||
|
||||
|
||||
@@ -17,13 +17,6 @@ DATASETS = (
|
||||
"valuation",
|
||||
"moneyflow",
|
||||
"auction",
|
||||
"limit_events",
|
||||
"popularity",
|
||||
"dragon_tiger",
|
||||
"sector_daily",
|
||||
"quotes",
|
||||
"index_quotes",
|
||||
"intraday",
|
||||
"status",
|
||||
)
|
||||
|
||||
@@ -35,13 +28,6 @@ ENV_DATASET = {
|
||||
"valuation": "VALUATION",
|
||||
"moneyflow": "MONEYFLOW",
|
||||
"auction": "AUCTION",
|
||||
"limit_events": "LIMIT_EVENTS",
|
||||
"popularity": "POPULARITY",
|
||||
"dragon_tiger": "DRAGON_TIGER",
|
||||
"sector_daily": "SECTOR_DAILY",
|
||||
"quotes": "QUOTES",
|
||||
"index_quotes": "INDEX_QUOTES",
|
||||
"intraday": "INTRADAY",
|
||||
"status": "STATUS",
|
||||
}
|
||||
|
||||
|
||||
@@ -37,9 +37,7 @@ class DataGateway:
|
||||
) -> TushareClient:
|
||||
if dataset_id:
|
||||
self.policy.assert_allowed(dataset_id, "tushare", usage)
|
||||
legacy = self.tushare_provider.client()
|
||||
legacy.realtime_aggregator = self.realtime_observer
|
||||
return DatahubAwareTushareClient(legacy, self.datahub)
|
||||
return DatahubAwareTushareClient(self.tushare_provider.client(), self.datahub)
|
||||
|
||||
def dataset_status(self, trade_date: str) -> list[dict[str, Any]] | None:
|
||||
return self.datahub.dataset_status(trade_date)
|
||||
@@ -87,13 +85,12 @@ def build_data_gateway(
|
||||
policy = DataSourcePolicy.load()
|
||||
settings = datahub_settings or DatahubSettings.load(credentials=credentials)
|
||||
datahub_client = DatahubClient(settings)
|
||||
datahub = DatahubBridge(settings, datahub_client)
|
||||
return DataGateway(
|
||||
policy=policy,
|
||||
quality=DataQualityGate.load(policy),
|
||||
tushare_provider=TushareProvider(token_supplier),
|
||||
ifind_provider=IfindProvider(ifind),
|
||||
chart_data=MarketChartClient(ifind, EastmoneyChartClient(), datahub),
|
||||
chart_data=MarketChartClient(ifind, EastmoneyChartClient()),
|
||||
realtime_observer=WebRealtimeAggregator(),
|
||||
datahub=datahub,
|
||||
datahub=DatahubBridge(settings, datahub_client),
|
||||
)
|
||||
|
||||
@@ -3,11 +3,7 @@ from __future__ import annotations
|
||||
from typing import Any
|
||||
|
||||
from backend.data.numbers import finite_number as _number
|
||||
from backend.data.providers.tushare_helpers import (
|
||||
_display_time,
|
||||
_prices_equal,
|
||||
calendar_is_open,
|
||||
)
|
||||
from backend.data.providers.tushare_helpers import _display_time, _prices_equal
|
||||
|
||||
|
||||
class DailyMarketMixin:
|
||||
@@ -21,11 +17,7 @@ class DailyMarketMixin:
|
||||
trade_date = requested
|
||||
else:
|
||||
row = requested_rows[0]
|
||||
trade_date = (
|
||||
row["cal_date"]
|
||||
if calendar_is_open(row.get("is_open"))
|
||||
else row.get("pretrade_date", requested)
|
||||
)
|
||||
trade_date = row["cal_date"] if row.get("is_open") == 1 else row.get("pretrade_date", requested)
|
||||
|
||||
resolved_rows = self.query(
|
||||
"trade_cal",
|
||||
|
||||
@@ -16,12 +16,6 @@ from backend.data.providers.tushare_transport import TushareError
|
||||
|
||||
|
||||
class DashboardMixin:
|
||||
def _now(self) -> datetime:
|
||||
clock = getattr(self, "clock", None)
|
||||
if callable(clock):
|
||||
return clock()
|
||||
return datetime.now().astimezone()
|
||||
|
||||
def dashboard(self, requested_date: str) -> dict[str, Any]:
|
||||
trade_date, previous_trade_date = self.resolve_trade_context(requested_date)
|
||||
if self.should_use_realtime(requested_date, trade_date):
|
||||
@@ -32,12 +26,11 @@ class DashboardMixin:
|
||||
)
|
||||
|
||||
daily = self._load_daily(trade_date)
|
||||
now = self._now()
|
||||
if (
|
||||
not daily
|
||||
and requested_date == now.strftime("%Y%m%d")
|
||||
and requested_date == datetime.now().astimezone().strftime("%Y%m%d")
|
||||
and trade_date == requested_date
|
||||
and now.time().replace(tzinfo=None) >= dt_time(9, 15)
|
||||
and datetime.now().astimezone().time().replace(tzinfo=None) >= dt_time(9, 15)
|
||||
):
|
||||
return self._realtime_dashboard(
|
||||
requested_date,
|
||||
@@ -105,14 +98,15 @@ class DashboardMixin:
|
||||
}
|
||||
return apply_sentiment_to_dashboard(dashboard)
|
||||
|
||||
def should_use_realtime(self, requested_date: str, trade_date: str) -> bool:
|
||||
"""Use live quotes for today's open session until official daily settles."""
|
||||
now = self._now()
|
||||
@staticmethod
|
||||
def should_use_realtime(requested_date: str, trade_date: str) -> bool:
|
||||
"""Use rt_k for today's open market until end-of-day datasets settle."""
|
||||
now = datetime.now().astimezone()
|
||||
today = now.strftime("%Y%m%d")
|
||||
return (
|
||||
requested_date == today
|
||||
and trade_date == today
|
||||
and dt_time(9, 15) <= now.time().replace(tzinfo=None) < dt_time(15, 5)
|
||||
and dt_time(9, 15) <= now.time().replace(tzinfo=None) < dt_time(16, 30)
|
||||
)
|
||||
|
||||
def _realtime_dashboard(
|
||||
@@ -128,7 +122,7 @@ class DashboardMixin:
|
||||
)
|
||||
if not codes:
|
||||
raise TushareError("No active stock codes available for rt_k")
|
||||
quotes, quote_source = self._load_realtime_quotes(codes, trade_date)
|
||||
quotes = self.query("rt_k", {"ts_code": codes})
|
||||
if not quotes:
|
||||
raise TushareError(f"No realtime data returned for {trade_date}")
|
||||
|
||||
@@ -184,30 +178,14 @@ class DashboardMixin:
|
||||
)
|
||||
sectors = _build_sectors(limits)
|
||||
previous_sectors = _build_sectors(previous_limits)
|
||||
now = self._now()
|
||||
now = datetime.now().astimezone()
|
||||
market_status = _realtime_market_status(now.time().replace(tzinfo=None))
|
||||
if quote_source == "eastmoney_clist":
|
||||
notice = (
|
||||
"盘中行情由东财免费实时快照计算;涨停原因、封板时间和开板次数以盘后榜单校正为准。"
|
||||
)
|
||||
source_name = "eastmoney"
|
||||
elif quote_source == "tencent_qt":
|
||||
notice = (
|
||||
"盘中行情由腾讯免费实时行情计算;涨停原因、封板时间和开板次数以盘后榜单校正为准。"
|
||||
)
|
||||
source_name = "tencent"
|
||||
else:
|
||||
notice = (
|
||||
"盘中行情由 Tushare rt_k 实时计算;涨停原因、封板时间和开板次数以盘后榜单校正为准。"
|
||||
)
|
||||
source_name = "tushare"
|
||||
dashboard = {
|
||||
"meta": {
|
||||
"requested_date": _display_date(requested_date),
|
||||
"trade_date": _display_date(trade_date),
|
||||
"previous_trade_date": _display_date(previous_trade_date),
|
||||
"source": source_name,
|
||||
"quote_source": quote_source,
|
||||
"source": "tushare",
|
||||
"mode": "realtime",
|
||||
"realtime": True,
|
||||
"market_status": market_status,
|
||||
@@ -215,8 +193,7 @@ class DashboardMixin:
|
||||
"auto_refresh": False,
|
||||
"quote_count": len(daily),
|
||||
"updated_at": now.isoformat(timespec="seconds"),
|
||||
"notice": notice,
|
||||
"indices": self._free_realtime_indices() if quote_source != "tushare_rt_k" else [],
|
||||
"notice": "盘中行情由 Tushare rt_k 实时计算;涨停原因、封板时间和开板次数以盘后榜单校正为准。",
|
||||
},
|
||||
"overview": _build_overview(daily, up_rows, down_rows, broken_rows),
|
||||
"limits": limits,
|
||||
@@ -230,67 +207,6 @@ class DashboardMixin:
|
||||
}
|
||||
return apply_sentiment_to_dashboard(dashboard)
|
||||
|
||||
def _realtime_aggregator(self):
|
||||
aggregator = getattr(self, "realtime_aggregator", None)
|
||||
if aggregator is None:
|
||||
raise TushareError("免费实时源未配置")
|
||||
return aggregator
|
||||
|
||||
def _load_realtime_quotes(
|
||||
self,
|
||||
codes: str,
|
||||
trade_date: str,
|
||||
) -> tuple[list[dict[str, Any]], str]:
|
||||
rt_error = ""
|
||||
try:
|
||||
quotes = self.query("rt_k", {"ts_code": codes})
|
||||
if quotes:
|
||||
return list(quotes), "tushare_rt_k"
|
||||
rt_error = f"No realtime data returned for {trade_date}"
|
||||
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"
|
||||
)
|
||||
return quotes, quote_source
|
||||
|
||||
def _free_realtime_quotes(
|
||||
self,
|
||||
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:
|
||||
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"
|
||||
|
||||
def _free_realtime_indices(self) -> list[dict[str, Any]]:
|
||||
try:
|
||||
return self._realtime_aggregator().eastmoney_indices()
|
||||
except Exception:
|
||||
return []
|
||||
|
||||
def _load_realtime_reference(
|
||||
self,
|
||||
trade_date: str,
|
||||
@@ -318,7 +234,7 @@ class DashboardMixin:
|
||||
{"trade_date": previous_trade_date},
|
||||
"ts_code,trade_date,total_share,float_share,free_share,total_mv,circ_mv",
|
||||
)
|
||||
if not basic_rows:
|
||||
if not basic_rows or not price_limits:
|
||||
raise TushareError(f"Realtime reference data is incomplete for {trade_date}")
|
||||
result = {
|
||||
"basic_rows": basic_rows,
|
||||
|
||||
@@ -6,17 +6,6 @@ from typing import Any
|
||||
from backend.data.numbers import finite_number as _number
|
||||
|
||||
|
||||
def calendar_is_open(value: Any) -> bool:
|
||||
if value in (True, 1, "1", "Y", "y"):
|
||||
return True
|
||||
if value in (False, 0, "0", "N", "n", None, ""):
|
||||
return False
|
||||
try:
|
||||
return int(value) == 1
|
||||
except (TypeError, ValueError):
|
||||
return False
|
||||
|
||||
|
||||
def _text(value: Any) -> str:
|
||||
if isinstance(value, (list, tuple, set)):
|
||||
return "、".join(str(item).strip() for item in value if str(item).strip())
|
||||
|
||||
@@ -59,12 +59,6 @@ class IndexMixin:
|
||||
}
|
||||
|
||||
def realtime_market_indices(self, requested_date: str) -> dict[str, Any]:
|
||||
try:
|
||||
return self._tushare_realtime_market_indices(requested_date)
|
||||
except TushareError:
|
||||
return self._free_realtime_market_indices(requested_date)
|
||||
|
||||
def _tushare_realtime_market_indices(self, requested_date: str) -> dict[str, Any]:
|
||||
trade_date, _ = self.resolve_trade_context(requested_date)
|
||||
index_names = {
|
||||
"000001.SH": "上证指数",
|
||||
@@ -122,52 +116,3 @@ class IndexMixin:
|
||||
"average_return_20d": 0,
|
||||
},
|
||||
}
|
||||
|
||||
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,
|
||||
},
|
||||
}
|
||||
|
||||
@@ -19,20 +19,8 @@ class RealtimeAggregateError(RuntimeError):
|
||||
|
||||
|
||||
EASTMONEY_INDEX_URL = "https://push2.eastmoney.com/api/qt/ulist.np/get"
|
||||
EASTMONEY_STOCK_URL = "https://push2.eastmoney.com/api/qt/stock/get"
|
||||
EASTMONEY_STOCK_FIELDS = "f43,f44,f45,f46,f47,f48,f57,f58,f60,f86,f168"
|
||||
EASTMONEY_SECTOR_URL = "https://push2.eastmoney.com/api/qt/clist/get"
|
||||
EASTMONEY_A_SHARE_BOARDS = (
|
||||
"m:0+t:6",
|
||||
"m:0+t:80",
|
||||
"m:1+t:2",
|
||||
"m:1+t:23",
|
||||
"m:0+t:81",
|
||||
)
|
||||
EASTMONEY_QUOTE_FIELDS = "f12,f13,f14,f2,f3,f4,f5,f6,f15,f16,f17,f18,f8,f124"
|
||||
EASTMONEY_MARKET_PAGE_SIZE = 100
|
||||
TENCENT_INDEX_URL = "https://qt.gtimg.cn/q=sh000001,sz399001,sz399006"
|
||||
TENCENT_QUOTE_URL = "https://qt.gtimg.cn/q="
|
||||
THS_LIMIT_URL = "https://data.10jqka.com.cn/dataapi/limit_up/limit_up_pool"
|
||||
XGB_POOL_URL = "https://flash-api.xuangubao.cn/api/pool/detail"
|
||||
BROWSER_USER_AGENT = (
|
||||
@@ -146,181 +134,6 @@ class WebRealtimeAggregator:
|
||||
raise RealtimeAggregateError(f"Eastmoney returned {len(result)}/3 indices")
|
||||
return result
|
||||
|
||||
def eastmoney_market_quotes(self, expected_date: str = "") -> list[dict[str, Any]]:
|
||||
"""Full A-share snapshot via Eastmoney clist, used when Tushare rt_k is unavailable."""
|
||||
now = time.time()
|
||||
cache_key = "assembled:eastmoney_market"
|
||||
with self._response_cache_lock:
|
||||
cached = self._response_cache.get(cache_key)
|
||||
cache_age = now - float((cached or {}).get("created_at") or 0)
|
||||
if cached and cache_age <= min(20, self.response_cache_ttl_seconds):
|
||||
quotes = list(cached.get("payload") or [])
|
||||
return self._filter_quotes_by_date(quotes, expected_date)
|
||||
|
||||
rows: list[dict[str, Any]] = []
|
||||
board_errors: list[str] = []
|
||||
for board in EASTMONEY_A_SHARE_BOARDS:
|
||||
try:
|
||||
rows.extend(self._eastmoney_board_quotes(board))
|
||||
except Exception as exc:
|
||||
board_errors.append(f"{board}:{exc}")
|
||||
quotes = []
|
||||
seen: set[str] = set()
|
||||
for row in rows:
|
||||
quote = _normalize_eastmoney_quote(row)
|
||||
ts_code = str((quote or {}).get("ts_code") or "")
|
||||
if not quote or ts_code in seen:
|
||||
continue
|
||||
seen.add(ts_code)
|
||||
quotes.append(quote)
|
||||
if len(quotes) < 200:
|
||||
detail = f";{'; '.join(board_errors)}" if board_errors else ""
|
||||
raise RealtimeAggregateError(
|
||||
f"Eastmoney market snapshot too small: {len(quotes)}{detail}"
|
||||
)
|
||||
quotes = self._filter_quotes_by_date(quotes, expected_date)
|
||||
with self._response_cache_lock:
|
||||
self._response_cache[cache_key] = {"created_at": now, "payload": quotes}
|
||||
return quotes
|
||||
|
||||
def _eastmoney_board_quotes(self, board: str) -> list[dict[str, Any]]:
|
||||
first = self._eastmoney_market_page(board, 1)
|
||||
data = first.get("data") or {}
|
||||
rows = _diff_rows(data)
|
||||
total = int(_number(data.get("total")))
|
||||
page_count = 1
|
||||
if total > 0:
|
||||
page_count = max(1, (total + EASTMONEY_MARKET_PAGE_SIZE - 1) // EASTMONEY_MARKET_PAGE_SIZE)
|
||||
for page in range(2, min(page_count, 40) + 1):
|
||||
payload = self._eastmoney_market_page(board, page)
|
||||
rows.extend(_diff_rows(payload.get("data") or {}))
|
||||
return rows
|
||||
|
||||
def _eastmoney_market_page(self, board: str, page: int) -> dict[str, Any]:
|
||||
return self._get_json(
|
||||
EASTMONEY_SECTOR_URL,
|
||||
{
|
||||
"pn": str(page),
|
||||
"pz": str(EASTMONEY_MARKET_PAGE_SIZE),
|
||||
"po": "1",
|
||||
"np": "1",
|
||||
"fltt": "2",
|
||||
"invt": "2",
|
||||
"fid": "f12",
|
||||
"fs": board,
|
||||
"fields": EASTMONEY_QUOTE_FIELDS,
|
||||
},
|
||||
referer="https://quote.eastmoney.com/center/gridlist.html",
|
||||
)
|
||||
|
||||
def _filter_quotes_by_date(
|
||||
self,
|
||||
quotes: list[dict[str, Any]],
|
||||
expected_date: str,
|
||||
) -> list[dict[str, Any]]:
|
||||
want = str(expected_date or "").replace("-", "")
|
||||
if not want or not quotes:
|
||||
return quotes
|
||||
dated = [item for item in quotes if str(item.get("quote_date") or "") == want]
|
||||
if dated and len(dated) >= max(100, int(len(quotes) * 0.2)):
|
||||
return dated
|
||||
if dated:
|
||||
return dated
|
||||
if all(not item.get("quote_date") for item in quotes):
|
||||
return quotes
|
||||
raise RealtimeAggregateError(f"Eastmoney quotes are not for {want}")
|
||||
|
||||
def tencent_market_quotes(
|
||||
self,
|
||||
codes: list[str],
|
||||
expected_date: str = "",
|
||||
) -> list[dict[str, Any]]:
|
||||
symbols: list[str] = []
|
||||
seen: set[str] = set()
|
||||
for raw in codes:
|
||||
ts = str(raw or "").strip().upper()
|
||||
if not ts:
|
||||
continue
|
||||
symbol = ts.split(".")[0]
|
||||
if not symbol.isdigit() or len(symbol) != 6 or symbol in seen:
|
||||
continue
|
||||
seen.add(symbol)
|
||||
if ts.endswith(".SH") or symbol.startswith(("5", "6", "9")):
|
||||
symbols.append(f"sh{symbol}")
|
||||
elif ts.endswith(".BJ") or symbol.startswith(("4", "8")):
|
||||
symbols.append(f"bj{symbol}")
|
||||
else:
|
||||
symbols.append(f"sz{symbol}")
|
||||
if not symbols:
|
||||
raise RealtimeAggregateError("No stock codes available for Tencent quotes")
|
||||
|
||||
quotes: list[dict[str, Any]] = []
|
||||
batch_size = 80
|
||||
|
||||
def load_batch(batch: list[str]) -> list[dict[str, Any]]:
|
||||
raw, _cache_age = self._get_text(
|
||||
f"{TENCENT_QUOTE_URL}{','.join(batch)}",
|
||||
referer="https://gu.qq.com/",
|
||||
encoding="gb18030",
|
||||
)
|
||||
return [
|
||||
quote
|
||||
for line in raw.splitlines()
|
||||
if (quote := _parse_tencent_stock_quote(line))
|
||||
]
|
||||
|
||||
batches = [symbols[index:index + batch_size] for index in range(0, len(symbols), batch_size)]
|
||||
errors: list[str] = []
|
||||
with ThreadPoolExecutor(max_workers=4) as executor:
|
||||
for result in executor.map(self._capture, [lambda batch=batch: load_batch(batch) for batch in batches]):
|
||||
rows, status = result
|
||||
if status.get("ok") and rows:
|
||||
quotes.extend(rows)
|
||||
elif not status.get("ok"):
|
||||
errors.append(str(status.get("error") or "batch failed"))
|
||||
if len(quotes) < 200:
|
||||
detail = f";{'; '.join(errors[:3])}" if errors else ""
|
||||
raise RealtimeAggregateError(
|
||||
f"Tencent market snapshot too small: {len(quotes)}{detail}"
|
||||
)
|
||||
return self._filter_quotes_by_date(quotes, expected_date)
|
||||
|
||||
def tencent_stock_quote(self, code: str, expected_date: str = "") -> dict[str, Any]:
|
||||
symbol, _secid, ts_code = _a_share_identity(code)
|
||||
raw, _cache_age = self._get_text(
|
||||
f"{TENCENT_QUOTE_URL}{symbol}",
|
||||
referer="https://gu.qq.com/",
|
||||
encoding="gb18030",
|
||||
)
|
||||
quote = next(
|
||||
(
|
||||
item
|
||||
for line in raw.splitlines()
|
||||
if (item := _parse_tencent_stock_quote(line))
|
||||
),
|
||||
None,
|
||||
)
|
||||
if not quote:
|
||||
raise RealtimeAggregateError(f"Tencent stock quote unavailable for {ts_code}")
|
||||
return _require_quote_date(quote, expected_date)
|
||||
|
||||
def eastmoney_stock_quote(self, code: str, expected_date: str = "") -> dict[str, Any]:
|
||||
_symbol, secid, ts_code = _a_share_identity(code)
|
||||
payload = self._get_json(
|
||||
EASTMONEY_STOCK_URL,
|
||||
{
|
||||
"secid": secid,
|
||||
"invt": "2",
|
||||
"fltt": "2",
|
||||
"fields": EASTMONEY_STOCK_FIELDS,
|
||||
},
|
||||
referer="https://quote.eastmoney.com/",
|
||||
)
|
||||
quote = _normalize_eastmoney_stock_quote(payload.get("data") or {}, ts_code)
|
||||
if not quote:
|
||||
raise RealtimeAggregateError(f"Eastmoney stock quote unavailable for {ts_code}")
|
||||
return _require_quote_date(quote, expected_date)
|
||||
|
||||
def tencent_indices(self) -> list[dict[str, Any]]:
|
||||
raw, cache_age = self._get_text(
|
||||
TENCENT_INDEX_URL,
|
||||
@@ -584,143 +397,6 @@ class WebRealtimeAggregator:
|
||||
) from last_error
|
||||
|
||||
|
||||
def _diff_rows(data: dict[str, Any]) -> list[dict[str, Any]]:
|
||||
diff = data.get("diff") or []
|
||||
if isinstance(diff, dict):
|
||||
return [row for row in diff.values() if isinstance(row, dict)]
|
||||
return [row for row in diff if isinstance(row, dict)]
|
||||
|
||||
|
||||
def _a_share_identity(code: str) -> tuple[str, str, str]:
|
||||
raw = str(code or "").strip().upper()
|
||||
symbol = raw.split(".")[0]
|
||||
if not symbol.isdigit() or len(symbol) != 6:
|
||||
raise RealtimeAggregateError("Invalid stock code")
|
||||
if raw.endswith(".SH") or symbol.startswith(("5", "6", "9")):
|
||||
return f"sh{symbol}", f"1.{symbol}", f"{symbol}.SH"
|
||||
if raw.endswith(".BJ") or symbol.startswith(("4", "8")):
|
||||
return f"bj{symbol}", f"0.{symbol}", f"{symbol}.BJ"
|
||||
return f"sz{symbol}", f"0.{symbol}", f"{symbol}.SZ"
|
||||
|
||||
|
||||
def _require_quote_date(quote: dict[str, Any], expected_date: str) -> dict[str, Any]:
|
||||
want = str(expected_date or "").replace("-", "")
|
||||
got = str(quote.get("quote_date") or "")
|
||||
if want and got != want:
|
||||
raise RealtimeAggregateError(f"quote date {got or 'empty'} is not {want}")
|
||||
return quote
|
||||
|
||||
|
||||
def _normalize_eastmoney_stock_quote(
|
||||
row: dict[str, Any], ts_code: str
|
||||
) -> dict[str, Any] | None:
|
||||
close = _number(row.get("f43"))
|
||||
previous_close = _number(row.get("f60"))
|
||||
if close <= 0 or previous_close <= 0:
|
||||
return None
|
||||
epoch = int(_number(row.get("f86")))
|
||||
quote_date = ""
|
||||
if epoch > 0:
|
||||
quote_date = datetime.fromtimestamp(epoch).astimezone().strftime("%Y%m%d")
|
||||
return {
|
||||
"ts_code": ts_code,
|
||||
"name": row.get("f58") or ts_code.split(".")[0],
|
||||
"pre_close": previous_close,
|
||||
"open": _number(row.get("f46")),
|
||||
"high": _number(row.get("f44")),
|
||||
"low": _number(row.get("f45")),
|
||||
"close": close,
|
||||
"vol": _number(row.get("f47")) * 100,
|
||||
"amount": _number(row.get("f48")),
|
||||
"num": 0,
|
||||
"quote_date": quote_date,
|
||||
"quote_time_epoch": epoch,
|
||||
"turnover_rate": _number(row.get("f168")),
|
||||
"source": "eastmoney_stock",
|
||||
}
|
||||
|
||||
|
||||
def _parse_tencent_stock_quote(line: str) -> dict[str, Any] | None:
|
||||
if '="' not in line:
|
||||
return None
|
||||
prefix, payload = line.split('="', 1)
|
||||
fields = payload.rsplit('";', 1)[0].split("~")
|
||||
if len(fields) < 38:
|
||||
return None
|
||||
symbol = fields[2]
|
||||
if not symbol.isdigit() or len(symbol) != 6:
|
||||
return None
|
||||
close = _number(fields[3])
|
||||
previous_close = _number(fields[4])
|
||||
if close <= 0 or previous_close <= 0:
|
||||
return None
|
||||
marker = prefix.lower()
|
||||
if "sh" in marker:
|
||||
ts_code = f"{symbol}.SH"
|
||||
elif "bj" in marker:
|
||||
ts_code = f"{symbol}.BJ"
|
||||
else:
|
||||
ts_code = f"{symbol}.SZ"
|
||||
try:
|
||||
quote_time = datetime.strptime(fields[30], "%Y%m%d%H%M%S")
|
||||
quote_date = quote_time.strftime("%Y%m%d")
|
||||
epoch = int(quote_time.timestamp())
|
||||
except ValueError:
|
||||
quote_date = ""
|
||||
epoch = 0
|
||||
return {
|
||||
"ts_code": ts_code,
|
||||
"name": fields[1] or symbol,
|
||||
"pre_close": previous_close,
|
||||
"open": _number(fields[5]),
|
||||
"high": _number(fields[33]),
|
||||
"low": _number(fields[34]),
|
||||
"close": close,
|
||||
"vol": _number(fields[6]) * 100,
|
||||
"amount": _number(fields[37]) * 10000,
|
||||
"num": 0,
|
||||
"quote_date": quote_date,
|
||||
"quote_time_epoch": epoch,
|
||||
"source": "tencent_qt",
|
||||
}
|
||||
|
||||
|
||||
def _normalize_eastmoney_quote(row: dict[str, Any]) -> dict[str, Any] | None:
|
||||
symbol = str(row.get("f12") or "").strip()
|
||||
if not symbol.isdigit() or len(symbol) != 6:
|
||||
return None
|
||||
close = _number(row.get("f2"))
|
||||
previous_close = _number(row.get("f18"))
|
||||
if close <= 0 or previous_close <= 0:
|
||||
return None
|
||||
market = int(_number(row.get("f13")))
|
||||
if market == 1 or symbol.startswith(("5", "6", "9")):
|
||||
ts_code = f"{symbol}.SH"
|
||||
elif symbol.startswith(("4", "8")):
|
||||
ts_code = f"{symbol}.BJ"
|
||||
else:
|
||||
ts_code = f"{symbol}.SZ"
|
||||
epoch = int(_number(row.get("f124")))
|
||||
quote_date = ""
|
||||
if epoch > 0:
|
||||
quote_date = datetime.fromtimestamp(epoch).astimezone().strftime("%Y%m%d")
|
||||
return {
|
||||
"ts_code": ts_code,
|
||||
"name": row.get("f14") or symbol,
|
||||
"pre_close": previous_close,
|
||||
"open": _number(row.get("f17")),
|
||||
"high": _number(row.get("f15")),
|
||||
"low": _number(row.get("f16")),
|
||||
"close": close,
|
||||
"vol": _number(row.get("f5")) * 100,
|
||||
"amount": _number(row.get("f6")),
|
||||
"num": 0,
|
||||
"quote_date": quote_date,
|
||||
"quote_time_epoch": epoch,
|
||||
"source": "eastmoney_clist",
|
||||
}
|
||||
|
||||
|
||||
def _normalize_sector(value: Any) -> str:
|
||||
text = str(value or "").strip().replace(" ", "")
|
||||
for suffix in ("板块", "概念", "行业", "Ⅱ", "Ⅲ", "(A股)", "(A股)"):
|
||||
|
||||
@@ -2,7 +2,6 @@ from __future__ import annotations
|
||||
|
||||
import http.client
|
||||
import json
|
||||
import logging
|
||||
import re
|
||||
import time
|
||||
import urllib.error
|
||||
@@ -16,15 +15,12 @@ from typing import Any, ClassVar
|
||||
from backend.bootstrap.config import tushare_code as _stock_market_code
|
||||
from backend.data.providers.ifind_client import IfindError, IfindHttpClient
|
||||
|
||||
LOGGER = logging.getLogger("xiaobai.charts")
|
||||
|
||||
|
||||
class ChartDataError(RuntimeError):
|
||||
pass
|
||||
|
||||
|
||||
TRENDS_URL = "https://push2delay.eastmoney.com/api/qt/stock/trends2/get"
|
||||
HIS_TRENDS_URL = "https://push2his.eastmoney.com/api/qt/stock/trends2/get"
|
||||
BOARD_LIST_URL = "https://push2delay.eastmoney.com/api/qt/clist/get"
|
||||
BROWSER_USER_AGENT = (
|
||||
"Mozilla/5.0 (Windows NT 10.0; Win64; x64) "
|
||||
@@ -41,23 +37,14 @@ INDEX_SECIDS = {
|
||||
class MarketChartClient:
|
||||
"""Prefer iFinD for display charts and retain Eastmoney as a last resort."""
|
||||
|
||||
def __init__(
|
||||
self,
|
||||
ifind: IfindHttpClient,
|
||||
fallback: "EastmoneyChartClient",
|
||||
datahub: Any = None,
|
||||
) -> None:
|
||||
def __init__(self, ifind: IfindHttpClient, fallback: "EastmoneyChartClient") -> None:
|
||||
self.ifind = ifind
|
||||
self.fallback = fallback
|
||||
self.datahub = datahub
|
||||
|
||||
def stock_intraday(self, code: str) -> dict[str, Any]:
|
||||
normalized = str(code or "").strip()
|
||||
if not re.fullmatch(r"\d{6}", normalized):
|
||||
raise ChartDataError("Invalid stock code")
|
||||
hub_chart = self._datahub_intraday(normalized)
|
||||
if hub_chart is not None:
|
||||
return hub_chart
|
||||
ifind_code = _stock_market_code(normalized)
|
||||
try:
|
||||
return self._ifind_intraday(ifind_code, "stock", normalized)
|
||||
@@ -86,29 +73,11 @@ class MarketChartClient:
|
||||
normalized = str(identifier or "").strip().upper()
|
||||
if normalized not in INDEX_SECIDS:
|
||||
raise ChartDataError("Unsupported index")
|
||||
hub_chart = self._datahub_intraday(normalized)
|
||||
if hub_chart is not None:
|
||||
return hub_chart
|
||||
try:
|
||||
return self._ifind_intraday(normalized, "index", normalized)
|
||||
except (IfindError, ChartDataError):
|
||||
return self.fallback.index_intraday(normalized)
|
||||
|
||||
def _datahub_intraday(self, code: str) -> dict[str, Any] | None:
|
||||
if self.datahub is None:
|
||||
return None
|
||||
try:
|
||||
chart = self.datahub.try_intraday(code)
|
||||
except Exception as exc:
|
||||
LOGGER.warning("datahub intraday unexpected error: %s", exc)
|
||||
return None
|
||||
if not chart:
|
||||
return None
|
||||
points = list(chart.get("points") or [])
|
||||
if not points:
|
||||
return None
|
||||
return chart
|
||||
|
||||
def board_intraday(self, identifier: str, name: str = "") -> dict[str, Any]:
|
||||
normalized = str(identifier or "").strip().upper()
|
||||
try:
|
||||
@@ -336,29 +305,21 @@ class EastmoneyChartClient:
|
||||
if cached is not None:
|
||||
return cached
|
||||
|
||||
params = {
|
||||
"secid": secid,
|
||||
"fields1": "f1,f2,f3,f4,f5,f6,f7,f8,f9,f10,f11,f12,f13",
|
||||
"fields2": "f51,f52,f53,f54,f55,f56,f57,f58",
|
||||
"iscr": "0",
|
||||
}
|
||||
last_error: Exception | None = None
|
||||
data: dict[str, Any] = {}
|
||||
points: list[dict[str, Any]] = []
|
||||
for url, ndays in ((TRENDS_URL, "1"), (TRENDS_URL, "5"), (HIS_TRENDS_URL, "5")):
|
||||
request_params = {**params, "ndays": ndays}
|
||||
try:
|
||||
payload = self._request_json(url, request_params, "https://quote.eastmoney.com/")
|
||||
except ChartDataError as exc:
|
||||
last_error = exc
|
||||
continue
|
||||
data = payload.get("data") or {}
|
||||
parsed = [point for raw in data.get("trends") or [] if (point := _parse_trend(raw))]
|
||||
points = _latest_session(parsed)
|
||||
if points:
|
||||
break
|
||||
payload = self._request_json(
|
||||
TRENDS_URL,
|
||||
{
|
||||
"secid": secid,
|
||||
"fields1": "f1,f2,f3,f4,f5,f6,f7,f8,f9,f10,f11,f12,f13",
|
||||
"fields2": "f51,f52,f53,f54,f55,f56,f57,f58",
|
||||
"iscr": "0",
|
||||
"ndays": "1",
|
||||
},
|
||||
"https://quote.eastmoney.com/",
|
||||
)
|
||||
data = payload.get("data") or {}
|
||||
points = [point for raw in data.get("trends") or [] if (point := _parse_trend(raw))]
|
||||
if not points:
|
||||
raise ChartDataError("No intraday chart data returned") from last_error
|
||||
raise ChartDataError("No intraday chart data returned")
|
||||
|
||||
result = {
|
||||
"entity_type": entity_type,
|
||||
@@ -472,15 +433,6 @@ class EastmoneyChartClient:
|
||||
raise ChartDataError("Intraday chart request failed") from last_error
|
||||
|
||||
|
||||
def _latest_session(points: list[dict[str, Any]]) -> list[dict[str, Any]]:
|
||||
if not points:
|
||||
return []
|
||||
latest = max(str(point.get("date") or "") for point in points)
|
||||
if not latest:
|
||||
return points
|
||||
return [point for point in points if str(point.get("date") or "") == latest]
|
||||
|
||||
|
||||
def _parse_trend(raw: Any) -> dict[str, Any] | None:
|
||||
fields = str(raw or "").split(",")
|
||||
if len(fields) < 8 or " " not in fields[0]:
|
||||
|
||||
@@ -15,7 +15,6 @@ from backend.bootstrap.config import (
|
||||
)
|
||||
from backend.data.providers.ifind_client import IfindError
|
||||
from backend.data.providers.tushare_client import TushareClient, TushareError
|
||||
from backend.data.realtime import RealtimeAggregateError
|
||||
from backend.features.market.backfill_history import (
|
||||
DEFAULT_RECENT_TRADING_DAYS,
|
||||
MAX_RANGE_TRADING_DAYS,
|
||||
@@ -43,7 +42,6 @@ SEARCH_TYPE_LABELS = {
|
||||
"theme": "题材",
|
||||
"index": "指数",
|
||||
}
|
||||
TODAY_DAILY_UNAVAILABLE_NOTICE = "今日日K暂不可用,仍显示最近收盘K线。"
|
||||
THS_SEARCH_TYPES = {
|
||||
"I": ("sector", "行业板块"),
|
||||
"R": ("sector", "地域板块"),
|
||||
@@ -65,37 +63,11 @@ class MarketServiceMixin:
|
||||
if gateway is not None:
|
||||
return gateway.tushare()
|
||||
# Compatibility for isolated legacy unit-test service stubs.
|
||||
client = TushareClient(self.token)
|
||||
aggregator = getattr(self, "realtime_aggregator", None)
|
||||
if aggregator is not None:
|
||||
client.realtime_aggregator = aggregator
|
||||
return client
|
||||
|
||||
def _now(self) -> datetime:
|
||||
clock = getattr(self, "clock", None)
|
||||
if callable(clock):
|
||||
return clock()
|
||||
return datetime.now().astimezone()
|
||||
|
||||
def _is_requested_open_session(self, requested_date: str) -> bool:
|
||||
now = self._now()
|
||||
if requested_date != now.strftime("%Y%m%d"):
|
||||
return False
|
||||
if now.time().replace(tzinfo=None) < dt_time(9, 15):
|
||||
return False
|
||||
client = self._tushare_client() if self.configured else None
|
||||
resolve = getattr(client, "resolve_trade_context", None) if client else None
|
||||
if resolve is None:
|
||||
return now.weekday() < 5
|
||||
try:
|
||||
trade_date, _ = resolve(requested_date)
|
||||
except Exception:
|
||||
return now.weekday() < 5
|
||||
return str(trade_date or "") == requested_date
|
||||
return TushareClient(self.token)
|
||||
|
||||
def get_dashboard(self, trade_date: str, force: bool = False) -> dict[str, Any]:
|
||||
normalized_date = normalize_date(trade_date)
|
||||
now = self._now()
|
||||
now = datetime.now().astimezone()
|
||||
if (
|
||||
normalized_date == now.strftime("%Y%m%d")
|
||||
and now.time().replace(tzinfo=None) < datetime.strptime("09:15", "%H:%M").time()
|
||||
@@ -202,14 +174,14 @@ class MarketServiceMixin:
|
||||
def _should_retry_incomplete_snapshot(
|
||||
self, snapshot: dict[str, Any], requested_date: str
|
||||
) -> bool:
|
||||
if requested_date != self._now().strftime("%Y%m%d"):
|
||||
if requested_date != date.today().strftime("%Y%m%d"):
|
||||
return False
|
||||
meta = snapshot.get("meta") or {}
|
||||
actual = str(meta.get("trade_date") or "").replace("-", "")
|
||||
stale_carry = bool(meta.get("carried_forward") or actual != requested_date)
|
||||
if stale_carry and self._is_requested_open_session(requested_date):
|
||||
return True
|
||||
incomplete = meta.get("limit_data_source") == "derived" or stale_carry
|
||||
incomplete = (
|
||||
meta.get("limit_data_source") == "derived"
|
||||
or bool(meta.get("carried_forward"))
|
||||
or str(meta.get("trade_date") or "").replace("-", "") != requested_date
|
||||
)
|
||||
return incomplete and self._snapshot_age_seconds(meta) >= 60
|
||||
|
||||
def _annotate_data_status(self, dashboard: dict[str, Any]) -> dict[str, Any]:
|
||||
@@ -227,9 +199,6 @@ class MarketServiceMixin:
|
||||
else:
|
||||
meta["data_status"] = "preparing"
|
||||
meta["display_notice"] = self._preparing_display_notice(actual, requested)
|
||||
elif meta.get("realtime"):
|
||||
meta["data_status"] = "intraday"
|
||||
meta.setdefault("display_notice", "")
|
||||
else:
|
||||
meta["data_status"] = "official"
|
||||
meta.setdefault("display_notice", "")
|
||||
@@ -256,9 +225,9 @@ class MarketServiceMixin:
|
||||
normalized_date: str,
|
||||
snapshot: dict[str, Any],
|
||||
) -> bool:
|
||||
if not self.configured or normalized_date != self._now().strftime("%Y%m%d"):
|
||||
if not self.configured or normalized_date != date.today().strftime("%Y%m%d"):
|
||||
return False
|
||||
now = self._now()
|
||||
now = datetime.now().astimezone()
|
||||
local_time = now.time().replace(tzinfo=None)
|
||||
realtime_start = datetime.strptime("09:15", "%H:%M").time()
|
||||
morning_end = datetime.strptime("11:35", "%H:%M").time()
|
||||
@@ -295,10 +264,7 @@ class MarketServiceMixin:
|
||||
raise TushareError("公共行情尚未配置")
|
||||
dashboard = self._tushare_client().dashboard(normalized_date)
|
||||
meta = dashboard.setdefault("meta", {})
|
||||
quote_source = str(meta.get("quote_source") or "")
|
||||
meta["source"] = source
|
||||
if quote_source:
|
||||
meta["quote_source"] = quote_source
|
||||
meta["requested_date"] = self._display_compact_date(normalized_date)
|
||||
if meta.get("limit_data_source") == "derived":
|
||||
meta.setdefault(
|
||||
@@ -310,12 +276,6 @@ class MarketServiceMixin:
|
||||
actual_date = normalize_date(
|
||||
str(dashboard.get("meta", {}).get("trade_date") or normalized_date)
|
||||
)
|
||||
if actual_date != normalized_date and self._is_requested_open_session(
|
||||
normalized_date
|
||||
):
|
||||
raise TushareError(
|
||||
f"Intraday dashboard resolved {actual_date} instead of {normalized_date}"
|
||||
)
|
||||
self.database.save_snapshot(actual_date, source, dashboard)
|
||||
if actual_date != normalized_date:
|
||||
dashboard.setdefault("meta", {}).update(
|
||||
@@ -337,30 +297,6 @@ class MarketServiceMixin:
|
||||
)
|
||||
return self._apply_reason_overrides(self._with_storage(dashboard, cached=False))
|
||||
except TushareError as exc:
|
||||
if self._is_requested_open_session(normalized_date):
|
||||
existing = self.database.get_snapshot(normalized_date)
|
||||
existing_date = str(
|
||||
((existing or {}).get("meta") or {}).get("trade_date") or ""
|
||||
).replace("-", "")
|
||||
if existing and existing_date == normalized_date:
|
||||
kept = copy.deepcopy(existing)
|
||||
kept.setdefault("meta", {}).update(
|
||||
{
|
||||
"requested_date": self._display_compact_date(normalized_date),
|
||||
}
|
||||
)
|
||||
self.database.finish_sync(
|
||||
sync_id,
|
||||
"fallback",
|
||||
self._record_count(kept),
|
||||
str(exc),
|
||||
"tushare",
|
||||
)
|
||||
return self._apply_reason_overrides(
|
||||
self._with_storage(kept, cached=True)
|
||||
)
|
||||
self.database.finish_sync(sync_id, "failed", message=str(exc))
|
||||
raise ValueError("当天盘中行情暂时不可用,请稍后重试。") from exc
|
||||
fallback = self.database.get_latest_real_snapshot(normalized_date)
|
||||
if fallback:
|
||||
actual = str((fallback.get("meta") or {}).get("trade_date") or "")
|
||||
@@ -816,27 +752,26 @@ class MarketServiceMixin:
|
||||
"trade_date": f"{actual_date[:4]}-{actual_date[4:6]}-{actual_date[6:]}",
|
||||
}
|
||||
today = now.strftime("%Y%m%d")
|
||||
latest_bar = (result.get("prices") or [{}])[-1] if result.get("prices") else {}
|
||||
official_today = (
|
||||
actual_date == today and not bool(latest_bar.get("realtime"))
|
||||
)
|
||||
after_close = now.time().replace(tzinfo=None) >= dt_time(15, 0)
|
||||
should_merge = (
|
||||
requested_date == today
|
||||
and actual_date <= today
|
||||
and now.weekday() < 5
|
||||
and now.time().replace(tzinfo=None) >= dt_time(9, 30)
|
||||
and not (official_today and after_close)
|
||||
)
|
||||
if should_merge:
|
||||
quote = self._resolve_today_daily_quote(code, today, result)
|
||||
quote = self._ifind_realtime_stock_quote(code)
|
||||
if quote and self._valid_realtime_stock_quote(quote, today):
|
||||
self._merge_realtime_stock_detail(result, quote, requested_date)
|
||||
elif actual_date < today:
|
||||
result["meta"] = {
|
||||
**(result.get("meta") or {}),
|
||||
"notice": TODAY_DAILY_UNAVAILABLE_NOTICE,
|
||||
}
|
||||
elif self.configured and actual_date < today:
|
||||
client = self._tushare_client()
|
||||
try:
|
||||
resolved_date, _ = client.resolve_trade_context(requested_date)
|
||||
if resolved_date == today:
|
||||
quote = client.realtime_stock_quote(tushare_code(code), requested_date)
|
||||
if self._valid_realtime_stock_quote(quote, today):
|
||||
self._merge_realtime_stock_detail(result, quote, requested_date)
|
||||
except TushareError:
|
||||
pass
|
||||
return self._enrich_stock_detail(result)
|
||||
|
||||
@staticmethod
|
||||
@@ -951,134 +886,6 @@ class MarketServiceMixin:
|
||||
"quote_time": str(row.get("time") or ""),
|
||||
}
|
||||
|
||||
def _resolve_today_daily_quote(
|
||||
self, code: str, today: str, payload: dict[str, Any]
|
||||
) -> dict[str, Any] | None:
|
||||
quote = self._ifind_realtime_stock_quote(code)
|
||||
if quote and self._valid_realtime_stock_quote(quote, today):
|
||||
return quote
|
||||
if self.configured:
|
||||
try:
|
||||
client = self._tushare_client()
|
||||
resolve = getattr(client, "resolve_trade_context", None)
|
||||
resolved = today
|
||||
if callable(resolve):
|
||||
resolved, _ = resolve(today)
|
||||
if str(resolved or "") == today:
|
||||
quote = client.realtime_stock_quote(tushare_code(code), today)
|
||||
if self._valid_realtime_stock_quote(quote, today):
|
||||
return quote
|
||||
except TushareError:
|
||||
pass
|
||||
quote = self._free_realtime_stock_quote(code, today)
|
||||
if quote and self._valid_realtime_stock_quote(quote, today):
|
||||
return quote
|
||||
return self._intraday_realtime_stock_quote(code, today, payload)
|
||||
|
||||
def _free_realtime_stock_quote(self, code: str, today: str) -> dict[str, Any] | None:
|
||||
aggregator = getattr(self, "realtime_aggregator", None)
|
||||
if aggregator is None:
|
||||
return None
|
||||
ts_code = tushare_code(code)
|
||||
for loader in (
|
||||
getattr(aggregator, "tencent_stock_quote", None),
|
||||
getattr(aggregator, "eastmoney_stock_quote", None),
|
||||
):
|
||||
if not callable(loader):
|
||||
continue
|
||||
try:
|
||||
row = loader(ts_code, expected_date=today)
|
||||
except (RealtimeAggregateError, Exception):
|
||||
continue
|
||||
quote = self._quote_from_free_row(code, today, row)
|
||||
if quote:
|
||||
return quote
|
||||
return None
|
||||
|
||||
def _quote_from_free_row(
|
||||
self, code: str, today: str, row: dict[str, Any]
|
||||
) -> dict[str, Any] | None:
|
||||
price = float(row.get("close") or 0)
|
||||
previous_close = float(row.get("pre_close") or 0)
|
||||
if price <= 0 or previous_close <= 0:
|
||||
return None
|
||||
try:
|
||||
name, sector = self._stock_identity(code, today)
|
||||
except Exception:
|
||||
name, sector = "--", "其他"
|
||||
epoch = int(row.get("quote_time_epoch") or 0)
|
||||
if epoch > 0:
|
||||
quote_time = datetime.fromtimestamp(epoch).astimezone().isoformat(timespec="seconds")
|
||||
else:
|
||||
quote_date = str(row.get("quote_date") or today)
|
||||
quote_time = f"{quote_date[:4]}-{quote_date[4:6]}-{quote_date[6:]}"
|
||||
return {
|
||||
"name": str(row.get("name") or name or "--"),
|
||||
"sector": sector,
|
||||
"price": price,
|
||||
"open": float(row.get("open") or 0),
|
||||
"high": float(row.get("high") or 0),
|
||||
"low": float(row.get("low") or 0),
|
||||
"change": round((price / previous_close - 1) * 100, 4),
|
||||
"volume": float(row.get("vol") or 0),
|
||||
"amount_billion": float(row.get("amount") or 0) / 100_000_000,
|
||||
"turnover_rate": float(row.get("turnover_rate") or 0),
|
||||
"quote_time": quote_time,
|
||||
}
|
||||
|
||||
def _intraday_realtime_stock_quote(
|
||||
self, code: str, today: str, payload: dict[str, Any]
|
||||
) -> dict[str, Any] | None:
|
||||
chart_data = getattr(self, "chart_data", None)
|
||||
if chart_data is None:
|
||||
return None
|
||||
try:
|
||||
chart = chart_data.stock_intraday(code)
|
||||
except (AttributeError, ChartDataError, Exception):
|
||||
return None
|
||||
points = [
|
||||
point
|
||||
for point in list(chart.get("points") or [])
|
||||
if str(point.get("date") or "").replace("-", "") == today
|
||||
]
|
||||
if not points:
|
||||
return None
|
||||
opens = [float(point.get("open") or 0) for point in points if float(point.get("open") or 0) > 0]
|
||||
highs = [float(point.get("high") or 0) for point in points if float(point.get("high") or 0) > 0]
|
||||
lows = [float(point.get("low") or 0) for point in points if float(point.get("low") or 0) > 0]
|
||||
closes = [float(point.get("close") or 0) for point in points if float(point.get("close") or 0) > 0]
|
||||
if not opens or not highs or not lows or not closes:
|
||||
return None
|
||||
price = closes[-1]
|
||||
previous_close = float(chart.get("previous_close") or 0)
|
||||
if previous_close <= 0:
|
||||
history = list(payload.get("prices") or [])
|
||||
previous_close = float((history[-1] if history else {}).get("close") or 0)
|
||||
if previous_close <= 0:
|
||||
return None
|
||||
volume = sum(float(point.get("volume") or 0) for point in points)
|
||||
amount = sum(float(point.get("amount") or 0) for point in points)
|
||||
if volume <= 0 and amount <= 0:
|
||||
return None
|
||||
try:
|
||||
name, sector = self._stock_identity(code, today)
|
||||
except Exception:
|
||||
name, sector = "--", "其他"
|
||||
return {
|
||||
"name": name,
|
||||
"sector": sector,
|
||||
"price": price,
|
||||
"open": opens[0],
|
||||
"high": max(highs),
|
||||
"low": min(lows),
|
||||
"change": round((price / previous_close - 1) * 100, 4),
|
||||
"volume": volume,
|
||||
"volume_unit": "lots",
|
||||
"amount_billion": amount / 100_000_000,
|
||||
"turnover_rate": 0.0,
|
||||
"quote_time": str(points[-1].get("date") or today),
|
||||
}
|
||||
|
||||
@staticmethod
|
||||
def _merge_realtime_stock_detail(
|
||||
payload: dict[str, Any], quote: dict[str, Any], trade_date: str
|
||||
@@ -1117,7 +924,6 @@ class MarketServiceMixin:
|
||||
**(payload.get("meta") or {}),
|
||||
"trade_date": display_date,
|
||||
"realtime": True,
|
||||
"notice": "",
|
||||
"updated_at": datetime.now().astimezone().isoformat(timespec="seconds"),
|
||||
}
|
||||
|
||||
|
||||
@@ -41,8 +41,6 @@ def official_catchup_due(today: str, snapshot: dict[str, object]) -> bool:
|
||||
actual == today
|
||||
and meta.get("limit_data_source") != "derived"
|
||||
and not meta.get("carried_forward")
|
||||
and not meta.get("realtime")
|
||||
and meta.get("mode") != "realtime"
|
||||
):
|
||||
return False
|
||||
return True
|
||||
|
||||
@@ -222,12 +222,12 @@
|
||||
{
|
||||
"provider": "eastmoney",
|
||||
"path": "backend/data/realtime.py",
|
||||
"runtime_role": "isolated realtime observation and intraday dashboard fallback"
|
||||
"runtime_role": "isolated realtime observation"
|
||||
},
|
||||
{
|
||||
"provider": "tencent",
|
||||
"path": "backend/data/realtime.py",
|
||||
"runtime_role": "index observation and intraday quote fallback"
|
||||
"runtime_role": "index observation fallback"
|
||||
}
|
||||
],
|
||||
"provider_domains": [
|
||||
@@ -508,8 +508,8 @@
|
||||
},
|
||||
{
|
||||
"path": "backend/data/providers/tushare_dashboard.py",
|
||||
"bytes": 31361,
|
||||
"lines": 732
|
||||
"bytes": 28234,
|
||||
"lines": 648
|
||||
},
|
||||
{
|
||||
"path": "backend/data/providers/tushare_industries.py",
|
||||
@@ -533,8 +533,8 @@
|
||||
},
|
||||
{
|
||||
"path": "frontend/pages/market/preview.js",
|
||||
"bytes": 18339,
|
||||
"lines": 450
|
||||
"bytes": 18178,
|
||||
"lines": 446
|
||||
},
|
||||
{
|
||||
"path": "backend/features/heaven/trend.py",
|
||||
@@ -551,11 +551,6 @@
|
||||
"bytes": 15311,
|
||||
"lines": 387
|
||||
},
|
||||
{
|
||||
"path": "frontend/shared/dashboard.js",
|
||||
"bytes": 15063,
|
||||
"lines": 321
|
||||
},
|
||||
{
|
||||
"path": "frontend/pages/pools/page.html",
|
||||
"bytes": 14942,
|
||||
@@ -566,6 +561,11 @@
|
||||
"bytes": 14743,
|
||||
"lines": 342
|
||||
},
|
||||
{
|
||||
"path": "frontend/shared/dashboard.js",
|
||||
"bytes": 14740,
|
||||
"lines": 316
|
||||
},
|
||||
{
|
||||
"path": "frontend/shared/admin.js",
|
||||
"bytes": 14410,
|
||||
@@ -631,11 +631,6 @@
|
||||
"bytes": 8357,
|
||||
"lines": 116
|
||||
},
|
||||
{
|
||||
"path": "backend/data/providers/tushare_indices.py",
|
||||
"bytes": 7823,
|
||||
"lines": 173
|
||||
},
|
||||
{
|
||||
"path": "backend/features/screener/formula.py",
|
||||
"bytes": 6983,
|
||||
@@ -643,8 +638,8 @@
|
||||
},
|
||||
{
|
||||
"path": "backend/data/providers/tushare_daily.py",
|
||||
"bytes": 6949,
|
||||
"lines": 168
|
||||
"bytes": 6837,
|
||||
"lines": 160
|
||||
},
|
||||
{
|
||||
"path": "backend/application.py",
|
||||
@@ -681,16 +676,21 @@
|
||||
"bytes": 6092,
|
||||
"lines": 138
|
||||
},
|
||||
{
|
||||
"path": "frontend/pages/market/stock-detail.js",
|
||||
"bytes": 6041,
|
||||
"lines": 134
|
||||
},
|
||||
{
|
||||
"path": "frontend/pages/dragon-tiger/page.html",
|
||||
"bytes": 5754,
|
||||
"lines": 85
|
||||
},
|
||||
{
|
||||
"path": "frontend/pages/market/stock-detail.js",
|
||||
"bytes": 5690,
|
||||
"lines": 124
|
||||
},
|
||||
{
|
||||
"path": "backend/data/providers/tushare_indices.py",
|
||||
"bytes": 5451,
|
||||
"lines": 118
|
||||
},
|
||||
{
|
||||
"path": "frontend/pages.config.js",
|
||||
"bytes": 5385,
|
||||
@@ -786,11 +786,6 @@
|
||||
"bytes": 2514,
|
||||
"lines": 63
|
||||
},
|
||||
{
|
||||
"path": "backend/data/providers/tushare_helpers.py",
|
||||
"bytes": 2360,
|
||||
"lines": 75
|
||||
},
|
||||
{
|
||||
"path": "backend/jobs/service.py",
|
||||
"bytes": 2337,
|
||||
@@ -816,6 +811,11 @@
|
||||
"bytes": 2165,
|
||||
"lines": 35
|
||||
},
|
||||
{
|
||||
"path": "backend/data/providers/tushare_helpers.py",
|
||||
"bytes": 2083,
|
||||
"lines": 64
|
||||
},
|
||||
{
|
||||
"path": "frontend/pages/market/breadth.js",
|
||||
"bytes": 2071,
|
||||
@@ -826,16 +826,16 @@
|
||||
"bytes": 1919,
|
||||
"lines": 45
|
||||
},
|
||||
{
|
||||
"path": "backend/jobs/refresh.py",
|
||||
"bytes": 1808,
|
||||
"lines": 48
|
||||
},
|
||||
{
|
||||
"path": "backend/features/system/routes.py",
|
||||
"bytes": 1791,
|
||||
"lines": 46
|
||||
},
|
||||
{
|
||||
"path": "backend/jobs/refresh.py",
|
||||
"bytes": 1728,
|
||||
"lines": 46
|
||||
},
|
||||
{
|
||||
"path": "backend/features/alerts/routes.py",
|
||||
"bytes": 1687,
|
||||
|
||||
@@ -13,13 +13,6 @@
|
||||
"valuation": { "read": false, "shadow": false },
|
||||
"moneyflow": { "read": false, "shadow": false },
|
||||
"auction": { "read": false, "shadow": false },
|
||||
"limit_events": { "read": false, "shadow": false },
|
||||
"popularity": { "read": false, "shadow": false },
|
||||
"dragon_tiger": { "read": false, "shadow": false },
|
||||
"sector_daily": { "read": false, "shadow": false },
|
||||
"quotes": { "read": false, "shadow": false },
|
||||
"index_quotes": { "read": false, "shadow": false },
|
||||
"intraday": { "read": false, "shadow": false },
|
||||
"status": { "read": false, "shadow": false }
|
||||
}
|
||||
}
|
||||
|
||||
@@ -213,12 +213,12 @@
|
||||
{
|
||||
"provider": "eastmoney",
|
||||
"path": "realtime_aggregator.py",
|
||||
"runtime_role": "isolated realtime observation and intraday dashboard fallback"
|
||||
"runtime_role": "isolated realtime observation"
|
||||
},
|
||||
{
|
||||
"provider": "tencent",
|
||||
"path": "realtime_aggregator.py",
|
||||
"runtime_role": "index observation and intraday quote fallback"
|
||||
"runtime_role": "index observation fallback"
|
||||
}
|
||||
],
|
||||
"llm_entrypoints": [
|
||||
|
||||
@@ -367,11 +367,7 @@ function selectStockPreviewChart(chart) {
|
||||
}
|
||||
} else if ((payload.prices || []).length) {
|
||||
setText("stockPreviewDate", payload.meta?.trade_date || "最新行情");
|
||||
const notice = String(payload.meta?.notice || "").trim();
|
||||
setText(
|
||||
"stockPreviewSource",
|
||||
notice ? `日 K 行情 · ${payload.prices.length} 个交易日 · ${notice}` : `日 K 行情 · ${payload.prices.length} 个交易日`,
|
||||
);
|
||||
setText("stockPreviewSource", `日 K 行情 · ${payload.prices.length} 个交易日`);
|
||||
drawDailyPreviewChart(payload.prices);
|
||||
} else {
|
||||
setText("stockPreviewDate", payload.meta?.trade_date || "最新行情");
|
||||
|
||||
@@ -52,11 +52,7 @@ async function openStock(code, fallback = null) {
|
||||
renderStockNotes(payload.notes || []);
|
||||
updateWatchButton();
|
||||
if (state.stockDetailChartMode === "daily") {
|
||||
const notice = String(payload.meta?.notice || "").trim();
|
||||
setText(
|
||||
"chartSource",
|
||||
notice ? `日 K 行情 · ${payload.prices.length} 个交易日 · ${notice}` : `日 K 行情 · ${payload.prices.length} 个交易日`,
|
||||
);
|
||||
setText("chartSource", `日 K 行情 · ${payload.prices.length} 个交易日`);
|
||||
requestAnimationFrame(() => drawPriceChart(payload.prices || []));
|
||||
}
|
||||
} catch (error) {
|
||||
@@ -73,13 +69,7 @@ async function selectStockDetailChart(mode) {
|
||||
syncDetailChartButtons("stock", selected);
|
||||
if (selected === "daily") {
|
||||
const prices = state.stockDetail?.prices || [];
|
||||
const notice = String(state.stockDetail?.meta?.notice || "").trim();
|
||||
setText(
|
||||
"chartSource",
|
||||
prices.length
|
||||
? (notice ? `日 K 行情 · ${prices.length} 个交易日 · ${notice}` : `日 K 行情 · ${prices.length} 个交易日`)
|
||||
: "正在加载行情",
|
||||
);
|
||||
setText("chartSource", prices.length ? `日 K 行情 · ${prices.length} 个交易日` : "正在加载行情");
|
||||
if (prices.length) requestAnimationFrame(() => drawPriceChart(prices));
|
||||
else clearPriceChart("正在加载日 K 数据");
|
||||
return;
|
||||
|
||||
@@ -67,11 +67,6 @@ async function startAdminRefresh() {
|
||||
const actualCompact = actualDate.replaceAll("-", "");
|
||||
const updated = formatTimestamp(meta.updated_at);
|
||||
const freshness = dashboardFreshnessMessage(meta);
|
||||
if (meta.realtime && actualCompact === requestedCompact && !meta.carried_forward) {
|
||||
setAdminRefreshStatus("success", `刷新成功:已获取 ${actualDate} 的盘中行情,更新时间 ${updated}`, "circle-check");
|
||||
showToast(`刷新成功:已获取 ${actualDate} 的盘中行情`);
|
||||
return;
|
||||
}
|
||||
if (freshness || actualCompact !== requestedCompact || meta.carried_forward || meta.limit_data_source === "derived") {
|
||||
setAdminRefreshStatus("warning", freshness || `部分正式数据尚未到齐,当前展示 ${actualDate || "最近可用数据"}`, "triangle-alert");
|
||||
setStatus(freshness || "部分正式数据尚未到齐,当前展示最近可用数据");
|
||||
|
||||
@@ -3,8 +3,7 @@ from __future__ import annotations
|
||||
import copy
|
||||
import threading
|
||||
import unittest
|
||||
from datetime import date, datetime, timedelta, timezone, time as dt_time
|
||||
from unittest.mock import patch
|
||||
from datetime import date, datetime, timedelta, timezone
|
||||
from pathlib import Path
|
||||
|
||||
from backend.features.market.service import MarketServiceMixin
|
||||
@@ -106,84 +105,18 @@ class FakeDerivedClient:
|
||||
}
|
||||
|
||||
|
||||
SHANGHAI = timezone(timedelta(hours=8))
|
||||
TRADE_DAY = date(2026, 9, 8)
|
||||
|
||||
|
||||
def at_clock(hour: int, minute: int, day: date = TRADE_DAY) -> datetime:
|
||||
return datetime(day.year, day.month, day.day, hour, minute, tzinfo=SHANGHAI)
|
||||
|
||||
|
||||
class FakeMissingDailyClient:
|
||||
def __init__(self, open_today: bool = True):
|
||||
self.open_today = open_today
|
||||
|
||||
def dashboard(self, trade_date: str):
|
||||
raise TushareError(f"No daily data returned for {trade_date}")
|
||||
|
||||
def resolve_trade_context(self, requested: str):
|
||||
if self.open_today:
|
||||
return requested, "20260907"
|
||||
return "20260907", "20260904"
|
||||
|
||||
|
||||
class FakeRealtimeTodayClient:
|
||||
def dashboard(self, trade_date: str):
|
||||
return {
|
||||
"meta": {
|
||||
"trade_date": f"{trade_date[:4]}-{trade_date[4:6]}-{trade_date[6:8]}",
|
||||
"requested_date": f"{trade_date[:4]}-{trade_date[4:6]}-{trade_date[6:8]}",
|
||||
"realtime": True,
|
||||
"mode": "realtime",
|
||||
"market_status": "trading",
|
||||
"notice": "盘中行情由 Tushare rt_k 实时计算;涨停原因、封板时间和开板次数以盘后榜单校正为准。",
|
||||
"updated_at": datetime.now().astimezone().isoformat(timespec="seconds"),
|
||||
},
|
||||
"overview": {"limit_up_count": 15},
|
||||
"limits": [{"code": "000001"}],
|
||||
"broken": [],
|
||||
"down_limits": [],
|
||||
"yesterday_limits": [],
|
||||
}
|
||||
|
||||
def resolve_trade_context(self, requested: str):
|
||||
return requested, "20260907"
|
||||
|
||||
|
||||
class FakeFreeRealtimeTodayClient:
|
||||
def dashboard(self, trade_date: str):
|
||||
return {
|
||||
"meta": {
|
||||
"trade_date": f"{trade_date[:4]}-{trade_date[4:6]}-{trade_date[6:8]}",
|
||||
"requested_date": f"{trade_date[:4]}-{trade_date[4:6]}-{trade_date[6:8]}",
|
||||
"realtime": True,
|
||||
"mode": "realtime",
|
||||
"quote_source": "eastmoney_clist",
|
||||
"source": "eastmoney",
|
||||
"market_status": "trading",
|
||||
"notice": "盘中行情由东财免费实时快照计算;涨停原因、封板时间和开板次数以盘后榜单校正为准。",
|
||||
"updated_at": datetime.now().astimezone().isoformat(timespec="seconds"),
|
||||
"indices": [{"code": "000001", "price": 3800.1, "change": 0.5}],
|
||||
},
|
||||
"overview": {"limit_up_count": 18, "up_count": 2100, "amount_billion": 12345.6},
|
||||
"limits": [{"code": "000001"}],
|
||||
"broken": [],
|
||||
"down_limits": [],
|
||||
"yesterday_limits": [],
|
||||
}
|
||||
|
||||
def resolve_trade_context(self, requested: str):
|
||||
return requested, "20260907"
|
||||
|
||||
|
||||
class SyncHarness(MarketServiceMixin):
|
||||
def __init__(self, client, latest=None, clock=None):
|
||||
def __init__(self, client, latest=None):
|
||||
self.configured = True
|
||||
self.sync_lock = threading.Lock()
|
||||
self.database = FakeSyncDatabase(latest)
|
||||
self._client = client
|
||||
self.current_user_id = 1
|
||||
self.clock = clock
|
||||
|
||||
def _tushare_client(self):
|
||||
return self._client
|
||||
@@ -209,161 +142,23 @@ class DashboardFreshnessTests(unittest.TestCase):
|
||||
self.assertEqual(harness.database.finished[0][0][1], "success")
|
||||
self.assertEqual(verified_dashboard_result(payload), payload)
|
||||
|
||||
def test_intraday_refresh_keeps_today_and_does_not_fall_back_to_yesterday(self):
|
||||
today = TRADE_DAY.strftime("%Y%m%d")
|
||||
def test_missing_official_data_keeps_previous_day_with_preparing_notice(self):
|
||||
today = date.today()
|
||||
previous = (today - timedelta(days=1)).strftime("%Y-%m-%d")
|
||||
latest = {
|
||||
"meta": {"trade_date": "2026-09-07", "source": "tushare"},
|
||||
"meta": {"trade_date": previous, "source": "tushare"},
|
||||
"overview": {"limit_up_count": 20},
|
||||
}
|
||||
harness = SyncHarness(
|
||||
FakeRealtimeTodayClient(),
|
||||
latest,
|
||||
clock=lambda: at_clock(10, 5),
|
||||
)
|
||||
payload = harness.sync_dashboard(today)
|
||||
harness = SyncHarness(FakeMissingDailyClient(), latest)
|
||||
payload = harness.sync_dashboard(today.strftime("%Y%m%d"))
|
||||
meta = payload["meta"]
|
||||
|
||||
self.assertFalse(meta.get("carried_forward"))
|
||||
self.assertTrue(meta["realtime"])
|
||||
self.assertEqual(meta["data_status"], "intraday")
|
||||
self.assertEqual(str(meta["trade_date"]).replace("-", ""), today)
|
||||
self.assertNotIn("今日数据正在准备", meta.get("display_notice") or "")
|
||||
self.assertEqual(harness.database.saved[0][0], today)
|
||||
|
||||
def test_intraday_free_source_keeps_today_and_indices(self):
|
||||
today = TRADE_DAY.strftime("%Y%m%d")
|
||||
latest = {
|
||||
"meta": {"trade_date": "2026-09-07", "source": "tushare"},
|
||||
"overview": {"limit_up_count": 20},
|
||||
}
|
||||
harness = SyncHarness(
|
||||
FakeFreeRealtimeTodayClient(),
|
||||
latest,
|
||||
clock=lambda: at_clock(10, 5),
|
||||
)
|
||||
payload = harness.sync_dashboard(today)
|
||||
meta = payload["meta"]
|
||||
self.assertFalse(meta.get("carried_forward"))
|
||||
self.assertTrue(meta["realtime"])
|
||||
self.assertEqual(meta["data_status"], "intraday")
|
||||
self.assertEqual(str(meta["trade_date"]).replace("-", ""), today)
|
||||
self.assertEqual(meta["quote_source"], "eastmoney_clist")
|
||||
self.assertEqual(payload["overview"]["amount_billion"], 12345.6)
|
||||
self.assertEqual(meta["indices"][0]["price"], 3800.1)
|
||||
self.assertEqual(harness.database.saved[0][0], today)
|
||||
|
||||
def test_intraday_missing_quotes_do_not_carry_yesterday(self):
|
||||
today = TRADE_DAY.strftime("%Y%m%d")
|
||||
latest = {
|
||||
"meta": {"trade_date": "2026-09-07", "source": "tushare"},
|
||||
"overview": {"limit_up_count": 20},
|
||||
}
|
||||
harness = SyncHarness(
|
||||
FakeMissingDailyClient(),
|
||||
latest,
|
||||
clock=lambda: at_clock(10, 5),
|
||||
)
|
||||
with self.assertRaises(ValueError) as ctx:
|
||||
harness.sync_dashboard(today)
|
||||
self.assertIn("当天盘中行情", str(ctx.exception))
|
||||
self.assertFalse(harness.database.saved)
|
||||
|
||||
def test_intraday_keeps_existing_today_snapshot_when_refresh_fails(self):
|
||||
today = TRADE_DAY.strftime("%Y%m%d")
|
||||
existing = {
|
||||
"meta": {
|
||||
"trade_date": "2026-09-08",
|
||||
"realtime": True,
|
||||
"mode": "realtime",
|
||||
"source": "tushare",
|
||||
},
|
||||
"overview": {"limit_up_count": 11},
|
||||
"limits": [{"code": "600000"}],
|
||||
"broken": [],
|
||||
"down_limits": [],
|
||||
"yesterday_limits": [],
|
||||
}
|
||||
harness = SyncHarness(
|
||||
FakeMissingDailyClient(),
|
||||
clock=lambda: at_clock(10, 5),
|
||||
)
|
||||
harness.database.get_snapshot = lambda *_args, **_kwargs: copy.deepcopy(existing)
|
||||
payload = harness.sync_dashboard(today)
|
||||
meta = payload["meta"]
|
||||
self.assertEqual(str(meta["trade_date"]).replace("-", ""), today)
|
||||
self.assertTrue(meta["realtime"])
|
||||
self.assertEqual(meta["data_status"], "intraday")
|
||||
self.assertFalse(meta.get("carried_forward"))
|
||||
|
||||
def test_lunch_and_after_hours_keep_today_until_official_arrives(self):
|
||||
today = TRADE_DAY.strftime("%Y%m%d")
|
||||
for clock in (lambda: at_clock(12, 0), lambda: at_clock(16, 10)):
|
||||
harness = SyncHarness(
|
||||
FakeRealtimeTodayClient(),
|
||||
clock=clock,
|
||||
)
|
||||
payload = harness.sync_dashboard(today)
|
||||
self.assertEqual(str(payload["meta"]["trade_date"]).replace("-", ""), today)
|
||||
self.assertFalse(payload["meta"].get("carried_forward"))
|
||||
|
||||
def test_preopen_and_weekend_still_carry_last_session(self):
|
||||
latest = {
|
||||
"meta": {"trade_date": "2026-09-07", "source": "tushare"},
|
||||
"overview": {"limit_up_count": 20},
|
||||
}
|
||||
preopen = SyncHarness(
|
||||
FakeMissingDailyClient(),
|
||||
latest,
|
||||
clock=lambda: at_clock(8, 30),
|
||||
)
|
||||
preopen_payload = preopen.sync_dashboard(TRADE_DAY.strftime("%Y%m%d"))
|
||||
self.assertTrue(preopen_payload["meta"]["carried_forward"])
|
||||
self.assertEqual(preopen_payload["meta"]["data_status"], "preparing")
|
||||
self.assertIn("今日数据正在准备,当前展示", preopen_payload["meta"]["display_notice"])
|
||||
|
||||
weekend = SyncHarness(
|
||||
FakeMissingDailyClient(open_today=False),
|
||||
latest,
|
||||
clock=lambda: at_clock(10, 5, date(2026, 9, 5)),
|
||||
)
|
||||
weekend_payload = weekend.sync_dashboard("20260905")
|
||||
self.assertTrue(weekend_payload["meta"]["carried_forward"])
|
||||
|
||||
def test_history_date_still_uses_official_or_preparing_notice(self):
|
||||
latest = {
|
||||
"meta": {"trade_date": "2026-09-01", "source": "tushare"},
|
||||
"overview": {"limit_up_count": 8},
|
||||
}
|
||||
harness = SyncHarness(
|
||||
FakeMissingDailyClient(),
|
||||
latest,
|
||||
clock=lambda: at_clock(10, 5),
|
||||
)
|
||||
payload = harness.sync_dashboard("20260902")
|
||||
self.assertTrue(payload["meta"]["carried_forward"])
|
||||
self.assertIn("所选日期数据尚未到齐", payload["meta"]["display_notice"])
|
||||
|
||||
def test_carried_today_snapshot_is_retried_immediately_in_session(self):
|
||||
today = TRADE_DAY.strftime("%Y%m%d")
|
||||
snapshot = {
|
||||
"meta": {
|
||||
"source": "tushare",
|
||||
"trade_date": "2026-09-07",
|
||||
"carried_forward": True,
|
||||
"requested_date": "2026-09-08",
|
||||
"updated_at": at_clock(10, 0).isoformat(),
|
||||
},
|
||||
"overview": {"limit_up_count": 1},
|
||||
}
|
||||
harness = SyncHarness(
|
||||
FakeRealtimeTodayClient(),
|
||||
clock=lambda: at_clock(10, 5),
|
||||
)
|
||||
harness.database.get_snapshot = lambda *_args, **_kwargs: copy.deepcopy(snapshot)
|
||||
payload = harness.get_dashboard(today)
|
||||
self.assertEqual(str(payload["meta"]["trade_date"]).replace("-", ""), today)
|
||||
self.assertEqual(payload["meta"]["data_status"], "intraday")
|
||||
self.assertTrue(harness.database.saved)
|
||||
self.assertTrue(meta["carried_forward"])
|
||||
self.assertEqual(meta["data_status"], "preparing")
|
||||
self.assertIn("今日数据正在准备,当前展示", meta["display_notice"])
|
||||
self.assertIn("月", meta["display_notice"])
|
||||
self.assertNotIn("No daily data", meta["display_notice"])
|
||||
self.assertNotEqual(verified_dashboard_result(payload).get("status"), "failed")
|
||||
|
||||
def test_weekend_carry_is_not_labeled_as_preparing(self):
|
||||
snapshot = {
|
||||
@@ -405,43 +200,19 @@ class DashboardFreshnessTests(unittest.TestCase):
|
||||
{"meta": {"trade_date": iso, "limit_data_source": "derived"}},
|
||||
)
|
||||
now = datetime.now().astimezone().time().replace(tzinfo=None)
|
||||
if dt_time(15, 5) <= now < dt_time(22, 0):
|
||||
if datetime.strptime("15:05", "%H:%M").time() <= now < datetime.strptime("22:00", "%H:%M").time():
|
||||
self.assertFalse(due)
|
||||
self.assertTrue(derived_due)
|
||||
else:
|
||||
self.assertFalse(due)
|
||||
self.assertFalse(derived_due)
|
||||
|
||||
def test_official_catchup_is_due_for_intraday_snapshot_after_close(self):
|
||||
today = TRADE_DAY.strftime("%Y%m%d")
|
||||
snapshot = {
|
||||
"meta": {
|
||||
"trade_date": "2026-09-08",
|
||||
"realtime": True,
|
||||
"mode": "realtime",
|
||||
}
|
||||
}
|
||||
with patch("backend.jobs.refresh.datetime") as mocked:
|
||||
mocked.now.return_value = at_clock(16, 10)
|
||||
mocked.strptime = datetime.strptime
|
||||
self.assertTrue(official_catchup_due(today, snapshot))
|
||||
official = {
|
||||
"meta": {
|
||||
"trade_date": "2026-09-08",
|
||||
"limit_data_source": "official",
|
||||
"realtime": False,
|
||||
}
|
||||
}
|
||||
self.assertFalse(official_catchup_due(today, official))
|
||||
|
||||
|
||||
class FrontendRefreshCopyTests(unittest.TestCase):
|
||||
def test_dashboard_script_distinguishes_partial_from_failure(self):
|
||||
script = (Path(__file__).resolve().parents[1] / "frontend" / "shared" / "dashboard.js").read_text(encoding="utf-8")
|
||||
self.assertIn("今日数据正在准备,当前展示", script)
|
||||
self.assertIn("部分正式数据尚未到齐", script)
|
||||
self.assertIn("盘中行情", script)
|
||||
self.assertIn("meta.realtime && actualCompact === requestedCompact", script)
|
||||
self.assertIn('job.status === "failed"', script)
|
||||
failed_block = script.split("if (job.status === \"failed\")", 1)[1].split("const query", 1)[0]
|
||||
self.assertIn("后台刷新失败", failed_block)
|
||||
|
||||
@@ -2,8 +2,7 @@ from __future__ import annotations
|
||||
|
||||
import unittest
|
||||
|
||||
from backend.data.providers.ifind_client import IfindHttpClient
|
||||
from backend.features.market.charts import ChartDataError, EastmoneyChartClient, HIS_TRENDS_URL, MarketChartClient, TRENDS_URL
|
||||
from backend.features.market.charts import ChartDataError, EastmoneyChartClient
|
||||
from server import DashboardService
|
||||
|
||||
|
||||
@@ -73,141 +72,6 @@ class ChartDataProviderTests(unittest.TestCase):
|
||||
self.client.stock_intraday("abc")
|
||||
|
||||
|
||||
class LookbackChartClient(EastmoneyChartClient):
|
||||
def __init__(self) -> None:
|
||||
super().__init__(cache_ttl_seconds=20)
|
||||
self.requests: list[tuple[str, dict[str, str]]] = []
|
||||
|
||||
def _request_json(self, url, params, referer):
|
||||
self.requests.append((url, params))
|
||||
if url == TRENDS_URL and params.get("ndays") == "1":
|
||||
return {"data": {"code": "601318", "name": "中国平安", "preClose": 56.0, "trends": []}}
|
||||
if url == TRENDS_URL and params.get("ndays") == "5":
|
||||
return {"data": {"code": "601318", "name": "中国平安", "preClose": 56.0, "trends": []}}
|
||||
if url == HIS_TRENDS_URL:
|
||||
return {
|
||||
"data": {
|
||||
"code": "601318",
|
||||
"name": "中国平安",
|
||||
"preClose": 55.8,
|
||||
"trends": [
|
||||
"2026-09-07 09:30,55.80,55.90,56.00,55.70,100,5580.00,55.900",
|
||||
"2026-09-07 15:00,56.10,56.20,56.30,56.00,200,11240.00,56.150",
|
||||
"2026-09-08 09:30,0,0,0,0,0,0.00,0",
|
||||
],
|
||||
}
|
||||
}
|
||||
raise ChartDataError("unexpected url")
|
||||
|
||||
|
||||
class ChartLookbackTests(unittest.TestCase):
|
||||
def setUp(self) -> None:
|
||||
EastmoneyChartClient._cache.clear()
|
||||
self.client = LookbackChartClient()
|
||||
|
||||
def test_empty_today_falls_back_to_latest_available_session(self):
|
||||
payload = self.client.stock_intraday("601318")
|
||||
urls = [url for url, _ in self.client.requests]
|
||||
self.assertEqual(urls[0], TRENDS_URL)
|
||||
self.assertEqual(self.client.requests[0][1]["ndays"], "1")
|
||||
self.assertEqual(urls[1], TRENDS_URL)
|
||||
self.assertEqual(self.client.requests[1][1]["ndays"], "5")
|
||||
self.assertEqual(urls[2], HIS_TRENDS_URL)
|
||||
self.assertEqual(payload["trade_date"], "2026-09-07")
|
||||
self.assertEqual([point["time"] for point in payload["points"]], ["09:30", "15:00"])
|
||||
self.assertEqual(payload["points"][0]["close"], 55.9)
|
||||
|
||||
def test_delay_multiday_can_recover_without_his(self):
|
||||
class DelayFive(EastmoneyChartClient):
|
||||
def __init__(self):
|
||||
super().__init__(cache_ttl_seconds=20)
|
||||
self.requests = []
|
||||
|
||||
def _request_json(self, url, params, referer):
|
||||
self.requests.append((url, params))
|
||||
if params.get("ndays") == "1":
|
||||
return {"data": {"code": "000001", "name": "平安银行", "preClose": 11.7, "trends": []}}
|
||||
return {
|
||||
"data": {
|
||||
"code": "000001",
|
||||
"name": "平安银行",
|
||||
"preClose": 11.5,
|
||||
"trends": [
|
||||
"2026-09-07 09:30,11.50,11.60,11.70,11.40,100,1160.00,11.600",
|
||||
"2026-09-07 15:00,11.70,11.80,11.90,11.60,200,2360.00,11.750",
|
||||
],
|
||||
}
|
||||
}
|
||||
|
||||
EastmoneyChartClient._cache.clear()
|
||||
client = DelayFive()
|
||||
payload = client.stock_intraday("000001")
|
||||
self.assertEqual(payload["trade_date"], "2026-09-07")
|
||||
self.assertEqual(len(payload["points"]), 2)
|
||||
self.assertEqual([url for url, _ in client.requests], [TRENDS_URL, TRENDS_URL])
|
||||
|
||||
def test_sh_sz_cyb_codes_use_correct_secid(self):
|
||||
for code, secid in (("601318", "1.601318"), ("000001", "0.000001"), ("300750", "0.300750")):
|
||||
EastmoneyChartClient._cache.clear()
|
||||
client = LookbackChartClient()
|
||||
client.stock_intraday(code)
|
||||
self.assertEqual(client.requests[0][1]["secid"], secid)
|
||||
|
||||
|
||||
class FakeHub:
|
||||
def __init__(self, chart=None, error=None):
|
||||
self.chart = chart
|
||||
self.error = error
|
||||
self.calls: list[str] = []
|
||||
|
||||
def try_intraday(self, code):
|
||||
self.calls.append(code)
|
||||
if self.error:
|
||||
raise self.error
|
||||
return self.chart
|
||||
|
||||
|
||||
class DatahubChartFallbackTests(unittest.TestCase):
|
||||
def setUp(self) -> None:
|
||||
EastmoneyChartClient._cache.clear()
|
||||
|
||||
def test_datahub_success_skips_old_channel(self):
|
||||
hub = FakeHub(
|
||||
{
|
||||
"entity_type": "stock",
|
||||
"identifier": "601318",
|
||||
"name": "中国平安",
|
||||
"code": "601318",
|
||||
"trade_date": "2026-09-08",
|
||||
"previous_close": 56.36,
|
||||
"points": [{"date": "2026-09-08", "time": "09:30", "close": 56.5, "average": 56.4}],
|
||||
"source": "datahub",
|
||||
}
|
||||
)
|
||||
fallback = LookbackChartClient()
|
||||
client = MarketChartClient(IfindHttpClient(), fallback, hub)
|
||||
payload = client.stock_intraday("601318")
|
||||
self.assertEqual(payload["source"], "datahub")
|
||||
self.assertEqual(hub.calls, ["601318"])
|
||||
self.assertEqual(fallback.requests, [])
|
||||
|
||||
def test_datahub_timeout_or_empty_falls_back_to_eastmoney(self):
|
||||
fallback = LookbackChartClient()
|
||||
for hub in (
|
||||
FakeHub(chart=None),
|
||||
FakeHub(error=RuntimeError("timeout")),
|
||||
FakeHub(error=RuntimeError("datahub exploded")),
|
||||
FakeHub(chart={"points": []}),
|
||||
):
|
||||
EastmoneyChartClient._cache.clear()
|
||||
fallback.requests.clear()
|
||||
client = MarketChartClient(IfindHttpClient(), fallback, hub)
|
||||
payload = client.stock_intraday("000001")
|
||||
self.assertEqual(payload["trade_date"], "2026-09-07")
|
||||
self.assertGreaterEqual(len(payload["points"]), 1)
|
||||
self.assertTrue(fallback.requests)
|
||||
|
||||
|
||||
class ChartServiceStub:
|
||||
@staticmethod
|
||||
def _payload(code: str, name: str):
|
||||
|
||||
@@ -64,11 +64,9 @@ class FakeClient(DatahubClient):
|
||||
meta={"tier": "official", "trade_date": "20240902", "stale": False, "staleness_seconds": 0},
|
||||
)
|
||||
self.paths: list[str] = []
|
||||
self.calls: list[tuple[str, dict[str, Any]]] = []
|
||||
|
||||
def get(self, path: str, params: dict[str, Any] | None = None) -> DatahubResponse:
|
||||
self.paths.append(path)
|
||||
self.calls.append((path, {key: value for key, value in (params or {}).items()}))
|
||||
if TOKEN in json.dumps(params or {}) or TOKEN in path:
|
||||
raise AssertionError("token leaked into url")
|
||||
if self.error:
|
||||
@@ -292,8 +290,8 @@ class DatahubBridgeTests(unittest.TestCase):
|
||||
self.assertEqual(canonical["vol"], 100000.0)
|
||||
self.assertEqual(canonical["amount"], 2000000.0)
|
||||
|
||||
def test_heaven_can_use_hub_when_dataset_flag_is_on(self) -> None:
|
||||
"""问天按数据依赖接入:已映射 API 跟随开关,不再整栈强制旧链路。"""
|
||||
def test_heaven_keeps_legacy_on_first_batch_even_when_read_flag_is_on(self) -> None:
|
||||
"""问天未永久冻结;首批只读接入仍走旧链路,后续迁移可以纳入。"""
|
||||
self.assertTrue(looks_like_heaven("backend.features.heaven.market_context", "backend/features/heaven/market_context.py"))
|
||||
self.assertFalse(looks_like_heaven("backend.features.market.service", "backend/features/market/service.py"))
|
||||
client = FakeClient()
|
||||
@@ -304,8 +302,7 @@ class DatahubBridgeTests(unittest.TestCase):
|
||||
)
|
||||
rows = wrapped.query("daily", {"trade_date": "20240902"}, "amount")
|
||||
self.assertEqual(rows[0]["amount"], 2000.0)
|
||||
self.assertEqual(client.paths, ["/v1/bars/daily"])
|
||||
self.assertEqual(legacy.calls, [])
|
||||
self.assertEqual(client.paths, [])
|
||||
|
||||
def test_status_flag_does_not_run_when_off_and_falls_back_when_on(self) -> None:
|
||||
off = DatahubBridge(flags(), FakeClient(error=DatahubError("UNAVAILABLE", "down")))
|
||||
@@ -351,69 +348,6 @@ class DatahubBridgeTests(unittest.TestCase):
|
||||
self.assertEqual(rows[0]["amount"], 2000.0)
|
||||
self.assertEqual(len(legacy.calls), 1)
|
||||
|
||||
def test_try_intraday_respects_switch_and_falls_back_on_bad_payload(self) -> None:
|
||||
closed = DatahubBridge(flags(), FakeClient(error=DatahubError("INTERNAL", "should not run")))
|
||||
self.assertIsNone(closed.try_intraday("601318"))
|
||||
|
||||
empty = DatahubBridge(
|
||||
flags(intraday=(True, False)),
|
||||
FakeClient(response=DatahubResponse(data={"points": []}, meta={"stale": False})),
|
||||
)
|
||||
self.assertIsNone(empty.try_intraday("601318"))
|
||||
|
||||
stale = DatahubBridge(
|
||||
flags(intraday=(True, False)),
|
||||
FakeClient(response=DatahubResponse(
|
||||
data={
|
||||
"entity_type": "stock",
|
||||
"code": "601318",
|
||||
"trade_date": "2026-09-07",
|
||||
"previous_close": 55.8,
|
||||
"points": [{"date": "2026-09-07", "time": "09:30", "close": 55.9, "avg_price": 55.85}],
|
||||
},
|
||||
meta={"stale": True},
|
||||
)),
|
||||
)
|
||||
self.assertIsNone(stale.try_intraday("601318"))
|
||||
|
||||
ok = DatahubBridge(
|
||||
flags(intraday=(True, False)),
|
||||
FakeClient(response=DatahubResponse(
|
||||
data={
|
||||
"entity_type": "stock",
|
||||
"identifier": "601318",
|
||||
"name": "中国平安",
|
||||
"code": "601318",
|
||||
"trade_date": "2026-09-08",
|
||||
"previous_close": 56.36,
|
||||
"points": [
|
||||
{"date": "2026-09-08", "time": "09:30", "close": 0},
|
||||
{"date": "2026-09-08", "time": "09:31", "close": 56.5, "avg_price": 56.4},
|
||||
],
|
||||
},
|
||||
meta={"stale": False},
|
||||
)),
|
||||
)
|
||||
chart = ok.try_intraday("601318")
|
||||
self.assertEqual(chart["source"], "datahub")
|
||||
self.assertEqual(len(chart["points"]), 1)
|
||||
self.assertEqual(chart["points"][0]["average"], 56.4)
|
||||
self.assertEqual(ok.client.paths, ["/v1/intraday/points"])
|
||||
self.assertEqual(ok.client.calls, [("/v1/intraday/points", {"code": "601318"})])
|
||||
self.assertNotIn("date", ok.client.calls[0][1])
|
||||
|
||||
timeout = DatahubBridge(
|
||||
flags(intraday=(True, False)),
|
||||
FakeClient(error=DatahubError("TIMEOUT", "datahub request timed out")),
|
||||
)
|
||||
self.assertIsNone(timeout.try_intraday("601318"))
|
||||
broken = DatahubBridge(
|
||||
flags(intraday=(True, False)),
|
||||
FakeClient(error=DatahubError("INTERNAL", "datahub exploded")),
|
||||
)
|
||||
self.assertIsNone(broken.try_intraday("601318"))
|
||||
self.assertFalse(DatahubSettings.load(environ={}, credentials={}).flags("intraday").read)
|
||||
|
||||
def test_features_do_not_import_datahub_client(self) -> None:
|
||||
violations = []
|
||||
for path in (ROOT / "backend" / "features").rglob("*.py"):
|
||||
|
||||
@@ -3,10 +3,9 @@ from __future__ import annotations
|
||||
import http.client
|
||||
import json
|
||||
import unittest
|
||||
from datetime import datetime
|
||||
from unittest.mock import MagicMock, patch
|
||||
|
||||
from backend.data.realtime import RealtimeAggregateError, WebRealtimeAggregator
|
||||
from backend.data.realtime import WebRealtimeAggregator
|
||||
from backend.features.heaven.engine import _market_line_scores, build_manual_market_hexagram
|
||||
from server import DashboardService
|
||||
from backend.data.providers.tushare_client import (
|
||||
@@ -378,87 +377,6 @@ class RealtimeAggregatorTests(unittest.TestCase):
|
||||
self.assertEqual(rows[0]["quote_time"][:10], "2026-07-20")
|
||||
self.assertAlmostEqual(rows[0]["amount_billion"], 12946.52)
|
||||
|
||||
@patch.object(WebRealtimeAggregator, "_get_json")
|
||||
def test_eastmoney_market_quotes_normalize_and_keep_expected_date(self, get_json: MagicMock):
|
||||
epoch = datetime(2026, 7, 20, 10, 5).timestamp()
|
||||
rows = []
|
||||
for index in range(200):
|
||||
sz = index < 100
|
||||
rows.append(
|
||||
{
|
||||
"f12": f"{index:06d}" if sz else f"{600000 + index - 100:06d}",
|
||||
"f13": 0 if sz else 1,
|
||||
"f14": f"股票{index}",
|
||||
"f2": 11.2,
|
||||
"f3": 2.0,
|
||||
"f5": 10,
|
||||
"f6": 50000000,
|
||||
"f15": 11.3,
|
||||
"f16": 11.0,
|
||||
"f17": 11.1,
|
||||
"f18": 11.0,
|
||||
"f124": epoch,
|
||||
}
|
||||
)
|
||||
def fake_get_json(_url, params, referer=""):
|
||||
page = int(params.get("pn") or 1)
|
||||
start = (page - 1) * 100
|
||||
return {"rc": 0, "data": {"total": 200, "diff": rows[start:start + 100]}}
|
||||
|
||||
get_json.side_effect = fake_get_json
|
||||
aggregator = WebRealtimeAggregator()
|
||||
aggregator._response_cache.clear()
|
||||
quotes = aggregator.eastmoney_market_quotes("20260720")
|
||||
self.assertEqual(len(quotes), 200)
|
||||
self.assertEqual(quotes[0]["ts_code"], "000000.SZ")
|
||||
self.assertTrue(quotes[100]["ts_code"].endswith(".SH"))
|
||||
self.assertEqual(quotes[0]["vol"], 1000)
|
||||
self.assertEqual(quotes[0]["quote_date"], "20260720")
|
||||
|
||||
@patch.object(WebRealtimeAggregator, "_get_text")
|
||||
def test_tencent_stock_quote_keeps_expected_date(self, get_text: MagicMock):
|
||||
fields = [""] * 38
|
||||
fields[1] = "浦发银行"
|
||||
fields[2] = "600000"
|
||||
fields[3] = "11.20"
|
||||
fields[4] = "11.00"
|
||||
fields[5] = "11.10"
|
||||
fields[6] = "1234"
|
||||
fields[30] = "20260720103000"
|
||||
fields[33] = "11.30"
|
||||
fields[34] = "11.00"
|
||||
fields[37] = "1380"
|
||||
get_text.return_value = (f'v_sh600000="{"~".join(fields)}";', 0)
|
||||
|
||||
quote = WebRealtimeAggregator().tencent_stock_quote("600000", "20260720")
|
||||
|
||||
self.assertEqual(quote["ts_code"], "600000.SH")
|
||||
self.assertEqual(quote["quote_date"], "20260720")
|
||||
self.assertEqual(quote["vol"], 123400)
|
||||
self.assertAlmostEqual(quote["amount"], 13_800_000)
|
||||
|
||||
@patch.object(WebRealtimeAggregator, "_get_json")
|
||||
def test_eastmoney_stock_quote_rejects_stale_date(self, get_json: MagicMock):
|
||||
epoch = datetime(2026, 7, 19, 15, 0).timestamp()
|
||||
get_json.return_value = {
|
||||
"rc": 0,
|
||||
"data": {
|
||||
"f43": 11.2,
|
||||
"f44": 11.3,
|
||||
"f45": 11.0,
|
||||
"f46": 11.1,
|
||||
"f47": 10,
|
||||
"f48": 50000000,
|
||||
"f57": "300750",
|
||||
"f58": "宁德时代",
|
||||
"f60": 11.0,
|
||||
"f86": epoch,
|
||||
},
|
||||
}
|
||||
|
||||
with self.assertRaises(RealtimeAggregateError):
|
||||
WebRealtimeAggregator().eastmoney_stock_quote("300750.SZ", "20260720")
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
unittest.main()
|
||||
|
||||
@@ -1,16 +1,8 @@
|
||||
from __future__ import annotations
|
||||
|
||||
import unittest
|
||||
from datetime import datetime, timedelta, timezone
|
||||
|
||||
from backend.data.providers.tushare_client import TushareClient
|
||||
from backend.data.providers.tushare_helpers import calendar_is_open
|
||||
from backend.data.providers.tushare_transport import TushareError
|
||||
from backend.data.realtime import (
|
||||
RealtimeAggregateError,
|
||||
_normalize_eastmoney_quote,
|
||||
_parse_tencent_stock_quote,
|
||||
)
|
||||
|
||||
|
||||
class FakeRealtimeClient(TushareClient):
|
||||
@@ -89,65 +81,6 @@ class FakeRealtimeClient(TushareClient):
|
||||
raise AssertionError(f"Unexpected API call: {api_name} {params}")
|
||||
|
||||
|
||||
FREE_QUOTES = [
|
||||
{
|
||||
"ts_code": "000001.SZ", "name": "甲", "pre_close": 10.0,
|
||||
"open": 10.1, "high": 11.0, "low": 10.0, "close": 11.0,
|
||||
"vol": 1000, "amount": 100000000, "num": 10,
|
||||
"quote_date": "20260720",
|
||||
},
|
||||
{
|
||||
"ts_code": "000002.SZ", "name": "乙", "pre_close": 20.0,
|
||||
"open": 19.5, "high": 20.0, "low": 18.0, "close": 18.0,
|
||||
"vol": 2000, "amount": 200000000, "num": 20,
|
||||
"quote_date": "20260720",
|
||||
},
|
||||
{
|
||||
"ts_code": "000003.SZ", "name": "丙", "pre_close": 30.0,
|
||||
"open": 31.0, "high": 33.0, "low": 30.0, "close": 32.0,
|
||||
"vol": 3000, "amount": 300000000, "num": 30,
|
||||
"quote_date": "20260720",
|
||||
},
|
||||
]
|
||||
|
||||
|
||||
class FakeFreeAggregator:
|
||||
def __init__(self, quotes=None, fail=False):
|
||||
self.quotes = list(quotes if quotes is not None else FREE_QUOTES)
|
||||
self.fail = fail
|
||||
self.calls = 0
|
||||
|
||||
def eastmoney_market_quotes(self, expected_date=""):
|
||||
self.calls += 1
|
||||
if self.fail:
|
||||
raise RealtimeAggregateError("eastmoney down")
|
||||
if expected_date and self.quotes:
|
||||
dated = [
|
||||
row for row in self.quotes
|
||||
if str(row.get("quote_date") or "") == str(expected_date).replace("-", "")
|
||||
]
|
||||
if dated:
|
||||
return dated
|
||||
return list(self.quotes)
|
||||
|
||||
def tencent_market_quotes(self, codes, expected_date=""):
|
||||
return self.eastmoney_market_quotes(expected_date)
|
||||
|
||||
def eastmoney_indices(self):
|
||||
return [
|
||||
{
|
||||
"code": "000001",
|
||||
"name": "上证指数",
|
||||
"price": 3800.12,
|
||||
"change": 0.85,
|
||||
"previous_close": 3768.0,
|
||||
"amount_billion": 4200.5,
|
||||
"quote_time": "2026-07-20T10:05:00+08:00",
|
||||
"source": "eastmoney_push2",
|
||||
}
|
||||
]
|
||||
|
||||
|
||||
class RealtimeDashboardTests(unittest.TestCase):
|
||||
def setUp(self):
|
||||
TushareClient._realtime_reference_cache.clear()
|
||||
@@ -197,173 +130,6 @@ class RealtimeDashboardTests(unittest.TestCase):
|
||||
self.assertEqual(dashboard["meta"]["limit_data_source"], "derived")
|
||||
self.assertIn("日线数据推算", dashboard["meta"]["notice"])
|
||||
|
||||
def test_calendar_open_flag_accepts_string_and_bool(self):
|
||||
self.assertTrue(calendar_is_open(1))
|
||||
self.assertTrue(calendar_is_open("1"))
|
||||
self.assertTrue(calendar_is_open(True))
|
||||
self.assertFalse(calendar_is_open(0))
|
||||
self.assertFalse(calendar_is_open("0"))
|
||||
self.assertFalse(calendar_is_open(False))
|
||||
|
||||
original_query = self.client.query
|
||||
|
||||
def query(api_name, params=None, fields=""):
|
||||
if api_name == "trade_cal":
|
||||
return [
|
||||
{
|
||||
"cal_date": params.get("start_date"),
|
||||
"is_open": "1",
|
||||
"pretrade_date": "20260907",
|
||||
}
|
||||
]
|
||||
return original_query(api_name, params, fields)
|
||||
|
||||
self.client.query = query
|
||||
trade_date, previous = self.client.resolve_trade_context("20260908")
|
||||
self.assertEqual(trade_date, "20260908")
|
||||
self.assertEqual(previous, "20260907")
|
||||
|
||||
def test_session_clock_uses_realtime_until_official_window(self):
|
||||
today = "20260908"
|
||||
self.client.clock = lambda: datetime(
|
||||
2026, 9, 8, 10, 5, tzinfo=timezone(timedelta(hours=8))
|
||||
)
|
||||
self.assertTrue(self.client.should_use_realtime(today, today))
|
||||
self.client.clock = lambda: datetime(
|
||||
2026, 9, 8, 16, 10, tzinfo=timezone(timedelta(hours=8))
|
||||
)
|
||||
self.assertFalse(self.client.should_use_realtime(today, today))
|
||||
|
||||
def test_realtime_dashboard_survives_missing_limit_table(self):
|
||||
original_query = self.client.query
|
||||
|
||||
def query(api_name, params=None, fields=""):
|
||||
if api_name == "stk_limit":
|
||||
return []
|
||||
return original_query(api_name, params, fields)
|
||||
|
||||
self.client.query = query
|
||||
TushareClient._realtime_reference_cache.clear()
|
||||
dashboard = self.client._realtime_dashboard("20260720", "20260720", "20260717")
|
||||
self.assertTrue(dashboard["meta"]["realtime"])
|
||||
self.assertEqual(dashboard["meta"]["quote_count"], 3)
|
||||
self.assertEqual(dashboard["overview"]["limit_up_count"], 0)
|
||||
|
||||
def test_rt_k_permission_error_falls_back_to_free_quotes(self):
|
||||
original_query = self.client.query
|
||||
|
||||
def query(api_name, params=None, fields=""):
|
||||
if api_name == "rt_k":
|
||||
raise TushareError("没有接口访问权限")
|
||||
return original_query(api_name, params, fields)
|
||||
|
||||
self.client.query = query
|
||||
self.client.realtime_aggregator = FakeFreeAggregator()
|
||||
TushareClient._realtime_reference_cache.clear()
|
||||
dashboard = self.client._realtime_dashboard("20260720", "20260720", "20260717")
|
||||
|
||||
self.assertTrue(dashboard["meta"]["realtime"])
|
||||
self.assertEqual(dashboard["meta"]["quote_source"], "eastmoney_clist")
|
||||
self.assertEqual(dashboard["meta"]["trade_date"], "2026-07-20")
|
||||
self.assertEqual(dashboard["meta"]["quote_count"], 3)
|
||||
self.assertEqual(dashboard["overview"]["limit_up_count"], 1)
|
||||
self.assertEqual(dashboard["overview"]["limit_down_count"], 1)
|
||||
self.assertEqual(dashboard["overview"]["amount_billion"], 6.0)
|
||||
self.assertIn("东财免费实时", dashboard["meta"]["notice"])
|
||||
self.assertEqual(dashboard["meta"]["indices"][0]["price"], 3800.12)
|
||||
|
||||
def test_rt_k_empty_result_falls_back_to_free_quotes(self):
|
||||
original_query = self.client.query
|
||||
|
||||
def query(api_name, params=None, fields=""):
|
||||
if api_name == "rt_k":
|
||||
return []
|
||||
return original_query(api_name, params, fields)
|
||||
|
||||
self.client.query = query
|
||||
self.client.realtime_aggregator = FakeFreeAggregator()
|
||||
TushareClient._realtime_reference_cache.clear()
|
||||
dashboard = self.client._realtime_dashboard("20260720", "20260720", "20260717")
|
||||
self.assertEqual(dashboard["meta"]["quote_source"], "eastmoney_clist")
|
||||
self.assertEqual(str(dashboard["meta"]["trade_date"]).replace("-", ""), "20260720")
|
||||
|
||||
def test_rt_k_and_free_source_failure_keeps_today_error(self):
|
||||
original_query = self.client.query
|
||||
|
||||
def query(api_name, params=None, fields=""):
|
||||
if api_name == "rt_k":
|
||||
raise TushareError("没有接口访问权限")
|
||||
return original_query(api_name, params, fields)
|
||||
|
||||
self.client.query = query
|
||||
self.client.realtime_aggregator = FakeFreeAggregator(fail=True)
|
||||
TushareClient._realtime_reference_cache.clear()
|
||||
with self.assertRaises(TushareError) as ctx:
|
||||
self.client._realtime_dashboard("20260720", "20260720", "20260717")
|
||||
self.assertIn("当天盘中实时行情不可用", str(ctx.exception))
|
||||
self.assertIn("没有接口访问权限", str(ctx.exception))
|
||||
|
||||
def test_rt_k_and_eastmoney_failure_falls_back_to_tencent(self):
|
||||
original_query = self.client.query
|
||||
|
||||
def query(api_name, params=None, fields=""):
|
||||
if api_name == "rt_k":
|
||||
raise TushareError("没有接口访问权限")
|
||||
return original_query(api_name, params, fields)
|
||||
|
||||
class TencentOnlyAggregator(FakeFreeAggregator):
|
||||
def eastmoney_market_quotes(self, expected_date=""):
|
||||
raise RealtimeAggregateError("eastmoney blocked")
|
||||
|
||||
def tencent_market_quotes(self, codes, expected_date=""):
|
||||
return list(FREE_QUOTES)
|
||||
|
||||
self.client.query = query
|
||||
self.client.realtime_aggregator = TencentOnlyAggregator()
|
||||
TushareClient._realtime_reference_cache.clear()
|
||||
dashboard = self.client._realtime_dashboard("20260720", "20260720", "20260717")
|
||||
self.assertEqual(dashboard["meta"]["quote_source"], "tencent_qt")
|
||||
self.assertEqual(str(dashboard["meta"]["trade_date"]).replace("-", ""), "20260720")
|
||||
self.assertIn("腾讯免费实时", dashboard["meta"]["notice"])
|
||||
self.assertEqual(dashboard["overview"]["amount_billion"], 6.0)
|
||||
|
||||
def test_normalize_eastmoney_quote_maps_units_and_exchange(self):
|
||||
quote = _normalize_eastmoney_quote(
|
||||
{
|
||||
"f12": "600000",
|
||||
"f13": 1,
|
||||
"f14": "浦发银行",
|
||||
"f2": 10.5,
|
||||
"f5": 12.0,
|
||||
"f6": 200000000,
|
||||
"f15": 10.8,
|
||||
"f16": 10.2,
|
||||
"f17": 10.3,
|
||||
"f18": 10.0,
|
||||
"f124": 1752986700,
|
||||
}
|
||||
)
|
||||
self.assertEqual(quote["ts_code"], "600000.SH")
|
||||
self.assertEqual(quote["vol"], 1200)
|
||||
self.assertEqual(quote["close"], 10.5)
|
||||
self.assertEqual(quote["pre_close"], 10.0)
|
||||
self.assertEqual(quote["source"], "eastmoney_clist")
|
||||
|
||||
def test_parse_tencent_stock_quote_keeps_today_and_units(self):
|
||||
line = (
|
||||
'v_sz000001="51~平安银行~000001~11.73~11.70~11.66~346232~0~0~0~0~0~0~0~0~0~0~0~0~0~0~0~0~0~0~0~0~0~0~'
|
||||
'~20260720100500~0.03~0.26~11.79~11.65~11.73/346232/406045563~346232~40605~0.18~5.24~~11.79~11.65~1.20~'
|
||||
'2276.29~2276.31~0.49~12.87~10.53~0.95~-3076~11.73~4.43~5.34~~~0.18~40604.5563~0.0000~0~";'
|
||||
)
|
||||
quote = _parse_tencent_stock_quote(line)
|
||||
self.assertEqual(quote["ts_code"], "000001.SZ")
|
||||
self.assertEqual(quote["quote_date"], "20260720")
|
||||
self.assertEqual(quote["close"], 11.73)
|
||||
self.assertEqual(quote["pre_close"], 11.70)
|
||||
self.assertEqual(quote["vol"], 34623200)
|
||||
self.assertEqual(quote["amount"], 406050000)
|
||||
self.assertEqual(quote["source"], "tencent_qt")
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
unittest.main()
|
||||
|
||||
@@ -5,10 +5,6 @@ import unittest
|
||||
from datetime import datetime, timedelta
|
||||
from unittest.mock import patch
|
||||
|
||||
from backend.data.providers.tushare_client import TushareError
|
||||
from backend.data.realtime import RealtimeAggregateError
|
||||
from backend.features.market.charts import ChartDataError
|
||||
from backend.features.market.service import TODAY_DAILY_UNAVAILABLE_NOTICE
|
||||
from server import DashboardService
|
||||
|
||||
|
||||
@@ -21,10 +17,6 @@ class DetailDatabaseStub:
|
||||
def list_notes(user_id, code=""):
|
||||
return []
|
||||
|
||||
@staticmethod
|
||||
def get_snapshot(trade_date):
|
||||
return {}
|
||||
|
||||
|
||||
class RealtimeClientStub:
|
||||
quote_calls = 0
|
||||
@@ -69,122 +61,6 @@ class FixedPreopenDatetime(datetime):
|
||||
return cls.fixed_now
|
||||
|
||||
|
||||
class FixedLunchDatetime(datetime):
|
||||
fixed_now = datetime(2026, 7, 31, 11, 45).astimezone()
|
||||
|
||||
@classmethod
|
||||
def now(cls, tz=None):
|
||||
return cls.fixed_now
|
||||
|
||||
|
||||
class FixedAfterCloseDatetime(datetime):
|
||||
fixed_now = datetime(2026, 7, 31, 15, 30).astimezone()
|
||||
|
||||
@classmethod
|
||||
def now(cls, tz=None):
|
||||
return cls.fixed_now
|
||||
|
||||
|
||||
class DeniedRealtimeClientStub:
|
||||
quote_calls = 0
|
||||
|
||||
def __init__(self, token):
|
||||
self.token = token
|
||||
|
||||
@staticmethod
|
||||
def resolve_trade_context(requested_date):
|
||||
return requested_date, requested_date
|
||||
|
||||
@classmethod
|
||||
def realtime_stock_quote(cls, ts_code, reference_date=""):
|
||||
cls.quote_calls += 1
|
||||
raise TushareError("没有接口访问权限")
|
||||
|
||||
|
||||
class FreeQuoteAggregator:
|
||||
def __init__(self, quote=None, fail=False):
|
||||
self.quote = quote
|
||||
self.fail = fail
|
||||
self.tencent_calls = 0
|
||||
self.eastmoney_calls = 0
|
||||
|
||||
def tencent_stock_quote(self, code, expected_date=""):
|
||||
self.tencent_calls += 1
|
||||
if self.fail:
|
||||
raise RealtimeAggregateError("tencent down")
|
||||
if self.quote and self.quote.get("source") == "eastmoney_stock":
|
||||
raise RealtimeAggregateError("tencent empty")
|
||||
if self.quote:
|
||||
return self.quote
|
||||
raise RealtimeAggregateError("tencent empty")
|
||||
|
||||
def eastmoney_stock_quote(self, code, expected_date=""):
|
||||
self.eastmoney_calls += 1
|
||||
if self.fail:
|
||||
raise RealtimeAggregateError("eastmoney down")
|
||||
if self.quote and self.quote.get("source") == "eastmoney_stock":
|
||||
return self.quote
|
||||
raise RealtimeAggregateError("eastmoney empty")
|
||||
|
||||
|
||||
class IntradayChartStub:
|
||||
def __init__(self, points, previous_close=10.0, trade_date="2026-07-31"):
|
||||
self.points = points
|
||||
self.previous_close = previous_close
|
||||
self.trade_date = trade_date
|
||||
|
||||
def stock_daily(self, code, end_date, limit=90):
|
||||
raise ChartDataError("iFinD daily unavailable")
|
||||
|
||||
def stock_intraday(self, code):
|
||||
return {
|
||||
"trade_date": self.trade_date,
|
||||
"previous_close": self.previous_close,
|
||||
"points": self.points,
|
||||
}
|
||||
|
||||
|
||||
def _history_payload(code="002141"):
|
||||
yesterday = (FixedMarketDatetime.fixed_now - timedelta(days=1)).strftime("%Y-%m-%d")
|
||||
return {
|
||||
"meta": {"trade_date": yesterday, "source": "tushare"},
|
||||
"stock": {"code": code, "name": "旧名称", "price": 10, "change": 7.1},
|
||||
"prices": [
|
||||
{
|
||||
"trade_date": yesterday,
|
||||
"open": 9.5,
|
||||
"high": 10.1,
|
||||
"low": 9.4,
|
||||
"close": 10,
|
||||
"change": 7.1,
|
||||
"volume": 100,
|
||||
"amount_billion": 1.1,
|
||||
}
|
||||
],
|
||||
"moneyflow": {},
|
||||
}
|
||||
|
||||
|
||||
def _free_quote(source="tencent_qt", **overrides):
|
||||
quote = {
|
||||
"ts_code": "002141.SZ",
|
||||
"name": "贤程科技",
|
||||
"pre_close": 10.0,
|
||||
"open": 10.2,
|
||||
"high": 10.8,
|
||||
"low": 10.1,
|
||||
"close": 10.6,
|
||||
"vol": 250000,
|
||||
"amount": 26_500_000,
|
||||
"quote_date": "20260731",
|
||||
"quote_time_epoch": int(datetime(2026, 7, 31, 10, 31).timestamp()),
|
||||
"source": source,
|
||||
"turnover_rate": 2.5,
|
||||
}
|
||||
quote.update(overrides)
|
||||
return quote
|
||||
|
||||
|
||||
class StockDetailRealtimeTests(unittest.TestCase):
|
||||
def setUp(self):
|
||||
self.service = DashboardService.__new__(DashboardService)
|
||||
@@ -192,11 +68,7 @@ class StockDetailRealtimeTests(unittest.TestCase):
|
||||
self.service.database = DetailDatabaseStub()
|
||||
self.service._request_context = threading.local()
|
||||
self.service._request_context.user_id = 1
|
||||
self.service.ifind = None
|
||||
self.service.realtime_aggregator = None
|
||||
self.service.chart_data = None
|
||||
RealtimeClientStub.quote_calls = 0
|
||||
DeniedRealtimeClientStub.quote_calls = 0
|
||||
|
||||
def test_today_detail_merges_rt_quote_without_mutating_daily_cache(self):
|
||||
today = FixedMarketDatetime.fixed_now.strftime("%Y%m%d")
|
||||
@@ -290,183 +162,6 @@ class StockDetailRealtimeTests(unittest.TestCase):
|
||||
self.assertEqual(result["stock"]["change"], 1.2)
|
||||
self.assertEqual(RealtimeClientStub.quote_calls, 0)
|
||||
|
||||
def test_today_detail_falls_back_to_tencent_quote_when_rt_k_denied(self):
|
||||
today = FixedMarketDatetime.fixed_now.strftime("%Y%m%d")
|
||||
aggregator = FreeQuoteAggregator(_free_quote())
|
||||
self.service.realtime_aggregator = aggregator
|
||||
DeniedRealtimeClientStub.quote_calls = 0
|
||||
|
||||
with patch("backend.features.market.service.datetime", FixedMarketDatetime), patch(
|
||||
"backend.features.market.service.TushareClient", DeniedRealtimeClientStub
|
||||
):
|
||||
result = self.service._prepare_stock_detail(_history_payload(), "002141", today)
|
||||
|
||||
bar = result["prices"][-1]
|
||||
self.assertEqual(bar["trade_date"], "2026-07-31")
|
||||
self.assertTrue(bar["realtime"])
|
||||
self.assertEqual(bar["open"], 10.2)
|
||||
self.assertEqual(bar["high"], 10.8)
|
||||
self.assertEqual(bar["low"], 10.1)
|
||||
self.assertEqual(bar["close"], 10.6)
|
||||
self.assertAlmostEqual(bar["change"], 6.0, places=4)
|
||||
self.assertEqual(bar["volume"], 2500)
|
||||
self.assertAlmostEqual(bar["amount_billion"], 0.265)
|
||||
self.assertEqual(len(result["prices"]), 2)
|
||||
self.assertEqual(result["meta"]["notice"], "")
|
||||
self.assertEqual(aggregator.tencent_calls, 1)
|
||||
self.assertEqual(DeniedRealtimeClientStub.quote_calls, 1)
|
||||
|
||||
def test_today_detail_falls_back_to_eastmoney_then_intraday(self):
|
||||
today = FixedMarketDatetime.fixed_now.strftime("%Y%m%d")
|
||||
aggregator = FreeQuoteAggregator(
|
||||
_free_quote("eastmoney_stock", ts_code="600000.SH", name="浦发银行"),
|
||||
)
|
||||
self.service.realtime_aggregator = aggregator
|
||||
DeniedRealtimeClientStub.quote_calls = 0
|
||||
|
||||
with patch("backend.features.market.service.datetime", FixedMarketDatetime), patch(
|
||||
"backend.features.market.service.TushareClient", DeniedRealtimeClientStub
|
||||
):
|
||||
result = self.service._prepare_stock_detail(_history_payload("600000"), "600000", today)
|
||||
|
||||
self.assertEqual(result["prices"][-1]["trade_date"], "2026-07-31")
|
||||
self.assertEqual(result["prices"][-1]["close"], 10.6)
|
||||
self.assertEqual(aggregator.tencent_calls, 1)
|
||||
self.assertEqual(aggregator.eastmoney_calls, 1)
|
||||
|
||||
aggregator = FreeQuoteAggregator(fail=True)
|
||||
self.service.realtime_aggregator = aggregator
|
||||
self.service.chart_data = IntradayChartStub(
|
||||
[
|
||||
{
|
||||
"date": "2026-07-31",
|
||||
"time": "09:30",
|
||||
"open": 10.1,
|
||||
"high": 10.2,
|
||||
"low": 10.0,
|
||||
"close": 10.15,
|
||||
"volume": 120,
|
||||
"amount": 121800,
|
||||
},
|
||||
{
|
||||
"date": "2026-07-31",
|
||||
"time": "10:05",
|
||||
"open": 10.15,
|
||||
"high": 10.5,
|
||||
"low": 9.9,
|
||||
"close": 10.4,
|
||||
"volume": 80,
|
||||
"amount": 83200,
|
||||
},
|
||||
]
|
||||
)
|
||||
with patch("backend.features.market.service.datetime", FixedMarketDatetime), patch(
|
||||
"backend.features.market.service.TushareClient", DeniedRealtimeClientStub
|
||||
):
|
||||
result = self.service._prepare_stock_detail(_history_payload("300750"), "300750", today)
|
||||
|
||||
bar = result["prices"][-1]
|
||||
self.assertEqual(bar["trade_date"], "2026-07-31")
|
||||
self.assertEqual(bar["open"], 10.1)
|
||||
self.assertEqual(bar["high"], 10.5)
|
||||
self.assertEqual(bar["low"], 9.9)
|
||||
self.assertEqual(bar["close"], 10.4)
|
||||
self.assertAlmostEqual(bar["change"], 4.0, places=4)
|
||||
self.assertEqual(bar["volume"], 200)
|
||||
self.assertTrue(bar["realtime"])
|
||||
|
||||
def test_today_detail_keeps_history_when_free_sources_fail(self):
|
||||
today = FixedMarketDatetime.fixed_now.strftime("%Y%m%d")
|
||||
self.service.realtime_aggregator = FreeQuoteAggregator(fail=True)
|
||||
self.service.chart_data = IntradayChartStub([], trade_date="2026-07-30")
|
||||
DeniedRealtimeClientStub.quote_calls = 0
|
||||
|
||||
with patch("backend.features.market.service.datetime", FixedMarketDatetime), patch(
|
||||
"backend.features.market.service.TushareClient", DeniedRealtimeClientStub
|
||||
):
|
||||
result = self.service._prepare_stock_detail(_history_payload(), "002141", today)
|
||||
|
||||
self.assertEqual(result["prices"][-1]["trade_date"], "2026-07-30")
|
||||
self.assertFalse(result["meta"].get("realtime", False))
|
||||
self.assertEqual(result["meta"]["notice"], TODAY_DAILY_UNAVAILABLE_NOTICE)
|
||||
self.assertEqual(len(result["prices"]), 1)
|
||||
|
||||
def test_lunch_keeps_morning_realtime_bar(self):
|
||||
today = FixedLunchDatetime.fixed_now.strftime("%Y%m%d")
|
||||
self.service.realtime_aggregator = FreeQuoteAggregator(
|
||||
_free_quote(quote_time_epoch=int(datetime(2026, 7, 31, 11, 30).timestamp()))
|
||||
)
|
||||
DeniedRealtimeClientStub.quote_calls = 0
|
||||
with patch("backend.features.market.service.datetime", FixedLunchDatetime), patch(
|
||||
"backend.features.market.service.TushareClient", DeniedRealtimeClientStub
|
||||
):
|
||||
result = self.service._prepare_stock_detail(_history_payload(), "002141", today)
|
||||
|
||||
self.assertEqual(result["prices"][-1]["trade_date"], "2026-07-31")
|
||||
self.assertTrue(result["meta"]["realtime"])
|
||||
|
||||
def test_after_close_keeps_forming_bar_until_official_ready(self):
|
||||
today = FixedAfterCloseDatetime.fixed_now.strftime("%Y%m%d")
|
||||
self.service.realtime_aggregator = FreeQuoteAggregator(_free_quote())
|
||||
DeniedRealtimeClientStub.quote_calls = 0
|
||||
with patch("backend.features.market.service.datetime", FixedAfterCloseDatetime), patch(
|
||||
"backend.features.market.service.TushareClient", DeniedRealtimeClientStub
|
||||
):
|
||||
forming = self.service._prepare_stock_detail(_history_payload(), "002141", today)
|
||||
self.assertEqual(forming["prices"][-1]["trade_date"], "2026-07-31")
|
||||
self.assertTrue(forming["prices"][-1]["realtime"])
|
||||
|
||||
official = _history_payload()
|
||||
official["prices"].append(
|
||||
{
|
||||
"trade_date": "2026-07-31",
|
||||
"open": 10.15,
|
||||
"high": 10.9,
|
||||
"low": 10.05,
|
||||
"close": 10.7,
|
||||
"change": 7.0,
|
||||
"volume": 1800,
|
||||
"amount_billion": 0.3,
|
||||
}
|
||||
)
|
||||
RealtimeClientStub.quote_calls = 0
|
||||
with patch("backend.features.market.service.datetime", FixedAfterCloseDatetime), patch(
|
||||
"backend.features.market.service.TushareClient", RealtimeClientStub
|
||||
):
|
||||
replaced = self.service._prepare_stock_detail(official, "002141", today)
|
||||
|
||||
self.assertEqual(replaced["prices"][-1]["close"], 10.7)
|
||||
self.assertFalse(replaced["prices"][-1].get("realtime", False))
|
||||
self.assertEqual(len(replaced["prices"]), 2)
|
||||
self.assertEqual(RealtimeClientStub.quote_calls, 0)
|
||||
|
||||
def test_same_date_bar_is_replaced_not_duplicated(self):
|
||||
today = FixedMarketDatetime.fixed_now.strftime("%Y%m%d")
|
||||
payload = _history_payload()
|
||||
payload["prices"].append(
|
||||
{
|
||||
"trade_date": "2026-07-31",
|
||||
"open": 10.0,
|
||||
"high": 10.1,
|
||||
"low": 9.9,
|
||||
"close": 10.05,
|
||||
"change": 0.5,
|
||||
"volume": 10,
|
||||
"amount_billion": 0.01,
|
||||
"realtime": True,
|
||||
}
|
||||
)
|
||||
self.service.realtime_aggregator = FreeQuoteAggregator(_free_quote())
|
||||
DeniedRealtimeClientStub.quote_calls = 0
|
||||
with patch("backend.features.market.service.datetime", FixedMarketDatetime), patch(
|
||||
"backend.features.market.service.TushareClient", DeniedRealtimeClientStub
|
||||
):
|
||||
result = self.service._prepare_stock_detail(payload, "002141", today)
|
||||
|
||||
self.assertEqual(len(result["prices"]), 2)
|
||||
self.assertEqual(result["prices"][-1]["close"], 10.6)
|
||||
self.assertEqual(result["prices"][-1]["trade_date"], "2026-07-31")
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
unittest.main()
|
||||
|
||||
@@ -221,8 +221,8 @@ def build() -> dict[str, Any]:
|
||||
{"provider": "datahub", "path": "backend/data/datahub/client.py", "runtime_role": "optional official EOD read path behind per-dataset flags"},
|
||||
{"provider": "ifind", "path": "backend/data/providers/ifind_client.py", "runtime_role": "realtime, charts, snapshots, enrichment"},
|
||||
{"provider": "eastmoney", "path": "backend/features/market/charts.py", "runtime_role": "display chart fallback"},
|
||||
{"provider": "eastmoney", "path": "backend/data/realtime.py", "runtime_role": "isolated realtime observation and intraday dashboard fallback"},
|
||||
{"provider": "tencent", "path": "backend/data/realtime.py", "runtime_role": "index observation and intraday quote fallback"},
|
||||
{"provider": "eastmoney", "path": "backend/data/realtime.py", "runtime_role": "isolated realtime observation"},
|
||||
{"provider": "tencent", "path": "backend/data/realtime.py", "runtime_role": "index observation fallback"},
|
||||
],
|
||||
"provider_domains": [
|
||||
{"provider": "tushare", "path": "backend/data/providers/tushare_transport.py", "responsibility": "HTTP transport and provider errors"},
|
||||
|
||||
@@ -6,12 +6,11 @@
|
||||
## 做什么
|
||||
|
||||
- SQLite WAL `datahub.db`,容器名 `xiaobai-datahub`,端口 `8766`
|
||||
- Tushare 盘后正式数据:交易日历、股票主档、daily、daily_basic、adj_factor、index_daily、moneyflow、stk_auction、limit_list_d、ths_hot/dc_hot、hm_detail、ths_daily/dc_index/sw_daily
|
||||
- 盘中观察(provisional):东财/腾讯指数报价、个股最新价、分时点(`/v1/quotes/latest` `/v1/indexes/quotes` `/v1/intraday/points`);永不写入 eod_* 正式表
|
||||
- Tushare 盘后正式数据:交易日历、股票主档、daily、daily_basic、adj_factor、index_daily、moneyflow、stk_auction
|
||||
- 暂存 → 校验 → 整批原子发布 → 可回滚
|
||||
- `/v1` 稳定接口(`X-Datahub-Token`)
|
||||
- `/admin/` 最小管理后台(总览 / 数据源 / 调度 / 发布 / 数据集 / 审计)
|
||||
- 同花顺/选股宝/AKShare/iFinD 适配器位仍预留;东财/腾讯已接入盘中观察
|
||||
- 东财/腾讯/同花顺/选股宝/AKShare/iFinD 适配器位已预留,本阶段不拉实时源
|
||||
|
||||
## 单位口径(相对现站)
|
||||
|
||||
@@ -124,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 份)。也可手动:
|
||||
|
||||
@@ -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 = `
|
||||
<div class="cards">
|
||||
<div class="card"><div class="muted">交易日</div><strong>${esc(data.trade_date)}</strong></div>
|
||||
<div class="card"><div class="muted">阶段</div><strong>${esc(data.session_phase)}</strong></div>
|
||||
<div class="card"><div class="muted">今日发布</div><strong>${data.publications.length}</strong></div>
|
||||
<div class="card"><div class="muted">盘后补跑</div><strong>${esc(eodLabels[eod.state] || eod.state || "-")}</strong><div class="muted">${eodExtra.join(" · ")}</div></div>
|
||||
<div class="card"><div class="muted">估值复核</div><strong>${esc(revLabels[rev.state] || rev.state || "-")}</strong><div class="muted">${revExtra.join(" · ")}</div></div>
|
||||
<div class="card"><div class="muted">异常批次</div><strong class="${data.anomalies.length ? "fail" : "ok"}">${data.anomalies.length}</strong></div>
|
||||
</div>
|
||||
<h2>最近调用</h2>
|
||||
|
||||
@@ -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",
|
||||
|
||||
@@ -1,13 +1,13 @@
|
||||
from datahub.adapters.akshare import ADAPTER as akshare
|
||||
from datahub.adapters.eastmoney import EastmoneyAdapter
|
||||
from datahub.adapters.eastmoney import ADAPTER as eastmoney
|
||||
from datahub.adapters.ifind import ADAPTER as ifind
|
||||
from datahub.adapters.tencent import TencentAdapter
|
||||
from datahub.adapters.tencent import ADAPTER as tencent
|
||||
from datahub.adapters.ths import ADAPTER as ths
|
||||
from datahub.adapters.xgb import ADAPTER as xgb
|
||||
|
||||
RESERVED = {
|
||||
"eastmoney": EastmoneyAdapter(),
|
||||
"tencent": TencentAdapter(),
|
||||
"eastmoney": eastmoney,
|
||||
"tencent": tencent,
|
||||
"ths": ths,
|
||||
"xgb": xgb,
|
||||
"akshare": akshare,
|
||||
|
||||
@@ -1,278 +1,3 @@
|
||||
from __future__ import annotations
|
||||
from datahub.adapters.base import ReservedAdapter
|
||||
|
||||
import json
|
||||
import time
|
||||
import urllib.error
|
||||
import urllib.parse
|
||||
import urllib.request
|
||||
from datetime import datetime
|
||||
from typing import Any
|
||||
|
||||
from datahub.adapters.base import AdapterError, MarketAdapter
|
||||
from datahub.numbers import finite_number, round4
|
||||
|
||||
EASTMONEY_INDEX_URL = "https://push2.eastmoney.com/api/qt/ulist.np/get"
|
||||
EASTMONEY_CLIST_URL = "https://push2.eastmoney.com/api/qt/clist/get"
|
||||
TRENDS_URL = "https://push2delay.eastmoney.com/api/qt/stock/trends2/get"
|
||||
HIS_TRENDS_URL = "https://push2his.eastmoney.com/api/qt/stock/trends2/get"
|
||||
BROWSER_UA = (
|
||||
"Mozilla/5.0 (Windows NT 10.0; Win64; x64) "
|
||||
"AppleWebKit/537.36 (KHTML, like Gecko) Chrome/138.0.0.0 Safari/537.36"
|
||||
)
|
||||
INDEX_SECIDS = {
|
||||
"000001.SH": "1.000001",
|
||||
"399001.SZ": "0.399001",
|
||||
"399006.SZ": "0.399006",
|
||||
}
|
||||
|
||||
|
||||
class EastmoneyAdapter(MarketAdapter):
|
||||
name = "eastmoney"
|
||||
|
||||
def __init__(self, timeout: int = 8) -> None:
|
||||
self.timeout = timeout
|
||||
|
||||
def probe(self) -> dict[str, Any]:
|
||||
started = time.perf_counter()
|
||||
try:
|
||||
rows = self.fetch_indices()
|
||||
state = "ok" if len(rows) == 3 else "empty"
|
||||
except AdapterError as exc:
|
||||
return {
|
||||
"provider": self.name,
|
||||
"configured": True,
|
||||
"state": "error",
|
||||
"message": str(exc),
|
||||
"latency_ms": round((time.perf_counter() - started) * 1000),
|
||||
}
|
||||
return {
|
||||
"provider": self.name,
|
||||
"configured": True,
|
||||
"state": state,
|
||||
"latency_ms": round((time.perf_counter() - started) * 1000),
|
||||
}
|
||||
|
||||
def fetch(self, dataset: str, params: dict[str, Any]) -> list[dict[str, Any]]:
|
||||
if dataset in {"indexes_quotes", "index_quotes"}:
|
||||
return self.fetch_indices()
|
||||
if dataset in {"quotes", "quotes_latest"}:
|
||||
codes = params.get("codes") or []
|
||||
if isinstance(codes, str):
|
||||
codes = [item.strip() for item in codes.split(",") if item.strip()]
|
||||
return self.fetch_quotes(list(codes))
|
||||
raise AdapterError(f"{self.name} unsupported dataset: {dataset}")
|
||||
|
||||
def normalize(self, dataset: str, rows: list[dict[str, Any]]) -> list[dict[str, Any]]:
|
||||
return list(rows)
|
||||
|
||||
def fetch_indices(self) -> list[dict[str, Any]]:
|
||||
payload = self._get_json(
|
||||
EASTMONEY_INDEX_URL,
|
||||
{
|
||||
"secids": "1.000001,0.399001,0.399006",
|
||||
"fltt": "2",
|
||||
"invt": "2",
|
||||
"fields": "f12,f14,f2,f3,f4,f15,f16,f17,f18,f6,f124",
|
||||
},
|
||||
referer="https://quote.eastmoney.com/",
|
||||
)
|
||||
rows = list((payload.get("data") or {}).get("diff") or [])
|
||||
result = []
|
||||
for row in rows:
|
||||
code = str(row.get("f12") or "")
|
||||
if code not in {"000001", "399001", "399006"}:
|
||||
continue
|
||||
epoch = int(finite_number(row.get("f124")) or 0)
|
||||
ts_code = f"{code}.SH" if code.startswith("0") and code == "000001" else f"{code}.SZ"
|
||||
if code == "000001":
|
||||
ts_code = "000001.SH"
|
||||
result.append(
|
||||
{
|
||||
"ts_code": ts_code,
|
||||
"code": code,
|
||||
"name": row.get("f14") or code,
|
||||
"price": round4(finite_number(row.get("f2"))),
|
||||
"pct_chg": round4(finite_number(row.get("f3"))),
|
||||
"change_amount": round4(finite_number(row.get("f4"))),
|
||||
"open": round4(finite_number(row.get("f17"))),
|
||||
"high": round4(finite_number(row.get("f15"))),
|
||||
"low": round4(finite_number(row.get("f16"))),
|
||||
"previous_close": round4(finite_number(row.get("f18"))),
|
||||
"amount": round4(finite_number(row.get("f6"))),
|
||||
"quote_time_epoch": epoch,
|
||||
"quote_time": (
|
||||
datetime.fromtimestamp(epoch).astimezone().isoformat(timespec="seconds")
|
||||
if epoch
|
||||
else ""
|
||||
),
|
||||
"source": "eastmoney_push2",
|
||||
}
|
||||
)
|
||||
if len(result) != 3:
|
||||
raise AdapterError(f"Eastmoney returned {len(result)}/3 indices")
|
||||
return result
|
||||
|
||||
def fetch_quotes(self, codes: list[str]) -> list[dict[str, Any]]:
|
||||
# Eastmoney clist does not accept arbitrary code lists well; use ulist.np for batches.
|
||||
secids = []
|
||||
for code in codes:
|
||||
ts = str(code or "").upper()
|
||||
symbol = ts.split(".")[0]
|
||||
if ts.endswith(".SH") or symbol.startswith(("5", "6", "9")):
|
||||
secids.append(f"1.{symbol}")
|
||||
else:
|
||||
secids.append(f"0.{symbol}")
|
||||
if not secids:
|
||||
return []
|
||||
payload = self._get_json(
|
||||
EASTMONEY_INDEX_URL,
|
||||
{
|
||||
"secids": ",".join(secids[:60]),
|
||||
"fltt": "2",
|
||||
"invt": "2",
|
||||
"fields": "f12,f14,f2,f3,f4,f15,f16,f17,f18,f5,f6,f8,f124",
|
||||
},
|
||||
referer="https://quote.eastmoney.com/",
|
||||
)
|
||||
rows = list((payload.get("data") or {}).get("diff") or [])
|
||||
result = []
|
||||
for row in rows:
|
||||
symbol = str(row.get("f12") or "")
|
||||
if not symbol:
|
||||
continue
|
||||
ts_code = f"{symbol}.SH" if symbol.startswith(("5", "6", "9")) else f"{symbol}.SZ"
|
||||
epoch = int(finite_number(row.get("f124")) or 0)
|
||||
result.append(
|
||||
{
|
||||
"ts_code": ts_code,
|
||||
"name": row.get("f14") or symbol,
|
||||
"price": round4(finite_number(row.get("f2"))),
|
||||
"pct_chg": round4(finite_number(row.get("f3"))),
|
||||
"change_amount": round4(finite_number(row.get("f4"))),
|
||||
"open": round4(finite_number(row.get("f17"))),
|
||||
"high": round4(finite_number(row.get("f15"))),
|
||||
"low": round4(finite_number(row.get("f16"))),
|
||||
"previous_close": round4(finite_number(row.get("f18"))),
|
||||
"volume": round4(finite_number(row.get("f5"))),
|
||||
"amount": round4(finite_number(row.get("f6"))),
|
||||
"turnover_rate": round4(finite_number(row.get("f8"))),
|
||||
"quote_time_epoch": epoch,
|
||||
"quote_time": (
|
||||
datetime.fromtimestamp(epoch).astimezone().isoformat(timespec="seconds")
|
||||
if epoch
|
||||
else ""
|
||||
),
|
||||
"source": "eastmoney_push2",
|
||||
}
|
||||
)
|
||||
return result
|
||||
|
||||
def fetch_intraday(self, ts_code: str, date: str = "") -> dict[str, Any]:
|
||||
code = str(ts_code or "").upper()
|
||||
if code in INDEX_SECIDS:
|
||||
secid = INDEX_SECIDS[code]
|
||||
entity = "index"
|
||||
identifier = code
|
||||
else:
|
||||
symbol = code.split(".")[0]
|
||||
market = "1" if symbol.startswith(("5", "6", "9")) else "0"
|
||||
secid = f"{market}.{symbol}"
|
||||
entity = "stock"
|
||||
identifier = symbol
|
||||
params = {
|
||||
"secid": secid,
|
||||
"fields1": "f1,f2,f3,f4,f5,f6,f7,f8,f9,f10,f11,f12,f13",
|
||||
"fields2": "f51,f52,f53,f54,f55,f56,f57,f58",
|
||||
"iscr": "0",
|
||||
}
|
||||
data: dict[str, Any] = {}
|
||||
points: list[dict[str, Any]] = []
|
||||
last_error: Exception | None = None
|
||||
for url, ndays in ((TRENDS_URL, "1"), (TRENDS_URL, "5"), (HIS_TRENDS_URL, "5")):
|
||||
try:
|
||||
payload = self._get_json(
|
||||
url,
|
||||
{**params, "ndays": ndays},
|
||||
referer="https://quote.eastmoney.com/",
|
||||
)
|
||||
except AdapterError as exc:
|
||||
last_error = exc
|
||||
continue
|
||||
data = payload.get("data") or {}
|
||||
parsed = [point for raw in data.get("trends") or [] if (point := _parse_trend(raw))]
|
||||
points = _preferred_session(parsed, date)
|
||||
if points:
|
||||
break
|
||||
if not points:
|
||||
raise AdapterError("No intraday chart data returned") from last_error
|
||||
return {
|
||||
"entity_type": entity,
|
||||
"identifier": identifier,
|
||||
"ts_code": code if "." in code else f"{identifier}.{'SH' if identifier.startswith(('5','6','9')) else 'SZ'}",
|
||||
"name": str(data.get("name") or ""),
|
||||
"code": str(data.get("code") or identifier),
|
||||
"trade_date": points[-1]["date"],
|
||||
"previous_close": round4(finite_number(data.get("preClose"))),
|
||||
"points": points,
|
||||
"source": "eastmoney_trends2",
|
||||
}
|
||||
|
||||
def _get_json(self, url: str, params: dict[str, str], referer: str) -> dict[str, Any]:
|
||||
request_url = f"{url}?{urllib.parse.urlencode(params)}"
|
||||
request = urllib.request.Request(
|
||||
request_url,
|
||||
headers={
|
||||
"Accept": "application/json,text/plain,*/*",
|
||||
"User-Agent": BROWSER_UA,
|
||||
"Referer": referer,
|
||||
},
|
||||
method="GET",
|
||||
)
|
||||
try:
|
||||
with urllib.request.urlopen(request, timeout=self.timeout) as response:
|
||||
return json.loads(response.read().decode("utf-8"))
|
||||
except Exception as exc:
|
||||
raise AdapterError(f"eastmoney request failed: {exc}") from exc
|
||||
|
||||
|
||||
def _preferred_session(points: list[dict[str, Any]], preferred_date: str = "") -> list[dict[str, Any]]:
|
||||
if not points:
|
||||
return []
|
||||
want = ""
|
||||
digits = str(preferred_date or "").replace("-", "")[:8]
|
||||
if len(digits) == 8 and digits.isdigit():
|
||||
want = f"{digits[:4]}-{digits[4:6]}-{digits[6:8]}"
|
||||
if want:
|
||||
matched = [point for point in points if str(point.get("date") or "") == want]
|
||||
if matched:
|
||||
return matched
|
||||
latest = max(str(point.get("date") or "") for point in points)
|
||||
if not latest:
|
||||
return points
|
||||
return [point for point in points if str(point.get("date") or "") == latest]
|
||||
|
||||
|
||||
def _parse_trend(raw: Any) -> dict[str, Any] | None:
|
||||
text = str(raw or "")
|
||||
parts = text.split(",")
|
||||
if len(parts) < 8:
|
||||
return None
|
||||
stamp = parts[0]
|
||||
try:
|
||||
when = datetime.strptime(stamp, "%Y-%m-%d %H:%M")
|
||||
except ValueError:
|
||||
return None
|
||||
close = round4(finite_number(parts[2]))
|
||||
if close <= 0:
|
||||
return None
|
||||
return {
|
||||
"time": when.strftime("%H:%M"),
|
||||
"date": when.strftime("%Y-%m-%d"),
|
||||
"open": round4(finite_number(parts[1])),
|
||||
"close": close,
|
||||
"high": round4(finite_number(parts[3])),
|
||||
"low": round4(finite_number(parts[4])),
|
||||
"avg_price": round4(finite_number(parts[7] if len(parts) > 7 else parts[2])),
|
||||
"volume": round4(finite_number(parts[5])),
|
||||
"amount": round4(finite_number(parts[6])),
|
||||
}
|
||||
ADAPTER = ReservedAdapter("eastmoney")
|
||||
|
||||
@@ -1,99 +1,3 @@
|
||||
from __future__ import annotations
|
||||
from datahub.adapters.base import ReservedAdapter
|
||||
|
||||
import time
|
||||
import urllib.error
|
||||
import urllib.request
|
||||
from datetime import datetime
|
||||
from typing import Any
|
||||
|
||||
from datahub.adapters.base import AdapterError, MarketAdapter
|
||||
from datahub.numbers import finite_number, round4
|
||||
|
||||
TENCENT_INDEX_URL = "https://qt.gtimg.cn/q=sh000001,sz399001,sz399006"
|
||||
BROWSER_UA = (
|
||||
"Mozilla/5.0 (Windows NT 10.0; Win64; x64) "
|
||||
"AppleWebKit/537.36 (KHTML, like Gecko) Chrome/138.0.0.0 Safari/537.36"
|
||||
)
|
||||
|
||||
|
||||
class TencentAdapter(MarketAdapter):
|
||||
name = "tencent"
|
||||
|
||||
def __init__(self, timeout: int = 8) -> None:
|
||||
self.timeout = timeout
|
||||
|
||||
def probe(self) -> dict[str, Any]:
|
||||
started = time.perf_counter()
|
||||
try:
|
||||
rows = self.fetch_indices()
|
||||
state = "ok" if len(rows) == 3 else "empty"
|
||||
except AdapterError as exc:
|
||||
return {
|
||||
"provider": self.name,
|
||||
"configured": True,
|
||||
"state": "error",
|
||||
"message": str(exc),
|
||||
"latency_ms": round((time.perf_counter() - started) * 1000),
|
||||
}
|
||||
return {
|
||||
"provider": self.name,
|
||||
"configured": True,
|
||||
"state": state,
|
||||
"latency_ms": round((time.perf_counter() - started) * 1000),
|
||||
}
|
||||
|
||||
def fetch(self, dataset: str, params: dict[str, Any]) -> list[dict[str, Any]]:
|
||||
if dataset in {"indexes_quotes", "index_quotes"}:
|
||||
return self.fetch_indices()
|
||||
raise AdapterError(f"{self.name} unsupported dataset: {dataset}")
|
||||
|
||||
def normalize(self, dataset: str, rows: list[dict[str, Any]]) -> list[dict[str, Any]]:
|
||||
return list(rows)
|
||||
|
||||
def fetch_indices(self) -> list[dict[str, Any]]:
|
||||
request = urllib.request.Request(
|
||||
TENCENT_INDEX_URL,
|
||||
headers={"User-Agent": BROWSER_UA, "Referer": "https://gu.qq.com/"},
|
||||
method="GET",
|
||||
)
|
||||
try:
|
||||
with urllib.request.urlopen(request, timeout=self.timeout) as response:
|
||||
raw = response.read().decode("gb18030", errors="ignore")
|
||||
except Exception as exc:
|
||||
raise AdapterError(f"tencent request failed: {exc}") from exc
|
||||
result = []
|
||||
for line in raw.splitlines():
|
||||
if '="' not in line:
|
||||
continue
|
||||
fields = line.split('="', 1)[1].rsplit('";', 1)[0].split("~")
|
||||
if len(fields) < 38:
|
||||
continue
|
||||
code = fields[2]
|
||||
if code not in {"000001", "399001", "399006"}:
|
||||
continue
|
||||
try:
|
||||
quote_time = datetime.strptime(fields[30], "%Y%m%d%H%M%S").astimezone()
|
||||
except ValueError as exc:
|
||||
raise AdapterError(f"Tencent invalid quote time for {code}") from exc
|
||||
ts_code = "000001.SH" if code == "000001" else f"{code}.SZ"
|
||||
result.append(
|
||||
{
|
||||
"ts_code": ts_code,
|
||||
"code": code,
|
||||
"name": fields[1] or code,
|
||||
"price": round4(finite_number(fields[3])),
|
||||
"pct_chg": round4(finite_number(fields[32])),
|
||||
"change_amount": round4(finite_number(fields[31])),
|
||||
"open": round4(finite_number(fields[5])),
|
||||
"high": round4(finite_number(fields[33])),
|
||||
"low": round4(finite_number(fields[34])),
|
||||
"previous_close": round4(finite_number(fields[4])),
|
||||
"amount": round4(finite_number(fields[37]) * 10000),
|
||||
"quote_time_epoch": int(quote_time.timestamp()),
|
||||
"quote_time": quote_time.isoformat(timespec="seconds"),
|
||||
"source": "tencent_qt",
|
||||
}
|
||||
)
|
||||
if len(result) != 3:
|
||||
raise AdapterError(f"Tencent returned {len(result)}/3 indices")
|
||||
return result
|
||||
ADAPTER = ReservedAdapter("tencent")
|
||||
|
||||
@@ -11,12 +11,8 @@ from datahub.normalize import (
|
||||
normalize_auction,
|
||||
normalize_calendar,
|
||||
normalize_daily,
|
||||
normalize_dragon_tiger,
|
||||
normalize_index_daily,
|
||||
normalize_limit_event,
|
||||
normalize_moneyflow,
|
||||
normalize_popularity,
|
||||
normalize_sector_daily,
|
||||
normalize_stock,
|
||||
normalize_valuation,
|
||||
)
|
||||
@@ -35,21 +31,6 @@ TUSHARE_FIELDS = {
|
||||
"buy_lg_amount,sell_lg_amount,buy_elg_amount,sell_elg_amount,net_mf_amount"
|
||||
),
|
||||
"stk_auction": "ts_code,trade_date,vol,price,amount,pre_close,turnover_rate,volume_ratio,float_share",
|
||||
"limit_list_d": (
|
||||
"trade_date,ts_code,industry,name,close,pct_chg,amount,limit_amount,"
|
||||
"float_mv,total_mv,turnover_ratio,fd_amount,first_time,last_time,"
|
||||
"open_times,up_stat,limit_times,limit_type"
|
||||
),
|
||||
"ths_hot": "ts_code,ts_name,hot,rank,pct_change,current_price,concept,data_type,trade_date",
|
||||
"dc_hot": "ts_code,ts_name,rank,pct_change,current_price,hot,concept,data_type,trade_date",
|
||||
"hm_detail": "trade_date,ts_code,ts_name,buy_amount,sell_amount,net_amount,hm_name,hm_orgs,tag",
|
||||
"hm_list": "name,desc,orgs",
|
||||
"top_list": "trade_date,ts_code,name,pct_change,reason",
|
||||
"top_inst": "trade_date,ts_code,exalter,buy,buy_rate,sell,sell_rate,net_buy,side,reason",
|
||||
"ths_index": "ts_code,name,count,exchange,list_date,type",
|
||||
"ths_daily": "ts_code,trade_date,open,high,low,close,pre_close,pct_change,vol,turnover_rate",
|
||||
"dc_index": "ts_code,trade_date,name,open,high,low,close,pre_close,pct_change,vol,amount,turnover_rate",
|
||||
"sw_daily": "ts_code,trade_date,name,open,high,low,close,pct_change,vol,amount",
|
||||
}
|
||||
|
||||
DATASET_API = {
|
||||
@@ -61,15 +42,12 @@ DATASET_API = {
|
||||
"index_daily": "index_daily",
|
||||
"moneyflow": "moneyflow",
|
||||
"auction": "stk_auction",
|
||||
"limit_events": "limit_list_d",
|
||||
"popularity": "ths_hot",
|
||||
"dragon_tiger": "hm_detail",
|
||||
"sector_daily": "ths_daily",
|
||||
}
|
||||
|
||||
# Website actual index usage: market cards / 90-day charts (SH/SZ/CYB) plus
|
||||
# screener 沪深300 benchmark (lookback up to 260 trading days).
|
||||
WEBSITE_INDEX_CODES = ("000001.SH", "399001.SZ", "399006.SZ", "000300.SH")
|
||||
DEFAULT_INDEX_CODES = WEBSITE_INDEX_CODES
|
||||
LIMIT_TYPES = ("U", "D", "Z")
|
||||
|
||||
|
||||
class TushareAdapter(MarketAdapter):
|
||||
@@ -107,14 +85,6 @@ class TushareAdapter(MarketAdapter):
|
||||
}
|
||||
|
||||
def fetch(self, dataset: str, params: dict[str, Any]) -> list[dict[str, Any]]:
|
||||
if dataset == "limit_events":
|
||||
return self.fetch_limit_events(str(params.get("trade_date") or ""))
|
||||
if dataset == "popularity":
|
||||
return self.fetch_popularity(str(params.get("trade_date") or ""))
|
||||
if dataset == "dragon_tiger":
|
||||
return self.fetch_dragon_tiger(str(params.get("trade_date") or ""))
|
||||
if dataset == "sector_daily":
|
||||
return self.fetch_sector_daily(str(params.get("trade_date") or ""))
|
||||
api_name = DATASET_API.get(dataset, dataset)
|
||||
fields = TUSHARE_FIELDS.get(api_name, "")
|
||||
query_params = dict(params)
|
||||
@@ -123,67 +93,10 @@ class TushareAdapter(MarketAdapter):
|
||||
if api_name == "trade_cal" and "exchange" not in query_params:
|
||||
query_params["exchange"] = "SSE"
|
||||
if api_name == "index_daily" and "ts_code" not in query_params:
|
||||
# Caller typically loops codes; a missing code would pull nothing useful.
|
||||
query_params.setdefault("ts_code", DEFAULT_INDEX_CODES[0])
|
||||
return self._query(api_name, query_params, fields)
|
||||
|
||||
def fetch_limit_events(self, trade_date: str) -> list[dict[str, Any]]:
|
||||
rows: list[dict[str, Any]] = []
|
||||
for limit_type in LIMIT_TYPES:
|
||||
part = self._query(
|
||||
"limit_list_d",
|
||||
{"trade_date": trade_date, "limit_type": limit_type},
|
||||
TUSHARE_FIELDS["limit_list_d"],
|
||||
)
|
||||
for row in part:
|
||||
row = dict(row)
|
||||
row.setdefault("limit_type", limit_type)
|
||||
rows.append(row)
|
||||
return rows
|
||||
|
||||
def fetch_popularity(self, trade_date: str) -> list[dict[str, Any]]:
|
||||
rows: list[dict[str, Any]] = []
|
||||
for api_name, source in (("ths_hot", "ths"), ("dc_hot", "dc")):
|
||||
for row in self._query(api_name, {"trade_date": trade_date}, TUSHARE_FIELDS[api_name]):
|
||||
item = dict(row)
|
||||
item["source"] = source
|
||||
item.setdefault("trade_date", trade_date)
|
||||
rows.append(item)
|
||||
return rows
|
||||
|
||||
def fetch_dragon_tiger(self, trade_date: str) -> list[dict[str, Any]]:
|
||||
details = self._query("hm_detail", {"trade_date": trade_date}, TUSHARE_FIELDS["hm_detail"])
|
||||
top_rows = self._query("top_list", {"trade_date": trade_date}, TUSHARE_FIELDS["top_list"])
|
||||
context = {
|
||||
str(row.get("ts_code") or ""): row
|
||||
for row in top_rows
|
||||
if str(row.get("ts_code") or "")
|
||||
}
|
||||
rows: list[dict[str, Any]] = []
|
||||
for row in details:
|
||||
item = dict(row)
|
||||
stock = context.get(str(item.get("ts_code") or ""), {})
|
||||
if item.get("pct_change") is None and stock.get("pct_change") is not None:
|
||||
item["pct_change"] = stock.get("pct_change")
|
||||
if not item.get("reason") and stock.get("reason"):
|
||||
item["reason"] = stock.get("reason")
|
||||
if not item.get("ts_name") and stock.get("name"):
|
||||
item["ts_name"] = stock.get("name")
|
||||
rows.append(item)
|
||||
return rows
|
||||
|
||||
def fetch_sector_daily(self, trade_date: str) -> list[dict[str, Any]]:
|
||||
rows: list[dict[str, Any]] = []
|
||||
for api_name, family in (("ths_daily", "ths"), ("dc_index", "dc"), ("sw_daily", "sw")):
|
||||
try:
|
||||
part = self._query(api_name, {"trade_date": trade_date}, TUSHARE_FIELDS[api_name])
|
||||
except AdapterError:
|
||||
part = []
|
||||
for row in part:
|
||||
item = dict(row)
|
||||
item["family"] = family
|
||||
rows.append(item)
|
||||
return rows
|
||||
|
||||
def fetch_index_daily(self, trade_date: str, codes: tuple[str, ...] = DEFAULT_INDEX_CODES) -> list[dict[str, Any]]:
|
||||
rows: list[dict[str, Any]] = []
|
||||
for ts_code in codes:
|
||||
@@ -191,17 +104,6 @@ class TushareAdapter(MarketAdapter):
|
||||
return rows
|
||||
|
||||
def normalize(self, dataset: str, rows: list[dict[str, Any]]) -> list[dict[str, Any]]:
|
||||
if dataset in {"limit_events", "limit_list_d"}:
|
||||
return [normalize_limit_event(row) for row in rows]
|
||||
if dataset == "popularity":
|
||||
return [normalize_popularity(row, source=str(row.get("source") or "")) for row in rows]
|
||||
if dataset == "dragon_tiger":
|
||||
return [normalize_dragon_tiger(row) for row in rows]
|
||||
if dataset == "sector_daily":
|
||||
return [
|
||||
normalize_sector_daily(row, family=str(row.get("family") or "ths"))
|
||||
for row in rows
|
||||
]
|
||||
mapping = {
|
||||
"calendar": normalize_calendar,
|
||||
"trade_cal": normalize_calendar,
|
||||
@@ -246,11 +148,12 @@ class TushareAdapter(MarketAdapter):
|
||||
try:
|
||||
with urllib.request.urlopen(request, timeout=self.timeout) as response:
|
||||
result = json.loads(response.read().decode("utf-8"))
|
||||
except (urllib.error.URLError, TimeoutError, json.JSONDecodeError) as exc:
|
||||
raise AdapterError(f"Tushare 请求失败: {exc}") from exc
|
||||
if result.get("code") not in (0, "0", None):
|
||||
raise AdapterError(str(result.get("msg") or f"Tushare error {result.get('code')}"))
|
||||
except json.JSONDecodeError:
|
||||
raise AdapterError("Tushare returned invalid json") from None
|
||||
except (urllib.error.URLError, TimeoutError) as exc:
|
||||
raise AdapterError(f"Tushare request failed: {exc}") from exc
|
||||
if result.get("code") != 0:
|
||||
raise AdapterError(result.get("msg") or "Tushare returned an unknown error")
|
||||
data = result.get("data") or {}
|
||||
items = data.get("items") or []
|
||||
fields_list = data.get("fields") or (fields.split(",") if fields else [])
|
||||
return [dict(zip(fields_list, item)) for item in items]
|
||||
columns = data.get("fields") or []
|
||||
return [dict(zip(columns, item)) for item in data.get("items") or []]
|
||||
|
||||
@@ -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 / 日志"},
|
||||
|
||||
@@ -1,180 +0,0 @@
|
||||
"""Extended EOD datasets beyond the first-batch A/B release groups.
|
||||
|
||||
These publish independently (soft): a failure here must not block daily/valuation
|
||||
release. Scheduler runs them after the core EOD window.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
from typing import Any
|
||||
|
||||
# Independent soft datasets (not part of A/B atomic groups).
|
||||
EXTENDED_SOFT_DATASETS = {
|
||||
"limit_events",
|
||||
"popularity",
|
||||
"dragon_tiger",
|
||||
"sector_daily",
|
||||
}
|
||||
|
||||
EXTENDED_SCHEMA = """
|
||||
CREATE TABLE IF NOT EXISTS eod_limit_events (
|
||||
ts_code TEXT NOT NULL, trade_date TEXT NOT NULL, limit_type TEXT NOT NULL,
|
||||
name TEXT, industry TEXT, close REAL, pct_chg REAL, amount REAL,
|
||||
limit_amount REAL, float_mv REAL, total_mv REAL, turnover_ratio REAL,
|
||||
fd_amount REAL, first_time TEXT, last_time TEXT,
|
||||
open_times INTEGER, up_stat TEXT, limit_times INTEGER,
|
||||
batch_id TEXT NOT NULL,
|
||||
PRIMARY KEY (ts_code, trade_date, limit_type, batch_id)
|
||||
) WITHOUT ROWID;
|
||||
|
||||
CREATE TABLE IF NOT EXISTS staging_limit_events (
|
||||
ts_code TEXT NOT NULL, trade_date TEXT NOT NULL, limit_type TEXT NOT NULL, batch_id TEXT NOT NULL,
|
||||
name TEXT, industry TEXT, close REAL, pct_chg REAL, amount REAL,
|
||||
limit_amount REAL, float_mv REAL, total_mv REAL, turnover_ratio REAL,
|
||||
fd_amount REAL, first_time TEXT, last_time TEXT,
|
||||
open_times INTEGER, up_stat TEXT, limit_times INTEGER,
|
||||
PRIMARY KEY (batch_id, ts_code, trade_date, limit_type)
|
||||
);
|
||||
|
||||
CREATE TABLE IF NOT EXISTS eod_popularity (
|
||||
ts_code TEXT NOT NULL, trade_date TEXT NOT NULL, source TEXT NOT NULL,
|
||||
ts_name TEXT, rank INTEGER, pct_change REAL, current_price REAL,
|
||||
hot REAL, concept TEXT, data_type TEXT,
|
||||
batch_id TEXT NOT NULL,
|
||||
PRIMARY KEY (ts_code, trade_date, source, batch_id)
|
||||
) WITHOUT ROWID;
|
||||
|
||||
CREATE TABLE IF NOT EXISTS staging_popularity (
|
||||
ts_code TEXT NOT NULL, trade_date TEXT NOT NULL, source TEXT NOT NULL, batch_id TEXT NOT NULL,
|
||||
ts_name TEXT, rank INTEGER, pct_change REAL, current_price REAL,
|
||||
hot REAL, concept TEXT, data_type TEXT,
|
||||
PRIMARY KEY (batch_id, ts_code, trade_date, source)
|
||||
);
|
||||
|
||||
CREATE TABLE IF NOT EXISTS eod_dragon_tiger (
|
||||
ts_code TEXT NOT NULL, trade_date TEXT NOT NULL, hm_name TEXT NOT NULL,
|
||||
ts_name TEXT, buy_amount REAL, sell_amount REAL, net_amount REAL,
|
||||
hm_orgs TEXT, tag TEXT, pct_change REAL, reason TEXT,
|
||||
batch_id TEXT NOT NULL,
|
||||
PRIMARY KEY (ts_code, trade_date, hm_name, batch_id)
|
||||
) WITHOUT ROWID;
|
||||
|
||||
CREATE TABLE IF NOT EXISTS staging_dragon_tiger (
|
||||
ts_code TEXT NOT NULL, trade_date TEXT NOT NULL, hm_name TEXT NOT NULL, batch_id TEXT NOT NULL,
|
||||
ts_name TEXT, buy_amount REAL, sell_amount REAL, net_amount REAL,
|
||||
hm_orgs TEXT, tag TEXT, pct_change REAL, reason TEXT,
|
||||
PRIMARY KEY (batch_id, ts_code, trade_date, hm_name)
|
||||
);
|
||||
|
||||
CREATE TABLE IF NOT EXISTS eod_sector_daily (
|
||||
ts_code TEXT NOT NULL, trade_date TEXT NOT NULL, family TEXT NOT NULL,
|
||||
name TEXT, open REAL, high REAL, low REAL, close REAL, pre_close REAL,
|
||||
pct_change REAL, vol REAL, turnover_rate REAL, amount REAL,
|
||||
batch_id TEXT NOT NULL,
|
||||
PRIMARY KEY (ts_code, trade_date, family, batch_id)
|
||||
) WITHOUT ROWID;
|
||||
|
||||
CREATE TABLE IF NOT EXISTS staging_sector_daily (
|
||||
ts_code TEXT NOT NULL, trade_date TEXT NOT NULL, family TEXT NOT NULL, batch_id TEXT NOT NULL,
|
||||
name TEXT, open REAL, high REAL, low REAL, close REAL, pre_close REAL,
|
||||
pct_change REAL, vol REAL, turnover_rate REAL, amount REAL,
|
||||
PRIMARY KEY (batch_id, ts_code, trade_date, family)
|
||||
);
|
||||
|
||||
CREATE TABLE IF NOT EXISTS sector_master (
|
||||
ts_code TEXT PRIMARY KEY,
|
||||
name TEXT,
|
||||
family TEXT NOT NULL,
|
||||
exchange TEXT,
|
||||
list_date TEXT,
|
||||
member_count INTEGER,
|
||||
type TEXT,
|
||||
updated_at TEXT NOT NULL
|
||||
);
|
||||
|
||||
CREATE INDEX IF NOT EXISTS idx_eod_limit_date ON eod_limit_events(trade_date, batch_id);
|
||||
CREATE INDEX IF NOT EXISTS idx_eod_pop_date ON eod_popularity(trade_date, batch_id);
|
||||
CREATE INDEX IF NOT EXISTS idx_eod_lhb_date ON eod_dragon_tiger(trade_date, batch_id);
|
||||
CREATE INDEX IF NOT EXISTS idx_eod_sector_date ON eod_sector_daily(trade_date, family, batch_id);
|
||||
"""
|
||||
|
||||
EXTENDED_DATASET_TABLES = {
|
||||
"limit_events": ("eod_limit_events", "staging_limit_events"),
|
||||
"popularity": ("eod_popularity", "staging_popularity"),
|
||||
"dragon_tiger": ("eod_dragon_tiger", "staging_dragon_tiger"),
|
||||
"sector_daily": ("eod_sector_daily", "staging_sector_daily"),
|
||||
}
|
||||
|
||||
EXTENDED_STAGING_INSERT: dict[str, tuple[str, Any]] = {
|
||||
"limit_events": (
|
||||
"INSERT INTO staging_limit_events("
|
||||
"ts_code,trade_date,limit_type,batch_id,name,industry,close,pct_chg,amount,"
|
||||
"limit_amount,float_mv,total_mv,turnover_ratio,fd_amount,first_time,last_time,"
|
||||
"open_times,up_stat,limit_times) VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?)",
|
||||
lambda r, b: (
|
||||
r["ts_code"], r["trade_date"], r["limit_type"], b,
|
||||
r.get("name"), r.get("industry"), r.get("close"), r.get("pct_chg"), r.get("amount"),
|
||||
r.get("limit_amount"), r.get("float_mv"), r.get("total_mv"), r.get("turnover_ratio"),
|
||||
r.get("fd_amount"), r.get("first_time"), r.get("last_time"),
|
||||
r.get("open_times"), r.get("up_stat"), r.get("limit_times"),
|
||||
),
|
||||
),
|
||||
"popularity": (
|
||||
"INSERT INTO staging_popularity("
|
||||
"ts_code,trade_date,source,batch_id,ts_name,rank,pct_change,current_price,hot,concept,data_type) "
|
||||
"VALUES (?,?,?,?,?,?,?,?,?,?,?)",
|
||||
lambda r, b: (
|
||||
r["ts_code"], r["trade_date"], r["source"], b,
|
||||
r.get("ts_name"), r.get("rank"), r.get("pct_change"), r.get("current_price"),
|
||||
r.get("hot"), r.get("concept"), r.get("data_type"),
|
||||
),
|
||||
),
|
||||
"dragon_tiger": (
|
||||
"INSERT INTO staging_dragon_tiger("
|
||||
"ts_code,trade_date,hm_name,batch_id,ts_name,buy_amount,sell_amount,net_amount,"
|
||||
"hm_orgs,tag,pct_change,reason) VALUES (?,?,?,?,?,?,?,?,?,?,?,?)",
|
||||
lambda r, b: (
|
||||
r["ts_code"], r["trade_date"], r["hm_name"], b,
|
||||
r.get("ts_name"), r.get("buy_amount"), r.get("sell_amount"), r.get("net_amount"),
|
||||
r.get("hm_orgs"), r.get("tag"), r.get("pct_change"), r.get("reason"),
|
||||
),
|
||||
),
|
||||
"sector_daily": (
|
||||
"INSERT INTO staging_sector_daily("
|
||||
"ts_code,trade_date,family,batch_id,name,open,high,low,close,pre_close,"
|
||||
"pct_change,vol,turnover_rate,amount) VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,?)",
|
||||
lambda r, b: (
|
||||
r["ts_code"], r["trade_date"], r["family"], b,
|
||||
r.get("name"), r.get("open"), r.get("high"), r.get("low"), r.get("close"),
|
||||
r.get("pre_close"), r.get("pct_change"), r.get("vol"), r.get("turnover_rate"),
|
||||
r.get("amount"),
|
||||
),
|
||||
),
|
||||
}
|
||||
|
||||
EXTENDED_EOD_COPY = {
|
||||
"limit_events": (
|
||||
"INSERT OR REPLACE INTO eod_limit_events "
|
||||
"SELECT ts_code,trade_date,limit_type,name,industry,close,pct_chg,amount,"
|
||||
"limit_amount,float_mv,total_mv,turnover_ratio,fd_amount,first_time,last_time,"
|
||||
"open_times,up_stat,limit_times,batch_id "
|
||||
"FROM staging_limit_events WHERE batch_id = ?"
|
||||
),
|
||||
"popularity": (
|
||||
"INSERT OR REPLACE INTO eod_popularity "
|
||||
"SELECT ts_code,trade_date,source,ts_name,rank,pct_change,current_price,hot,concept,data_type,batch_id "
|
||||
"FROM staging_popularity WHERE batch_id = ?"
|
||||
),
|
||||
"dragon_tiger": (
|
||||
"INSERT OR REPLACE INTO eod_dragon_tiger "
|
||||
"SELECT ts_code,trade_date,hm_name,ts_name,buy_amount,sell_amount,net_amount,"
|
||||
"hm_orgs,tag,pct_change,reason,batch_id "
|
||||
"FROM staging_dragon_tiger WHERE batch_id = ?"
|
||||
),
|
||||
"sector_daily": (
|
||||
"INSERT OR REPLACE INTO eod_sector_daily "
|
||||
"SELECT ts_code,trade_date,family,name,open,high,low,close,pre_close,"
|
||||
"pct_change,vol,turnover_rate,amount,batch_id "
|
||||
"FROM staging_sector_daily WHERE batch_id = ?"
|
||||
),
|
||||
}
|
||||
@@ -7,10 +7,9 @@ from contextlib import contextmanager
|
||||
from pathlib import Path
|
||||
from typing import Any
|
||||
|
||||
from datahub.datasets_ext import EXTENDED_DATASET_TABLES, EXTENDED_SCHEMA
|
||||
from datahub.timeutil import isoformat
|
||||
|
||||
_BASE_SCHEMA = """
|
||||
SCHEMA = """
|
||||
CREATE TABLE IF NOT EXISTS schema_migrations (
|
||||
version INTEGER PRIMARY KEY,
|
||||
applied_at TEXT NOT NULL
|
||||
@@ -239,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,
|
||||
@@ -283,8 +295,6 @@ CREATE INDEX IF NOT EXISTS idx_eod_bars_date ON eod_bars(trade_date, batch_id);
|
||||
CREATE INDEX IF NOT EXISTS idx_calendar_open ON trade_calendar(is_open, cal_date);
|
||||
"""
|
||||
|
||||
SCHEMA = _BASE_SCHEMA + EXTENDED_SCHEMA
|
||||
|
||||
DATASET_TABLES = {
|
||||
"daily": ("eod_bars", "staging_bars"),
|
||||
"valuation": ("eod_valuation", "staging_valuation"),
|
||||
@@ -292,7 +302,6 @@ DATASET_TABLES = {
|
||||
"auction": ("eod_auction", "staging_auction"),
|
||||
"index_daily": ("eod_index_bars", "staging_index_bars"),
|
||||
"stocks": ("eod_stocks", "staging_stocks"),
|
||||
**EXTENDED_DATASET_TABLES,
|
||||
}
|
||||
|
||||
|
||||
|
||||
@@ -156,95 +156,6 @@ def normalize_stock(row: dict[str, Any]) -> dict[str, Any]:
|
||||
}
|
||||
|
||||
|
||||
def normalize_limit_event(row: dict[str, Any]) -> dict[str, Any]:
|
||||
"""limit_list_d. float_mv/total_mv/limit_amount are 万元 → yuan; amount/fd_amount already yuan."""
|
||||
return {
|
||||
"ts_code": _code(row.get("ts_code")),
|
||||
"trade_date": _date(row.get("trade_date")),
|
||||
"limit_type": str(row.get("limit_type") or "").strip().upper() or "U",
|
||||
"name": str(row.get("name") or "").strip() or None,
|
||||
"industry": str(row.get("industry") or "").strip() or None,
|
||||
"close": round4(finite_number(row.get("close"))),
|
||||
"pct_chg": round4(finite_number(row.get("pct_chg"))),
|
||||
"amount": round4(finite_number(row.get("amount"))),
|
||||
"limit_amount": round4(_scale(row.get("limit_amount"), AMOUNT_WAN_YUAN)),
|
||||
"float_mv": round4(_scale(row.get("float_mv"), AMOUNT_WAN_YUAN)),
|
||||
"total_mv": round4(_scale(row.get("total_mv"), AMOUNT_WAN_YUAN)),
|
||||
"turnover_ratio": round4(finite_number(row.get("turnover_ratio"))),
|
||||
"fd_amount": round4(finite_number(row.get("fd_amount"))),
|
||||
"first_time": str(row.get("first_time") or "").strip() or None,
|
||||
"last_time": str(row.get("last_time") or "").strip() or None,
|
||||
"open_times": _optional_int(row.get("open_times")),
|
||||
"up_stat": str(row.get("up_stat") or "").strip() or None,
|
||||
"limit_times": _optional_int(row.get("limit_times")),
|
||||
}
|
||||
|
||||
|
||||
def normalize_popularity(row: dict[str, Any], source: str = "") -> dict[str, Any]:
|
||||
src = str(source or row.get("source") or "").strip().lower() or "ths"
|
||||
return {
|
||||
"ts_code": _code(row.get("ts_code")),
|
||||
"trade_date": _date(row.get("trade_date")),
|
||||
"source": src,
|
||||
"ts_name": str(row.get("ts_name") or row.get("name") or "").strip() or None,
|
||||
"rank": _optional_int(row.get("rank")),
|
||||
"pct_change": round4(
|
||||
finite_number(row.get("pct_change") if row.get("pct_change") is not None else row.get("pct_chg"))
|
||||
),
|
||||
"current_price": round4(finite_number(row.get("current_price") or row.get("price"))),
|
||||
"hot": round4(finite_number(row.get("hot"))),
|
||||
"concept": str(row.get("concept") or "").strip() or None,
|
||||
"data_type": str(row.get("data_type") or "").strip() or None,
|
||||
}
|
||||
|
||||
|
||||
def normalize_dragon_tiger(row: dict[str, Any]) -> dict[str, Any]:
|
||||
"""hm_detail amounts are 万元 → yuan."""
|
||||
return {
|
||||
"ts_code": _code(row.get("ts_code")),
|
||||
"trade_date": _date(row.get("trade_date")),
|
||||
"hm_name": str(row.get("hm_name") or "未命名游资").strip() or "未命名游资",
|
||||
"ts_name": str(row.get("ts_name") or row.get("name") or "").strip() or None,
|
||||
"buy_amount": round4(_scale(row.get("buy_amount"), AMOUNT_WAN_YUAN)),
|
||||
"sell_amount": round4(_scale(row.get("sell_amount"), AMOUNT_WAN_YUAN)),
|
||||
"net_amount": round4(_scale(row.get("net_amount"), AMOUNT_WAN_YUAN)),
|
||||
"hm_orgs": str(row.get("hm_orgs") or "").strip() or None,
|
||||
"tag": str(row.get("tag") or "").strip() or None,
|
||||
"pct_change": round4(finite_number(row.get("pct_change"))),
|
||||
"reason": str(row.get("reason") or "").strip() or None,
|
||||
}
|
||||
|
||||
|
||||
def normalize_sector_daily(row: dict[str, Any], family: str = "ths") -> dict[str, Any]:
|
||||
fam = str(family or row.get("family") or "ths").strip().lower()
|
||||
return {
|
||||
"ts_code": _code(row.get("ts_code")),
|
||||
"trade_date": _date(row.get("trade_date")),
|
||||
"family": fam,
|
||||
"name": str(row.get("name") or "").strip() or None,
|
||||
"open": round4(finite_number(row.get("open"))),
|
||||
"high": round4(finite_number(row.get("high"))),
|
||||
"low": round4(finite_number(row.get("low"))),
|
||||
"close": round4(finite_number(row.get("close"))),
|
||||
"pre_close": round4(finite_number(row.get("pre_close"))),
|
||||
"pct_change": round4(
|
||||
finite_number(row.get("pct_change") if row.get("pct_change") is not None else row.get("pct_chg"))
|
||||
),
|
||||
"vol": round4(finite_number(row.get("vol"))),
|
||||
"turnover_rate": round4(finite_number(row.get("turnover_rate"))),
|
||||
"amount": round4(finite_number(row.get("amount"))),
|
||||
}
|
||||
|
||||
|
||||
def _optional_int(value: Any) -> int | None:
|
||||
if value in (None, ""):
|
||||
return None
|
||||
try:
|
||||
return int(float(value))
|
||||
except (TypeError, ValueError):
|
||||
return None
|
||||
|
||||
|
||||
def apply_qfq(price: float | None, factor: float | None, latest_factor: float | None) -> float | None:
|
||||
if price is None:
|
||||
return None
|
||||
@@ -273,11 +184,6 @@ NORMALIZERS = {
|
||||
"calendar": normalize_calendar,
|
||||
"stock_basic": normalize_stock,
|
||||
"stocks": normalize_stock,
|
||||
"limit_events": normalize_limit_event,
|
||||
"limit_list_d": normalize_limit_event,
|
||||
"popularity": normalize_popularity,
|
||||
"dragon_tiger": normalize_dragon_tiger,
|
||||
"sector_daily": normalize_sector_daily,
|
||||
}
|
||||
|
||||
|
||||
|
||||
@@ -9,33 +9,30 @@ from typing import Any
|
||||
|
||||
from datahub.adapters.base import AdapterError
|
||||
from datahub.adapters.tushare import DEFAULT_INDEX_CODES, WEBSITE_INDEX_CODES, TushareAdapter
|
||||
from datahub.datasets_ext import (
|
||||
EXTENDED_EOD_COPY,
|
||||
EXTENDED_SOFT_DATASETS,
|
||||
EXTENDED_STAGING_INSERT,
|
||||
)
|
||||
from datahub.db import DATASET_TABLES, HubDB
|
||||
from datahub.governance.circuit import CircuitBreaker
|
||||
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
|
||||
|
||||
LOGGER = get_logger()
|
||||
|
||||
HARD_DATASETS = {"daily", "valuation", "index_daily"}
|
||||
SOFT_DATASETS = {"moneyflow", "auction"} | EXTENDED_SOFT_DATASETS
|
||||
OFFICIAL_DATASETS = HARD_DATASETS | {"moneyflow", "auction"} # A/B retry scope unchanged
|
||||
SOFT_DATASETS = {"moneyflow", "auction"}
|
||||
OFFICIAL_DATASETS = HARD_DATASETS | SOFT_DATASETS
|
||||
STOCKS_DATASET = "stocks"
|
||||
STOCK_SNAPSHOT_FIELDS = ("ts_code", "symbol", "name", "area", "industry", "market", "list_status", "list_date")
|
||||
EOD_A_DATASETS = ("daily", "valuation", "moneyflow", "auction")
|
||||
EOD_B_DATASETS = ("index_daily",)
|
||||
EOD_C_DATASETS = ("limit_events",)
|
||||
EOD_D_DATASETS = ("dragon_tiger",)
|
||||
EOD_E_DATASETS = ("sector_daily",)
|
||||
EOD_F_DATASETS = ("popularity",)
|
||||
EMPTY_BATCH_ERROR = "empty official batch: 0 valid rows"
|
||||
|
||||
STAGING_INSERT = {
|
||||
@@ -89,7 +86,6 @@ STAGING_INSERT = {
|
||||
r.get("close"), r.get("pct_chg"), r.get("volume"), r.get("amount"),
|
||||
),
|
||||
),
|
||||
**EXTENDED_STAGING_INSERT,
|
||||
}
|
||||
|
||||
EOD_COPY = {
|
||||
@@ -124,7 +120,6 @@ EOD_COPY = {
|
||||
"SELECT ts_code,trade_date,open,high,low,close,pct_chg,volume,amount,batch_id "
|
||||
"FROM staging_index_bars WHERE batch_id = ?"
|
||||
),
|
||||
**EXTENDED_EOD_COPY,
|
||||
}
|
||||
|
||||
|
||||
@@ -683,71 +678,255 @@ class Pipeline:
|
||||
def run_eod_batch_b(self, trade_date: str, force: bool = False) -> dict[str, Any]:
|
||||
return self.run_release_group(EOD_B_DATASETS, trade_date, force=force)
|
||||
|
||||
def run_extended_soft(self, datasets: tuple[str, ...], trade_date: str, force: bool = False) -> dict[str, Any]:
|
||||
"""Publish extended soft datasets independently (not A/B atomic)."""
|
||||
results: dict[str, Any] = {}
|
||||
day = yyyymmdd(trade_date)
|
||||
for dataset in datasets:
|
||||
if not force and self.active_batch(dataset, day):
|
||||
results[dataset] = {
|
||||
"dataset": dataset,
|
||||
"trade_date": day,
|
||||
"state": "skipped",
|
||||
"reason": "already_published",
|
||||
}
|
||||
continue
|
||||
try:
|
||||
rows = self._fetch_dataset(dataset, day)
|
||||
if not rows and dataset in {"popularity", "dragon_tiger"}:
|
||||
results[dataset] = {
|
||||
"dataset": dataset,
|
||||
"trade_date": day,
|
||||
"state": "skipped",
|
||||
"reason": "upstream_empty",
|
||||
"rows": 0,
|
||||
}
|
||||
continue
|
||||
results[dataset] = self.run_dataset(dataset, day, prepared_rows=rows)
|
||||
except Exception as exc:
|
||||
results[dataset] = {
|
||||
"dataset": dataset,
|
||||
"trade_date": day,
|
||||
"state": "failed",
|
||||
"error": str(exc),
|
||||
}
|
||||
LOGGER.exception("extended soft publish failed dataset=%s date=%s", dataset, day)
|
||||
return results
|
||||
|
||||
def run_eod_batch_c(self, trade_date: str, force: bool = False) -> dict[str, Any]:
|
||||
return self.run_extended_soft(EOD_C_DATASETS, trade_date, force=force)
|
||||
|
||||
def run_eod_batch_d(self, trade_date: str, force: bool = False) -> dict[str, Any]:
|
||||
return self.run_extended_soft(EOD_D_DATASETS, trade_date, force=force)
|
||||
|
||||
def run_eod_batch_e(self, trade_date: str, force: bool = False) -> dict[str, Any]:
|
||||
return self.run_extended_soft(EOD_E_DATASETS, trade_date, force=force)
|
||||
|
||||
def run_eod_batch_f(self, trade_date: str, force: bool = False) -> dict[str, Any]:
|
||||
return self.run_extended_soft(EOD_F_DATASETS, trade_date, force=force)
|
||||
|
||||
def force_republish_boundary(self, dataset: str, trade_date: str) -> dict[str, Any]:
|
||||
"""Force-republish the full A/B consistency boundary that owns ``dataset``.
|
||||
|
||||
CLI ``eod-refresh --force`` and admin manual backfill must not publish a
|
||||
single official member alone — that would mix old and new batches inside
|
||||
the same trade date. Naming any A-group member (or stocks) rebuilds the
|
||||
whole A group; naming ``index_daily`` rebuilds B. Extended soft datasets
|
||||
republish independently.
|
||||
whole A group; naming ``index_daily`` rebuilds B.
|
||||
"""
|
||||
name = str(dataset or "").strip()
|
||||
if name in EOD_A_DATASETS or name == STOCKS_DATASET:
|
||||
return self.run_eod_batch_a(trade_date, force=True)
|
||||
if name in EOD_B_DATASETS:
|
||||
return self.run_eod_batch_b(trade_date, force=True)
|
||||
if name in EXTENDED_SOFT_DATASETS:
|
||||
return self.run_extended_soft((name,), 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, ...],
|
||||
@@ -1083,16 +1262,7 @@ class Pipeline:
|
||||
)
|
||||
listed_n = int((listed or {}).get("n") or 0)
|
||||
row_n = len(rows)
|
||||
if dataset == "limit_events":
|
||||
keys = [(row.get("ts_code"), row.get("trade_date"), row.get("limit_type")) for row in rows]
|
||||
elif dataset == "popularity":
|
||||
keys = [(row.get("ts_code"), row.get("trade_date"), row.get("source")) for row in rows]
|
||||
elif dataset == "dragon_tiger":
|
||||
keys = [(row.get("ts_code"), row.get("trade_date"), row.get("hm_name")) for row in rows]
|
||||
elif dataset == "sector_daily":
|
||||
keys = [(row.get("ts_code"), row.get("trade_date"), row.get("family")) for row in rows]
|
||||
else:
|
||||
keys = [(row.get("ts_code"), row.get("trade_date")) for row in rows]
|
||||
keys = [(row.get("ts_code"), row.get("trade_date")) for row in rows]
|
||||
dup = row_n - len(set(keys))
|
||||
if dup:
|
||||
errors.append(f"duplicate keys: {dup}")
|
||||
@@ -1113,8 +1283,7 @@ class Pipeline:
|
||||
errors.append(EMPTY_BATCH_ERROR)
|
||||
field_report = self._field_gate(dataset, trade_date, rows, errors)
|
||||
if dataset in SOFT_DATASETS:
|
||||
allow_empty = dataset in {"popularity", "dragon_tiger", "moneyflow", "auction"}
|
||||
hard_fail = bool(dup or bad_date or (empty and not allow_empty))
|
||||
hard_fail = bool(dup or bad_date or empty)
|
||||
else:
|
||||
hard_fail = bool(errors) and (dataset in HARD_DATASETS or dataset == STOCKS_DATASET)
|
||||
report = {
|
||||
|
||||
@@ -1,233 +0,0 @@
|
||||
"""Provisional (盘中观察) serving: quotes, index quotes, intraday points.
|
||||
|
||||
Free sources only. Never writes official eod_* tables. Uses rt_cache + LKG.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import json
|
||||
import time
|
||||
from datetime import datetime
|
||||
from typing import Any
|
||||
|
||||
from datahub.adapters.eastmoney import EastmoneyAdapter
|
||||
from datahub.adapters.tencent import TencentAdapter
|
||||
from datahub.codes import resolve_code
|
||||
from datahub.db import HubDB
|
||||
from datahub.governance.lkg import LastKnownGood
|
||||
from datahub.timeutil import isoformat, now_shanghai, yyyymmdd
|
||||
|
||||
QUOTE_TTL = 60
|
||||
INDEX_TTL = 60
|
||||
INTRADAY_TTL = 20
|
||||
|
||||
|
||||
class RealtimeApiError(RuntimeError):
|
||||
def __init__(self, code: str, message: str) -> None:
|
||||
super().__init__(message)
|
||||
self.code = code
|
||||
self.message = message
|
||||
|
||||
|
||||
def _envelope(data: Any, meta: dict[str, Any]) -> dict[str, Any]:
|
||||
from datahub import SCHEMA_VERSION
|
||||
|
||||
return {"schema_version": SCHEMA_VERSION, "data": data, "meta": meta}
|
||||
|
||||
|
||||
def fetch_index_quotes(db: HubDB) -> dict[str, Any]:
|
||||
cache_key = "indexes:quotes"
|
||||
cached = _read_cache(db, cache_key)
|
||||
if cached is not None:
|
||||
return cached
|
||||
eastmoney = EastmoneyAdapter()
|
||||
try:
|
||||
rows = eastmoney.fetch_indices()
|
||||
source = "eastmoney:ulist"
|
||||
except Exception:
|
||||
rows = TencentAdapter().fetch_indices()
|
||||
source = "tencent:qt"
|
||||
if len(rows) < 3:
|
||||
raise RealtimeApiError("SOURCE_UNAVAILABLE", "index quotes incomplete")
|
||||
payload = _envelope(
|
||||
rows,
|
||||
{
|
||||
"tier": "provisional",
|
||||
"trade_date": yyyymmdd(now_shanghai()),
|
||||
"source": source,
|
||||
"stale": False,
|
||||
"staleness_seconds": 0,
|
||||
"published_at": isoformat(now_shanghai()),
|
||||
},
|
||||
)
|
||||
_write_cache(db, cache_key, payload, INDEX_TTL, source)
|
||||
return payload
|
||||
|
||||
|
||||
def fetch_quotes(db: HubDB, codes: list[str]) -> dict[str, Any]:
|
||||
if not codes:
|
||||
raise RealtimeApiError("INVALID_ARGUMENT", "codes is required")
|
||||
resolved: list[str] = []
|
||||
for code in codes[:60]:
|
||||
item = resolve_code(db, code) or _guess_ts_code(code)
|
||||
if item:
|
||||
resolved.append(item)
|
||||
if not resolved:
|
||||
raise RealtimeApiError("INVALID_ARGUMENT", "no resolvable codes")
|
||||
cache_key = "quotes:" + ",".join(sorted(resolved))
|
||||
cached = _read_cache(db, cache_key)
|
||||
if cached is not None:
|
||||
return cached
|
||||
adapter = EastmoneyAdapter()
|
||||
try:
|
||||
rows = adapter.fetch_quotes(resolved)
|
||||
source = "eastmoney:clist"
|
||||
except Exception as exc:
|
||||
raise RealtimeApiError("SOURCE_UNAVAILABLE", f"quotes unavailable: {exc}") from exc
|
||||
payload = _envelope(
|
||||
rows,
|
||||
{
|
||||
"tier": "provisional",
|
||||
"trade_date": yyyymmdd(now_shanghai()),
|
||||
"source": source,
|
||||
"stale": False,
|
||||
"staleness_seconds": 0,
|
||||
"published_at": isoformat(now_shanghai()),
|
||||
},
|
||||
)
|
||||
_write_cache(db, cache_key, payload, QUOTE_TTL, source)
|
||||
return payload
|
||||
|
||||
|
||||
def fetch_intraday(db: HubDB, code: str, date: str = "") -> dict[str, Any]:
|
||||
ts_code = resolve_code(db, code) or _guess_ts_code(code)
|
||||
if not ts_code:
|
||||
raise RealtimeApiError("INVALID_ARGUMENT", f"ambiguous code: {code}")
|
||||
cache_key = f"intraday:{ts_code}:{date or 'today'}"
|
||||
cached = _read_cache(db, cache_key)
|
||||
if cached is not None:
|
||||
return cached
|
||||
adapter = EastmoneyAdapter()
|
||||
try:
|
||||
payload_data = adapter.fetch_intraday(ts_code, date)
|
||||
source = "eastmoney:trends2"
|
||||
except Exception as exc:
|
||||
recovered = _load_intraday_lkg(db, ts_code, date)
|
||||
if recovered is None:
|
||||
raise RealtimeApiError("SOURCE_UNAVAILABLE", f"intraday unavailable: {exc}") from exc
|
||||
return recovered
|
||||
payload = _envelope(
|
||||
payload_data,
|
||||
{
|
||||
"tier": "provisional",
|
||||
"trade_date": yyyymmdd(payload_data.get("trade_date") or date or now_shanghai()),
|
||||
"source": source,
|
||||
"stale": False,
|
||||
"staleness_seconds": 0,
|
||||
"published_at": isoformat(now_shanghai()),
|
||||
},
|
||||
)
|
||||
_write_cache(db, cache_key, payload, INTRADAY_TTL, source)
|
||||
return payload
|
||||
|
||||
|
||||
def _load_intraday_lkg(db: HubDB, ts_code: str, date: str = "") -> dict[str, Any] | None:
|
||||
store = LastKnownGood(db)
|
||||
keys = [f"intraday:{ts_code}:{date or 'today'}"]
|
||||
if date:
|
||||
keys.append(f"intraday:{ts_code}:today")
|
||||
for key in keys:
|
||||
item = store.load(key)
|
||||
payload = _lkg_payload(item)
|
||||
if payload is not None:
|
||||
return payload
|
||||
row = db.fetchone(
|
||||
"SELECT * FROM last_known_good WHERE cache_key LIKE ? ORDER BY stored_at DESC LIMIT 1",
|
||||
(f"intraday:{ts_code}:%",),
|
||||
)
|
||||
if not row:
|
||||
return None
|
||||
try:
|
||||
raw = json.loads(row["payload"])
|
||||
except json.JSONDecodeError:
|
||||
return None
|
||||
return _mark_stale(raw) if isinstance(raw, dict) else None
|
||||
|
||||
|
||||
def _lkg_payload(item: dict[str, Any] | None) -> dict[str, Any] | None:
|
||||
if not item:
|
||||
return None
|
||||
payload = item.get("payload")
|
||||
return _mark_stale(payload) if isinstance(payload, dict) else None
|
||||
|
||||
|
||||
def _mark_stale(payload: dict[str, Any]) -> dict[str, Any] | None:
|
||||
data = payload.get("data")
|
||||
if not isinstance(data, dict) or not data.get("points"):
|
||||
return None
|
||||
stamped = dict(payload)
|
||||
meta = dict(stamped.get("meta") or {})
|
||||
meta["stale"] = True
|
||||
stamped["meta"] = meta
|
||||
return stamped
|
||||
|
||||
|
||||
def _guess_ts_code(code: str) -> str | None:
|
||||
raw = str(code or "").strip().upper()
|
||||
if "." in raw:
|
||||
return raw
|
||||
if len(raw) == 6 and raw.isdigit():
|
||||
if raw.startswith(("5", "6", "9")):
|
||||
return f"{raw}.SH"
|
||||
return f"{raw}.SZ"
|
||||
return None
|
||||
|
||||
|
||||
def _read_cache(db: HubDB, cache_key: str) -> dict[str, Any] | None:
|
||||
row = db.fetchone("SELECT * FROM rt_cache WHERE cache_key = ?", (cache_key,))
|
||||
if not row:
|
||||
return None
|
||||
expires = str(row.get("expires_at") or "")
|
||||
now = isoformat(now_shanghai())
|
||||
if expires and expires < now:
|
||||
return None
|
||||
try:
|
||||
payload = json.loads(row["payload"])
|
||||
except json.JSONDecodeError:
|
||||
return None
|
||||
if isinstance(payload, dict) and isinstance(payload.get("meta"), dict):
|
||||
stored = str(row.get("stored_at") or "")
|
||||
try:
|
||||
age = max(0, int(time.time() - datetime.fromisoformat(stored).timestamp()))
|
||||
except Exception:
|
||||
age = 0
|
||||
payload["meta"]["staleness_seconds"] = age
|
||||
payload["meta"]["stale"] = age > QUOTE_TTL
|
||||
return payload
|
||||
|
||||
|
||||
def _write_cache(db: HubDB, cache_key: str, payload: dict[str, Any], ttl: int, source: str) -> None:
|
||||
from datetime import timedelta
|
||||
|
||||
now = now_shanghai()
|
||||
stored = isoformat(now)
|
||||
expires = isoformat(now + timedelta(seconds=ttl))
|
||||
db.execute(
|
||||
"""
|
||||
INSERT INTO rt_cache(cache_key, payload, source, stored_at, expires_at)
|
||||
VALUES (?,?,?,?,?)
|
||||
ON CONFLICT(cache_key) DO UPDATE SET
|
||||
payload=excluded.payload, source=excluded.source,
|
||||
stored_at=excluded.stored_at, expires_at=excluded.expires_at
|
||||
""",
|
||||
(cache_key, json.dumps(payload, ensure_ascii=False), source, stored, expires),
|
||||
)
|
||||
db.execute(
|
||||
"""
|
||||
INSERT INTO last_known_good(cache_key, payload, source, stored_at)
|
||||
VALUES (?,?,?,?)
|
||||
ON CONFLICT(cache_key) DO UPDATE SET
|
||||
payload=excluded.payload, source=excluded.source, stored_at=excluded.stored_at
|
||||
""",
|
||||
(cache_key, json.dumps(payload, ensure_ascii=False), source, stored),
|
||||
)
|
||||
@@ -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]
|
||||
@@ -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()
|
||||
|
||||
@@ -48,11 +64,8 @@ class Scheduler:
|
||||
"precheck": self._precheck,
|
||||
"eod_a": self._eod_a,
|
||||
"eod_b": self._eod_b,
|
||||
"eod_c": self._eod_c,
|
||||
"eod_d": self._eod_d,
|
||||
"eod_e": self._eod_e,
|
||||
"eod_f": self._eod_f,
|
||||
"eod_retry": self._eod_retry,
|
||||
"eod_revise": self._eod_revise,
|
||||
"stocks_refresh": self._stocks_refresh,
|
||||
"cleanup": self._cleanup,
|
||||
"backup": self._backup,
|
||||
@@ -91,10 +104,6 @@ class Scheduler:
|
||||
("precheck", time(8, 45)),
|
||||
("eod_a", time(15, 5)),
|
||||
("eod_b", time(15, 10)),
|
||||
("eod_c", time(16, 40)),
|
||||
("eod_d", time(16, 45)),
|
||||
("eod_e", time(18, 5)),
|
||||
("eod_f", time(22, 40)),
|
||||
("cleanup", time(0, 30)),
|
||||
("backup", time(0, 40)),
|
||||
]
|
||||
@@ -107,7 +116,7 @@ class Scheduler:
|
||||
key = (job_id, day, at.strftime("%H%M"))
|
||||
if key in self._fired:
|
||||
continue
|
||||
if job_id in {"eod_a", "eod_b", "eod_c", "eod_d", "eod_e", "eod_f", "stocks_refresh"} and not open_day:
|
||||
if job_id in {"eod_a", "eod_b", "stocks_refresh"} and not open_day:
|
||||
self._fired.add(key)
|
||||
continue
|
||||
self._fired.add(key)
|
||||
@@ -118,7 +127,7 @@ class Scheduler:
|
||||
try:
|
||||
self.run_job(job_id, day)
|
||||
except Exception:
|
||||
if job_id not in {"eod_a", "eod_b", "eod_c", "eod_d", "eod_e", "eod_f", "stocks_refresh"}:
|
||||
if job_id not in {"eod_a", "eod_b", "stocks_refresh"}:
|
||||
raise
|
||||
# Keep the tick alive; evening retries take over.
|
||||
LOGGER.exception("scheduled job %s failed for %s", job_id, day)
|
||||
@@ -126,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
|
||||
|
||||
# ------------------------------------------------------------------
|
||||
@@ -227,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
|
||||
@@ -325,21 +531,12 @@ class Scheduler:
|
||||
def _eod_b(self, trade_date: str) -> dict[str, Any]:
|
||||
return self.pipeline.run_eod_batch_b(trade_date)
|
||||
|
||||
def _eod_c(self, trade_date: str) -> dict[str, Any]:
|
||||
return self.pipeline.run_eod_batch_c(trade_date)
|
||||
|
||||
def _eod_d(self, trade_date: str) -> dict[str, Any]:
|
||||
return self.pipeline.run_eod_batch_d(trade_date)
|
||||
|
||||
def _eod_e(self, trade_date: str) -> dict[str, Any]:
|
||||
return self.pipeline.run_eod_batch_e(trade_date)
|
||||
|
||||
def _eod_f(self, trade_date: str) -> dict[str, Any]:
|
||||
return self.pipeline.run_eod_batch_f(trade_date)
|
||||
|
||||
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)
|
||||
|
||||
|
||||
@@ -73,20 +73,6 @@ class V1API:
|
||||
return self.moneyflow(q)
|
||||
if path == "/v1/auction":
|
||||
return self.auction(q)
|
||||
if path == "/v1/limit-events":
|
||||
return self.limit_events(q)
|
||||
if path == "/v1/popularity":
|
||||
return self.popularity(q)
|
||||
if path == "/v1/dragon-tiger":
|
||||
return self.dragon_tiger(q)
|
||||
if path == "/v1/sectors":
|
||||
return self.sectors(q)
|
||||
if path == "/v1/quotes/latest":
|
||||
return self.quotes_latest(q)
|
||||
if path == "/v1/indexes/quotes":
|
||||
return self.index_quotes(q)
|
||||
if path == "/v1/intraday/points":
|
||||
return self.intraday_points(q)
|
||||
if path == "/v1/datasets/status":
|
||||
return self.dataset_status(q.get("date") or "")
|
||||
if path == "/v1/batches":
|
||||
@@ -211,80 +197,9 @@ class V1API:
|
||||
def auction(self, q: dict[str, str]) -> dict[str, Any]:
|
||||
return self._published_rows(dataset="auction", table="eod_auction", q=q, source="tushare:stk_auction")
|
||||
|
||||
def limit_events(self, q: dict[str, str]) -> dict[str, Any]:
|
||||
return self._published_rows(
|
||||
dataset="limit_events",
|
||||
table="eod_limit_events",
|
||||
q=q,
|
||||
source="tushare:limit_list_d",
|
||||
extra_filters={"limit_type": q.get("limit_type") or ""},
|
||||
)
|
||||
|
||||
def popularity(self, q: dict[str, str]) -> dict[str, Any]:
|
||||
return self._published_rows(
|
||||
dataset="popularity",
|
||||
table="eod_popularity",
|
||||
q=q,
|
||||
source="tushare:ths_hot+dc_hot",
|
||||
extra_filters={"source": q.get("source") or ""},
|
||||
)
|
||||
|
||||
def dragon_tiger(self, q: dict[str, str]) -> dict[str, Any]:
|
||||
return self._published_rows(
|
||||
dataset="dragon_tiger",
|
||||
table="eod_dragon_tiger",
|
||||
q=q,
|
||||
source="tushare:hm_detail",
|
||||
)
|
||||
|
||||
def sectors(self, q: dict[str, str]) -> dict[str, Any]:
|
||||
return self._published_rows(
|
||||
dataset="sector_daily",
|
||||
table="eod_sector_daily",
|
||||
q=q,
|
||||
source="tushare:ths_daily+dc_index+sw_daily",
|
||||
extra_filters={"family": q.get("family") or ""},
|
||||
)
|
||||
|
||||
def quotes_latest(self, q: dict[str, str]) -> dict[str, Any]:
|
||||
from datahub.realtime_serve import RealtimeApiError, fetch_quotes
|
||||
|
||||
codes = [item.strip() for item in str(q.get("codes") or "").split(",") if item.strip()]
|
||||
try:
|
||||
return fetch_quotes(self.db, codes)
|
||||
except RealtimeApiError as exc:
|
||||
raise ApiError(exc.code, exc.message) from exc
|
||||
|
||||
def index_quotes(self, q: dict[str, str]) -> dict[str, Any]:
|
||||
from datahub.realtime_serve import RealtimeApiError, fetch_index_quotes
|
||||
|
||||
try:
|
||||
return fetch_index_quotes(self.db)
|
||||
except RealtimeApiError as exc:
|
||||
raise ApiError(exc.code, exc.message) from exc
|
||||
|
||||
def intraday_points(self, q: dict[str, str]) -> dict[str, Any]:
|
||||
from datahub.realtime_serve import RealtimeApiError, fetch_intraday
|
||||
|
||||
code = str(q.get("code") or "").strip()
|
||||
if not code:
|
||||
raise ApiError("INVALID_ARGUMENT", "code is required")
|
||||
raw_date = str(q.get("date") or "").strip()
|
||||
try:
|
||||
trade_date = yyyymmdd(raw_date or now_shanghai())
|
||||
except ValueError as exc:
|
||||
raise ApiError("INVALID_ARGUMENT", str(exc)) from exc
|
||||
try:
|
||||
return fetch_intraday(self.db, code, trade_date)
|
||||
except RealtimeApiError as exc:
|
||||
raise ApiError(exc.code, exc.message) from exc
|
||||
|
||||
def dataset_status(self, date: str) -> dict[str, Any]:
|
||||
trade_date = yyyymmdd(date or now_shanghai())
|
||||
datasets = (
|
||||
"daily", "valuation", "moneyflow", "auction", "index_daily", "stocks",
|
||||
"limit_events", "popularity", "dragon_tiger", "sector_daily",
|
||||
)
|
||||
datasets = ("daily", "valuation", "moneyflow", "auction", "index_daily", "stocks")
|
||||
items = []
|
||||
for dataset in datasets:
|
||||
pub = self.db.fetchone(
|
||||
@@ -329,7 +244,6 @@ class V1API:
|
||||
source: str,
|
||||
adjust: str = "none",
|
||||
default_code: str = "",
|
||||
extra_filters: dict[str, str] | None = None,
|
||||
) -> dict[str, Any]:
|
||||
trade_date = q.get("date") or q.get("trade_date") or ""
|
||||
code = q.get("code") or default_code
|
||||
@@ -350,7 +264,6 @@ class V1API:
|
||||
if resolved is None:
|
||||
raise ApiError("INVALID_ARGUMENT", f"ambiguous code: {code}")
|
||||
ts_code = resolved
|
||||
filters = {key: value for key, value in (extra_filters or {}).items() if value}
|
||||
# For a range, use per-date published batch. Single-date is the common path.
|
||||
if start == end:
|
||||
pub = self.db.fetchone(
|
||||
@@ -369,9 +282,6 @@ class V1API:
|
||||
if ts_code:
|
||||
sql += " AND ts_code = ?"
|
||||
params.append(ts_code)
|
||||
for key, value in filters.items():
|
||||
sql += f" AND {key} = ?"
|
||||
params.append(value)
|
||||
sql += " ORDER BY ts_code LIMIT ? OFFSET ?"
|
||||
params.extend([limit, offset])
|
||||
rows = [dict(row) for row in self.db.fetchall(sql, tuple(params))]
|
||||
@@ -407,9 +317,6 @@ class V1API:
|
||||
if ts_code:
|
||||
sql += " AND ts_code = ?"
|
||||
params.append(ts_code)
|
||||
for key, value in filters.items():
|
||||
sql += f" AND {key} = ?"
|
||||
params.append(value)
|
||||
sql += " ORDER BY ts_code"
|
||||
rows.extend(self.db.fetchall(sql, tuple(params)))
|
||||
sliced = rows[offset: offset + limit]
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -45,30 +45,6 @@ RAW = {
|
||||
{"ts_code": "600000.SH", "trade_date": "20240902", "vol": 100, "price": 10.15, "amount": 1500000, "pre_close": 10.00, "turnover_rate": 0.1, "volume_ratio": 1.2, "float_share": 2000},
|
||||
{"ts_code": "000001.SZ", "trade_date": "20240902", "vol": 80, "price": 11.05, "amount": 1200000, "pre_close": 11.10, "turnover_rate": 0.2, "volume_ratio": 0.9, "float_share": 1800},
|
||||
],
|
||||
"limit_list_d": [
|
||||
{"trade_date": "20240902", "ts_code": "600000.SH", "industry": "银行", "name": "浦发银行", "close": 10.2, "pct_chg": 9.95, "amount": 1e8, "limit_amount": 5000, "float_mv": 800, "total_mv": 1000, "turnover_ratio": 5.0, "fd_amount": 2e7, "first_time": "09:30:01", "last_time": "14:55:00", "open_times": 0, "up_stat": "1/1", "limit_times": 1, "limit_type": "U"},
|
||||
],
|
||||
"ths_hot": [
|
||||
{"ts_code": "600000.SH", "ts_name": "浦发银行", "hot": 90.0, "rank": 1, "pct_change": 1.2, "current_price": 10.2, "concept": "银行", "data_type": "热股", "trade_date": "20240902"},
|
||||
],
|
||||
"dc_hot": [
|
||||
{"ts_code": "600000.SH", "ts_name": "浦发银行", "rank": 2, "pct_change": 1.2, "current_price": 10.2, "hot": 80.0, "concept": "银行", "data_type": "A股市场", "trade_date": "20240902"},
|
||||
],
|
||||
"hm_detail": [
|
||||
{"trade_date": "20240902", "ts_code": "600000.SH", "ts_name": "浦发银行", "buy_amount": 1000, "sell_amount": 200, "net_amount": 800, "hm_name": "测试游资", "hm_orgs": "某某营业部", "tag": "超买"},
|
||||
],
|
||||
"top_list": [
|
||||
{"trade_date": "20240902", "ts_code": "600000.SH", "name": "浦发银行", "pct_change": 9.95, "reason": "涨幅偏离值达7%"},
|
||||
],
|
||||
"ths_daily": [
|
||||
{"ts_code": "885811.TI", "trade_date": "20240902", "open": 1000, "high": 1010, "low": 990, "close": 1005, "pre_close": 995, "pct_change": 1.0, "vol": 100, "turnover_rate": 1.2},
|
||||
],
|
||||
"dc_index": [
|
||||
{"ts_code": "BK0475", "trade_date": "20240902", "name": "银行", "open": 100, "high": 101, "low": 99, "close": 100.5, "pre_close": 99.5, "pct_change": 1.0, "vol": 10, "amount": 1e8, "turnover_rate": 0.5},
|
||||
],
|
||||
"sw_daily": [
|
||||
{"ts_code": "801780.SI", "trade_date": "20240902", "name": "银行", "open": 2000, "high": 2010, "low": 1990, "close": 2005, "pct_change": 0.8, "vol": 50, "amount": 2e8},
|
||||
],
|
||||
}
|
||||
|
||||
|
||||
@@ -90,9 +66,4 @@ def fake_transport(api_name: str, params: dict, fields: str):
|
||||
start = str(params.get("start_date") or "")
|
||||
end = str(params.get("end_date") or "99999999")
|
||||
return [row for row in RAW["trade_cal"] if start <= row["cal_date"] <= end]
|
||||
rows = list(RAW.get(api_name) or [])
|
||||
if api_name == "limit_list_d":
|
||||
limit_type = str(params.get("limit_type") or "")
|
||||
if limit_type:
|
||||
rows = [row for row in rows if str(row.get("limit_type") or "") == limit_type]
|
||||
return rows
|
||||
return list(RAW.get(api_name) or [])
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -1,65 +0,0 @@
|
||||
from __future__ import annotations
|
||||
|
||||
import tempfile
|
||||
import unittest
|
||||
from pathlib import Path
|
||||
|
||||
from datahub.adapters.tushare import TushareAdapter
|
||||
from datahub.crypto import SecretVault
|
||||
from datahub.hub import Hub
|
||||
from datahub.settings import Settings
|
||||
from tests.fixtures import TRADE_DATE, fake_transport
|
||||
|
||||
|
||||
class ExtendedEodTests(unittest.TestCase):
|
||||
def setUp(self) -> None:
|
||||
self.tmp = tempfile.TemporaryDirectory()
|
||||
key = SecretVault.generate_key()
|
||||
settings = Settings(
|
||||
host="127.0.0.1",
|
||||
port=0,
|
||||
encryption_key=key,
|
||||
api_token="k" * 32,
|
||||
admin_password="StartPass1",
|
||||
tushare_token="tushare-secret",
|
||||
db_path=Path(self.tmp.name) / "hub.db",
|
||||
backup_dir=Path(self.tmp.name) / "backups",
|
||||
scheduler_enabled=False,
|
||||
quality={"daily_row_ratio": 0.5, "null_rate_max": 0.5, "list_limit_default": 5000, "list_limit_max": 5000},
|
||||
)
|
||||
adapter = TushareAdapter("tushare-secret", transport=fake_transport)
|
||||
self.hub = Hub(settings, adapter=adapter)
|
||||
self.hub.pipeline.ingest_reference(TRADE_DATE)
|
||||
for dataset in ("daily", "valuation", "moneyflow", "auction", "index_daily"):
|
||||
self.hub.pipeline.run_dataset(dataset, TRADE_DATE)
|
||||
|
||||
def tearDown(self) -> None:
|
||||
self.hub.stop()
|
||||
self.tmp.cleanup()
|
||||
|
||||
def test_extended_soft_datasets_publish_and_serve(self) -> None:
|
||||
results = self.hub.pipeline.run_extended_soft(
|
||||
("limit_events", "popularity", "dragon_tiger", "sector_daily"),
|
||||
TRADE_DATE,
|
||||
)
|
||||
for name in ("limit_events", "popularity", "dragon_tiger", "sector_daily"):
|
||||
self.assertEqual(results[name]["state"], "published", results[name])
|
||||
api = self.hub.api
|
||||
limits = api.handle("/v1/limit-events", {"date": [TRADE_DATE]})
|
||||
self.assertGreaterEqual(len(limits["data"]), 1)
|
||||
self.assertEqual(limits["meta"]["tier"], "official")
|
||||
pop = api.handle("/v1/popularity", {"date": [TRADE_DATE], "source": ["ths"]})
|
||||
self.assertEqual(pop["data"][0]["source"], "ths")
|
||||
lhb = api.handle("/v1/dragon-tiger", {"date": [TRADE_DATE]})
|
||||
self.assertEqual(lhb["data"][0]["hm_name"], "测试游资")
|
||||
# hub stores 万元→元
|
||||
self.assertEqual(lhb["data"][0]["buy_amount"], 10_000_000.0)
|
||||
sectors = api.handle("/v1/sectors", {"date": [TRADE_DATE], "family": ["ths"]})
|
||||
self.assertEqual(sectors["data"][0]["family"], "ths")
|
||||
status = api.handle("/v1/datasets/status", {"date": [TRADE_DATE]})
|
||||
names = {item["dataset"] for item in status["data"]}
|
||||
self.assertTrue({"limit_events", "popularity", "dragon_tiger", "sector_daily"} <= names)
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
unittest.main()
|
||||
@@ -19,17 +19,11 @@ class LayoutTests(unittest.TestCase):
|
||||
def test_reserved_adapters_present(self) -> None:
|
||||
from datahub.adapters import RESERVED
|
||||
|
||||
for name in ("ths", "xgb", "akshare", "ifind"):
|
||||
for name in ("eastmoney", "tencent", "ths", "xgb", "akshare", "ifind"):
|
||||
self.assertIn(name, RESERVED)
|
||||
probe = RESERVED[name].probe()
|
||||
self.assertEqual(probe["state"], "reserved")
|
||||
self.assertFalse(probe["configured"])
|
||||
for name in ("eastmoney", "tencent"):
|
||||
self.assertIn(name, RESERVED)
|
||||
probe = RESERVED[name].probe()
|
||||
# Live free adapters: probe may be ok/error/empty depending on network.
|
||||
self.assertIn(probe["state"], {"ok", "empty", "error"})
|
||||
self.assertTrue(probe["configured"])
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
|
||||
@@ -1,176 +0,0 @@
|
||||
from __future__ import annotations
|
||||
|
||||
import tempfile
|
||||
import unittest
|
||||
from pathlib import Path
|
||||
from unittest.mock import patch
|
||||
|
||||
from datahub.adapters.base import AdapterError
|
||||
from datahub.adapters.eastmoney import HIS_TRENDS_URL, TRENDS_URL, EastmoneyAdapter
|
||||
from datahub.db import HubDB
|
||||
from datahub.realtime_serve import fetch_intraday
|
||||
from datahub.serving import ApiError, V1API
|
||||
from datahub.timeutil import now_shanghai, yyyymmdd
|
||||
|
||||
|
||||
class FakeEastmoney(EastmoneyAdapter):
|
||||
def __init__(self) -> None:
|
||||
super().__init__(timeout=2)
|
||||
self.urls: list[str] = []
|
||||
|
||||
def _get_json(self, url, params, referer):
|
||||
self.urls.append(f"{url}|{params.get('ndays')}")
|
||||
if url == TRENDS_URL:
|
||||
return {"data": {"name": "中国平安", "code": "601318", "preClose": 56.36, "trends": []}}
|
||||
if url == HIS_TRENDS_URL:
|
||||
return {
|
||||
"data": {
|
||||
"name": "中国平安",
|
||||
"code": "601318",
|
||||
"preClose": 55.8,
|
||||
"trends": [
|
||||
"2026-09-07 09:30,55.80,55.90,56.00,55.70,100,5580.00,55.900",
|
||||
"2026-09-07 15:00,56.10,56.20,56.30,56.00,200,11240.00,56.150",
|
||||
"2026-09-08 09:30,0,0,0,0,0,0.00,0",
|
||||
],
|
||||
}
|
||||
}
|
||||
raise AdapterError(f"unexpected url {url}")
|
||||
|
||||
|
||||
class EastmoneyIntradayLookbackTests(unittest.TestCase):
|
||||
def test_empty_today_uses_latest_available_session(self):
|
||||
adapter = FakeEastmoney()
|
||||
payload = adapter.fetch_intraday("601318.SH")
|
||||
self.assertEqual(adapter.urls, [f"{TRENDS_URL}|1", f"{TRENDS_URL}|5", f"{HIS_TRENDS_URL}|5"])
|
||||
self.assertEqual(payload["trade_date"], "2026-09-07")
|
||||
self.assertEqual([point["time"] for point in payload["points"]], ["09:30", "15:00"])
|
||||
self.assertEqual(payload["points"][0]["close"], 55.9)
|
||||
|
||||
def test_preferred_date_keeps_that_session(self):
|
||||
adapter = FakeEastmoney()
|
||||
payload = adapter.fetch_intraday("601318.SH", "20260907")
|
||||
self.assertEqual(payload["trade_date"], "2026-09-07")
|
||||
self.assertEqual(len(payload["points"]), 2)
|
||||
|
||||
|
||||
class IntradayLkgTests(unittest.TestCase):
|
||||
def setUp(self) -> None:
|
||||
self.tmp = tempfile.TemporaryDirectory()
|
||||
self.db = HubDB(Path(self.tmp.name) / "hub.db")
|
||||
|
||||
def tearDown(self) -> None:
|
||||
self.tmp.cleanup()
|
||||
|
||||
def test_source_failure_returns_last_known_good(self):
|
||||
from datahub.realtime_serve import _envelope, _write_cache
|
||||
|
||||
payload = _envelope(
|
||||
{
|
||||
"entity_type": "stock",
|
||||
"ts_code": "601318.SH",
|
||||
"trade_date": "2026-09-07",
|
||||
"previous_close": 55.8,
|
||||
"points": [{"date": "2026-09-07", "time": "09:30", "close": 55.9}],
|
||||
},
|
||||
{
|
||||
"tier": "provisional",
|
||||
"trade_date": "20260907",
|
||||
"source": "eastmoney:trends2",
|
||||
"stale": False,
|
||||
},
|
||||
)
|
||||
_write_cache(self.db, "intraday:601318.SH:today", payload, 20, "eastmoney:trends2")
|
||||
self.db.execute(
|
||||
"UPDATE rt_cache SET expires_at = ? WHERE cache_key = ?",
|
||||
("2000-01-01T00:00:00+08:00", "intraday:601318.SH:today"),
|
||||
)
|
||||
|
||||
with patch("datahub.realtime_serve.EastmoneyAdapter") as mocked:
|
||||
mocked.return_value.fetch_intraday.side_effect = AdapterError("down")
|
||||
recovered = fetch_intraday(self.db, "601318.SH")
|
||||
self.assertTrue(recovered["meta"]["stale"])
|
||||
self.assertEqual(recovered["data"]["points"][0]["close"], 55.9)
|
||||
|
||||
def test_source_failure_without_lkg_raises(self):
|
||||
with patch("datahub.realtime_serve.EastmoneyAdapter") as mocked:
|
||||
mocked.return_value.fetch_intraday.side_effect = AdapterError("down")
|
||||
with self.assertRaises(Exception) as ctx:
|
||||
fetch_intraday(self.db, "000001.SZ")
|
||||
self.assertIn("intraday unavailable", str(ctx.exception))
|
||||
|
||||
|
||||
class ServingIntradayDateTests(unittest.TestCase):
|
||||
def setUp(self) -> None:
|
||||
self.tmp = tempfile.TemporaryDirectory()
|
||||
self.db = HubDB(Path(self.tmp.name) / "hub.db")
|
||||
self.api = V1API(self.db, pipeline=None, settings=None)
|
||||
|
||||
def tearDown(self) -> None:
|
||||
self.tmp.cleanup()
|
||||
|
||||
def _assert_usable_intraday(self, payload: dict) -> None:
|
||||
data = payload["data"]
|
||||
points = [point for point in data.get("points") or [] if float(point.get("close") or 0) > 0]
|
||||
self.assertGreaterEqual(len(points), 1)
|
||||
self.assertTrue(str(data.get("trade_date") or ""))
|
||||
self.assertFalse((payload.get("meta") or {}).get("stale"))
|
||||
|
||||
def test_serving_omitted_or_empty_date_uses_today_and_returns_points(self) -> None:
|
||||
today = yyyymmdd(now_shanghai())
|
||||
omitted = self.api.handle("/v1/intraday/points", {"code": ["601318"]})
|
||||
empty = self.api.handle("/v1/intraday/points", {"code": ["601318"], "date": [""]})
|
||||
explicit = self.api.handle("/v1/intraday/points", {"code": ["601318"], "date": [today]})
|
||||
self._assert_usable_intraday(omitted)
|
||||
self._assert_usable_intraday(empty)
|
||||
self._assert_usable_intraday(explicit)
|
||||
self.assertEqual(omitted["data"]["trade_date"], empty["data"]["trade_date"])
|
||||
self.assertEqual(explicit["data"]["trade_date"], omitted["data"]["trade_date"])
|
||||
|
||||
def test_serving_normalizes_empty_date_to_today_and_keeps_history(self) -> None:
|
||||
today = yyyymmdd(now_shanghai())
|
||||
captured: list[str] = []
|
||||
|
||||
def fake_fetch(db, code, date=""):
|
||||
captured.append(date)
|
||||
return {
|
||||
"schema_version": 1,
|
||||
"data": {
|
||||
"trade_date": f"{date[:4]}-{date[4:6]}-{date[6:8]}",
|
||||
"points": [{"date": f"{date[:4]}-{date[4:6]}-{date[6:8]}", "time": "09:30", "close": 55.9}],
|
||||
},
|
||||
"meta": {"stale": False, "trade_date": date},
|
||||
}
|
||||
|
||||
with patch("datahub.realtime_serve.fetch_intraday", side_effect=fake_fetch):
|
||||
omitted = self.api.handle("/v1/intraday/points", {"code": ["601318"]})
|
||||
empty = self.api.handle("/v1/intraday/points", {"code": ["601318"], "date": [" "]})
|
||||
history = self.api.handle("/v1/intraday/points", {"code": ["601318"], "date": ["20260907"]})
|
||||
self.assertEqual(captured, [today, today, "20260907"])
|
||||
self.assertEqual(omitted["data"]["trade_date"], f"{today[:4]}-{today[4:6]}-{today[6:8]}")
|
||||
self.assertEqual(empty["data"]["trade_date"], omitted["data"]["trade_date"])
|
||||
self.assertEqual(history["data"]["trade_date"], "2026-09-07")
|
||||
|
||||
def test_serving_invalid_date_is_invalid_argument(self) -> None:
|
||||
with self.assertRaises(ApiError) as ctx:
|
||||
self.api.handle("/v1/intraday/points", {"code": ["601318"], "date": ["not-a-date"]})
|
||||
self.assertEqual(ctx.exception.code, "INVALID_ARGUMENT")
|
||||
self.assertIn("invalid trade_date", ctx.exception.message)
|
||||
|
||||
def test_serving_missing_code_is_invalid_argument(self) -> None:
|
||||
with self.assertRaises(ApiError) as ctx:
|
||||
self.api.handle("/v1/intraday/points", {"date": [yyyymmdd(now_shanghai())]})
|
||||
self.assertEqual(ctx.exception.code, "INVALID_ARGUMENT")
|
||||
self.assertIn("code is required", ctx.exception.message)
|
||||
|
||||
def test_serving_no_data_keeps_source_unavailable(self) -> None:
|
||||
with patch("datahub.realtime_serve.EastmoneyAdapter") as mocked:
|
||||
mocked.return_value.fetch_intraday.side_effect = AdapterError("No intraday chart data returned")
|
||||
with self.assertRaises(ApiError) as ctx:
|
||||
self.api.handle("/v1/intraday/points", {"code": ["000001"]})
|
||||
self.assertEqual(ctx.exception.code, "SOURCE_UNAVAILABLE")
|
||||
self.assertIn("intraday unavailable", ctx.exception.message)
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
unittest.main()
|
||||
@@ -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()
|
||||
Reference in New Issue
Block a user