Compare commits
6
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
605f97e5df | ||
|
|
16ba83ec01 | ||
|
|
1c740a9d48 | ||
|
|
75c2e33b68 | ||
|
|
32f565ecb9 | ||
|
|
16841e9ae3 |
@@ -7,6 +7,9 @@ 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,11 +21,14 @@ from backend.data.providers.tushare_client import TushareClient
|
||||
|
||||
LOGGER = logging.getLogger("xiaobai.datahub")
|
||||
ShadowSink = Callable[[dict[str, Any]], None]
|
||||
EMPTY_FAIL_DATASETS = {"stocks", "daily", "index_daily", "valuation", "moneyflow", "auction"}
|
||||
EMPTY_FAIL_DATASETS = {
|
||||
"stocks", "daily", "index_daily", "valuation", "moneyflow", "auction",
|
||||
"limit_events", "sector_daily",
|
||||
}
|
||||
|
||||
|
||||
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
|
||||
|
||||
@@ -96,8 +99,8 @@ class DatahubBridge:
|
||||
legacy_query: Callable[..., list[dict[str, Any]]],
|
||||
) -> list[dict[str, Any]]:
|
||||
dataset = API_TO_DATASET.get(api_name)
|
||||
# 问天允许后续纳入 datahub;首批只读接入仍保持旧链路,避免误切。
|
||||
if not dataset or self.heaven_guard():
|
||||
# 问天按实际数据依赖接入:已映射到 hub 的 API 跟随开关;未映射的继续旧链路。
|
||||
if not dataset:
|
||||
return legacy_query(api_name, params, fields)
|
||||
flags = self.settings.flags(dataset)
|
||||
if not flags.read and not flags.shadow:
|
||||
@@ -108,7 +111,7 @@ class DatahubBridge:
|
||||
hub_error: str | None = None
|
||||
hub_canonical: list[dict[str, Any]] = []
|
||||
try:
|
||||
response = self._fetch_dataset(dataset, params or {})
|
||||
response = self._fetch_dataset(dataset, params or {}, api_name=api_name)
|
||||
hub_canonical = self._extract_rows(dataset, response, params or {})
|
||||
hub_rows = to_native_rows(dataset, hub_canonical)
|
||||
hub_meta = dict(response.meta)
|
||||
@@ -122,10 +125,12 @@ class DatahubBridge:
|
||||
legacy_rows = legacy_query(api_name, params, fields)
|
||||
except Exception as exc:
|
||||
if flags.read and hub_rows is not None and hub_error is None:
|
||||
self._emit_shadow(compare_rows(dataset, [], hub_canonical, hub_meta, self._error_text(exc)))
|
||||
self._emit_shadow(
|
||||
compare_rows(dataset, [], hub_canonical, hub_meta, self._error_text(exc), fields)
|
||||
)
|
||||
return project_fields(hub_rows, fields)
|
||||
raise
|
||||
self._emit_shadow(compare_rows(dataset, legacy_rows, hub_canonical, hub_meta, hub_error))
|
||||
self._emit_shadow(compare_rows(dataset, legacy_rows, hub_canonical, hub_meta, hub_error, fields))
|
||||
if flags.read and hub_rows is not None and hub_error is None:
|
||||
return project_fields(hub_rows, fields)
|
||||
return legacy_rows
|
||||
@@ -134,7 +139,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]) -> DatahubResponse:
|
||||
def _fetch_dataset(self, dataset: str, params: dict[str, Any], api_name: str = "") -> 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)
|
||||
@@ -151,6 +156,10 @@ 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] = {}
|
||||
@@ -165,6 +174,23 @@ 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,6 +60,27 @@ 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})
|
||||
|
||||
|
||||
@@ -5,6 +5,7 @@ from typing import Any
|
||||
from backend.data.datahub.native import SCALE_FIELDS, row_key, to_canonical_row, yyyymmdd
|
||||
|
||||
NUMERIC_TOLERANCE = 1e-4
|
||||
CANONICAL_ALIASES = {"volume": "vol"}
|
||||
|
||||
|
||||
def compare_rows(
|
||||
@@ -13,8 +14,10 @@ def compare_rows(
|
||||
hub_rows: list[dict[str, Any]] | None,
|
||||
hub_meta: dict[str, Any] | None = None,
|
||||
hub_error: str | None = None,
|
||||
fields: str = "",
|
||||
) -> dict[str, Any]:
|
||||
hub = hub_rows or []
|
||||
requested = _requested_fields(fields)
|
||||
legacy_map = {row_key(dataset, row): row for row in legacy_rows}
|
||||
hub_map = {row_key(dataset, _align_hub_row(row)): row for row in hub}
|
||||
missing_hub = sorted(key for key in legacy_map if key not in hub_map)
|
||||
@@ -26,7 +29,7 @@ def compare_rows(
|
||||
hub_row = hub_map.get(key)
|
||||
if hub_row is None:
|
||||
continue
|
||||
field_report = _compare_fields(dataset, legacy, hub_row)
|
||||
field_report = _compare_fields(dataset, legacy, hub_row, requested)
|
||||
if field_report["unit_conversion"]:
|
||||
unit_conversion.append({"key": list(key), "fields": field_report["unit_conversion"]})
|
||||
if field_report["value_diff"]:
|
||||
@@ -53,6 +56,7 @@ def compare_rows(
|
||||
"published_at": (hub_meta or {}).get("published_at"),
|
||||
"trade_date": yyyymmdd((hub_meta or {}).get("trade_date")),
|
||||
"hub_error": hub_error,
|
||||
"fields_compared": sorted(requested) if requested is not None else None,
|
||||
"equal": (
|
||||
not hub_error
|
||||
and not missing_hub
|
||||
@@ -71,13 +75,36 @@ def _align_hub_row(row: dict[str, Any]) -> dict[str, Any]:
|
||||
return aligned
|
||||
|
||||
|
||||
def _compare_fields(dataset: str, legacy: dict[str, Any], hub: dict[str, Any]) -> dict[str, list[dict[str, Any]]]:
|
||||
def _requested_fields(fields: str) -> list[str] | None:
|
||||
"""Fields the website actually asked for; None means "no projection"."""
|
||||
keys = [item.strip() for item in str(fields or "").split(",") if item.strip()]
|
||||
if not keys:
|
||||
return None
|
||||
seen: list[str] = []
|
||||
for key in keys:
|
||||
canonical = CANONICAL_ALIASES.get(key, key)
|
||||
if canonical not in seen:
|
||||
seen.append(canonical)
|
||||
return seen
|
||||
|
||||
|
||||
def _compare_fields(
|
||||
dataset: str,
|
||||
legacy: dict[str, Any],
|
||||
hub: dict[str, Any],
|
||||
requested: list[str] | None = None,
|
||||
) -> dict[str, list[dict[str, Any]]]:
|
||||
canonical_legacy = to_canonical_row(dataset, legacy)
|
||||
hub_canonical = _hub_canonical(dataset, hub)
|
||||
native_hub = _align_hub_row(hub)
|
||||
value_diff: list[dict[str, Any]] = []
|
||||
unit_conversion: list[dict[str, Any]] = []
|
||||
keys = (set(canonical_legacy) | set(hub_canonical)) - {"batch_id", "updated_at", "volume"}
|
||||
if requested is not None:
|
||||
# Compare only what the website asked for. Extra hub columns are
|
||||
# transport detail, not business differences; a requested field still
|
||||
# alarms when it is missing or holds a different value.
|
||||
keys = set(requested) - {"batch_id", "updated_at", "volume"}
|
||||
scales = SCALE_FIELDS.get(dataset) or {}
|
||||
for field in sorted(keys):
|
||||
left = canonical_legacy.get(field)
|
||||
|
||||
@@ -17,6 +17,13 @@ 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 = {
|
||||
@@ -35,6 +42,16 @@ 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,
|
||||
},
|
||||
}
|
||||
|
||||
|
||||
@@ -67,6 +84,16 @@ 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
|
||||
|
||||
|
||||
@@ -96,6 +123,30 @@ 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,6 +17,13 @@ DATASETS = (
|
||||
"valuation",
|
||||
"moneyflow",
|
||||
"auction",
|
||||
"limit_events",
|
||||
"popularity",
|
||||
"dragon_tiger",
|
||||
"sector_daily",
|
||||
"quotes",
|
||||
"index_quotes",
|
||||
"intraday",
|
||||
"status",
|
||||
)
|
||||
|
||||
@@ -28,6 +35,13 @@ 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",
|
||||
}
|
||||
|
||||
|
||||
@@ -13,6 +13,13 @@
|
||||
"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 }
|
||||
}
|
||||
}
|
||||
|
||||
@@ -201,6 +201,87 @@ class DatahubBridgeTests(unittest.TestCase):
|
||||
skew = compare_rows("daily", [LEGACY_DAILY], [HUB_DAILY], {"stale": False, "staleness_seconds": 12})
|
||||
self.assertTrue(skew["time_skew"])
|
||||
|
||||
def test_shadow_extra_hub_columns_are_not_false_diffs_when_projected(self) -> None:
|
||||
hub_full = {**HUB_DAILY, "adj_factor": 1.1}
|
||||
legacy_close_only = {k: LEGACY_DAILY[k] for k in ("ts_code", "trade_date", "close")}
|
||||
report = compare_rows(
|
||||
"daily", [legacy_close_only], [hub_full],
|
||||
{"stale": False, "staleness_seconds": 0},
|
||||
fields="ts_code,trade_date,close",
|
||||
)
|
||||
self.assertTrue(report["equal"])
|
||||
self.assertEqual(report["value_diff_count"], 0)
|
||||
self.assertEqual(report["fields_compared"], ["close", "trade_date", "ts_code"])
|
||||
# without projection the same pair shows the historic false diff
|
||||
unprojected = compare_rows("daily", [legacy_close_only], [hub_full])
|
||||
self.assertFalse(unprojected["equal"])
|
||||
|
||||
legacy_stocks = {"ts_code": "600000.SH", "name": "浦发银行"}
|
||||
hub_stocks = {
|
||||
"ts_code": "600000.SH", "symbol": "600000", "name": "浦发银行", "area": "上海",
|
||||
"industry": "银行", "market": "主板", "list_status": "L", "list_date": "19991110",
|
||||
}
|
||||
stocks = compare_rows("stocks", [legacy_stocks], [hub_stocks], {}, fields="ts_code,name")
|
||||
self.assertTrue(stocks["equal"])
|
||||
|
||||
legacy_cal = {"cal_date": "20240902", "is_open": 1}
|
||||
hub_cal = {
|
||||
"cal_date": "20240902", "is_open": True,
|
||||
"pretrade_date": "20240830", "prev_open": "20240830",
|
||||
}
|
||||
calendar = compare_rows(
|
||||
"calendar", [legacy_cal], [hub_cal], {}, fields="cal_date,is_open"
|
||||
)
|
||||
self.assertTrue(calendar["equal"])
|
||||
|
||||
def test_shadow_projection_still_alarms_on_requested_field_problems(self) -> None:
|
||||
hub_missing_field = {k: v for k, v in HUB_DAILY.items() if k != "close"}
|
||||
legacy_close_only = {k: LEGACY_DAILY[k] for k in ("ts_code", "trade_date", "close")}
|
||||
lost = compare_rows(
|
||||
"daily", [legacy_close_only], [hub_missing_field], fields="ts_code,trade_date,close"
|
||||
)
|
||||
self.assertFalse(lost["equal"])
|
||||
self.assertEqual(lost["value_diff_count"], 1)
|
||||
|
||||
changed = compare_rows(
|
||||
"daily", [legacy_close_only], [{**HUB_DAILY, "close": 99.0}],
|
||||
fields="ts_code,trade_date,close",
|
||||
)
|
||||
self.assertFalse(changed["equal"])
|
||||
self.assertEqual(changed["value_diff_count"], 1)
|
||||
self.assertEqual(changed["value_diffs"][0]["fields"][0]["field"], "close")
|
||||
|
||||
gone = compare_rows("daily", [LEGACY_DAILY], [], fields="ts_code,trade_date,close")
|
||||
self.assertEqual(gone["missing_hub_count"], 1)
|
||||
self.assertFalse(gone["equal"])
|
||||
|
||||
unit = compare_rows(
|
||||
"daily", [LEGACY_DAILY], [{**HUB_DAILY, "amount": 2000.0, "volume": 1000.0}],
|
||||
fields="ts_code,trade_date,vol,amount",
|
||||
)
|
||||
self.assertGreater(unit["unit_conversion_count"], 0)
|
||||
self.assertFalse(unit["equal"])
|
||||
|
||||
def test_bridge_shadow_report_uses_website_request_fields(self) -> None:
|
||||
hub_full = {**HUB_DAILY, "adj_factor": 1.1}
|
||||
legacy_close_only = {k: LEGACY_DAILY[k] for k in ("ts_code", "trade_date", "close", "vol", "amount")}
|
||||
reports: list[dict[str, Any]] = []
|
||||
client = FakeClient(
|
||||
response=DatahubResponse(
|
||||
data=[hub_full],
|
||||
meta={"tier": "official", "trade_date": "20240902", "stale": False, "staleness_seconds": 0},
|
||||
)
|
||||
)
|
||||
wrapped = DatahubAwareTushareClient(
|
||||
FakeLegacy([legacy_close_only]),
|
||||
DatahubBridge(flags(daily=(False, True)), client, shadow_sink=reports.append),
|
||||
)
|
||||
rows = wrapped.query("daily", {"trade_date": "20240902"}, "ts_code,trade_date,close,vol,amount")
|
||||
self.assertEqual(rows[0]["close"], 10.20)
|
||||
self.assertEqual(rows[0]["vol"], 1000.0)
|
||||
self.assertTrue(reports[0]["equal"])
|
||||
self.assertEqual(reports[0]["matched"], 1)
|
||||
|
||||
def test_native_roundtrip_matches_known_scales(self) -> None:
|
||||
native = to_native_row("daily", HUB_DAILY)
|
||||
self.assertEqual(native["vol"], 1000.0)
|
||||
@@ -209,8 +290,8 @@ class DatahubBridgeTests(unittest.TestCase):
|
||||
self.assertEqual(canonical["vol"], 100000.0)
|
||||
self.assertEqual(canonical["amount"], 2000000.0)
|
||||
|
||||
def test_heaven_keeps_legacy_on_first_batch_even_when_read_flag_is_on(self) -> None:
|
||||
"""问天未永久冻结;首批只读接入仍走旧链路,后续迁移可以纳入。"""
|
||||
def test_heaven_can_use_hub_when_dataset_flag_is_on(self) -> None:
|
||||
"""问天按数据依赖接入:已映射 API 跟随开关,不再整栈强制旧链路。"""
|
||||
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()
|
||||
@@ -221,7 +302,8 @@ class DatahubBridgeTests(unittest.TestCase):
|
||||
)
|
||||
rows = wrapped.query("daily", {"trade_date": "20240902"}, "amount")
|
||||
self.assertEqual(rows[0]["amount"], 2000.0)
|
||||
self.assertEqual(client.paths, [])
|
||||
self.assertEqual(client.paths, ["/v1/bars/daily"])
|
||||
self.assertEqual(legacy.calls, [])
|
||||
|
||||
def test_status_flag_does_not_run_when_off_and_falls_back_when_on(self) -> None:
|
||||
off = DatahubBridge(flags(), FakeClient(error=DatahubError("UNAVAILABLE", "down")))
|
||||
|
||||
@@ -6,11 +6,12 @@
|
||||
## 做什么
|
||||
|
||||
- SQLite WAL `datahub.db`,容器名 `xiaobai-datahub`,端口 `8766`
|
||||
- Tushare 盘后正式数据:交易日历、股票主档、daily、daily_basic、adj_factor、index_daily、moneyflow、stk_auction
|
||||
- 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_* 正式表
|
||||
- 暂存 → 校验 → 整批原子发布 → 可回滚
|
||||
- `/v1` 稳定接口(`X-Datahub-Token`)
|
||||
- `/admin/` 最小管理后台(总览 / 数据源 / 调度 / 发布 / 数据集 / 审计)
|
||||
- 东财/腾讯/同花顺/选股宝/AKShare/iFinD 适配器位已预留,本阶段不拉实时源
|
||||
- 同花顺/选股宝/AKShare/iFinD 适配器位仍预留;东财/腾讯已接入盘中观察
|
||||
|
||||
## 单位口径(相对现站)
|
||||
|
||||
@@ -83,6 +84,16 @@ python -m datahub history-backfill
|
||||
|
||||
`hub-quality.config.json` 的 `field_gates` 按数据集配置关键字段:非空率下限(支持按字段覆盖,如 `dv_ttm` 合法高空值)、非有限值比例上限、以及相对上一已发布批次的非空率塌陷保护。字段大面积为空的批次会被拒绝发布、保留上一份正常正式数据,失败原因逐字段写入 `batches.error` / `quality_json`。被拒后数据集仍视为缺失,盘后自动重试(HEL-435 机制)会继续尝试直到成功或截止。配置对任意数据集生效,不写死单日或单字段。
|
||||
|
||||
## 整批原子发布(release group)
|
||||
|
||||
盘后发布/重发(eod_a、eod_retry、`eod-refresh`、跨数据集重发)不再逐数据集各自切换,而是走整批原子可见机制:
|
||||
|
||||
- 一致性边界:日 K、估值、资金流、竞价同属 A 组整批;指数日 K 为 B 组;当日股票主档快照随 A 组一同切换(主档 `stock_master` 的 UPSERT 与快照发布同一事务,不会出现主档先行/滞后)。
|
||||
- 流程:组内全部成员先在暂存表完成拉取、字段质量门、覆盖检查和跨数据集交叉校验(`cross_gates` 配置 ts_code 覆盖重叠率下限),全部达标后才在**一个 SQLite 事务**里复制正式表并翻转全部 `publications` 指针。
|
||||
- 任一成员失败(拉取失败、质量门拒绝、交叉校验不过、切换事务中断)→ 整批不切换,对外继续提供上一份完整正式版本,失败原因写入 `batches.error` 与 `audit_log`(`action=release-group`),等待晚间自动重试。
|
||||
- 读取侧任何时刻只会看到"旧完整版本"或"新完整版本":发布指针在单事务内统一翻转,容器重启/事务中断自动回滚,不暴露字段残缺或跨数据集混合版本。
|
||||
- 幂等:仅当一致性边界内全部成员都已发布时才整组跳过;边界内任有缺失则整组重暂存后统一切换,避免旧批次与新批次混在同一次重发中。重复执行、并发重试不会在完整边界已就绪时生成重复批次(调度器另有 EOD 互斥锁)。
|
||||
|
||||
## 股票主档每日刷新与发布
|
||||
|
||||
交易日 20:00 与 23:10(`stocks_refresh_times` 可配)自动刷新股票主档并发布版本化快照(`eod_stocks` + `publications.dataset='stocks'`),覆盖当日新上市、证券简称变化和上市首日 N/C 前缀摘除;无变化则跳过,重复执行幂等。`/v1/stocks` 从最新已发布快照提供数据并带 `batch_id` / `published_at`;`/v1/datasets/status` 同步展示 stocks 状态。
|
||||
@@ -105,11 +116,13 @@ python -m datahub moneyflow-backfill # --trading-days 60 --end-date --for
|
||||
|
||||
```bash
|
||||
cd xiaobai-datahub
|
||||
python -m datahub eod-refresh --trade-date 20260904 # 只补缺失数据集
|
||||
python -m datahub eod-refresh --trade-date 20260904 # 补不完整的 A/B 边界
|
||||
python -m datahub eod-refresh --trade-date 20260904 --force --dataset valuation
|
||||
# 强制重取重发:仍走全部质量门,生成新批次,上一批次保留可回滚
|
||||
# --force 按一致性边界整组重发:valuation/daily/moneyflow/auction/stocks → A 组;
|
||||
# index_daily → B 组。不可再单独切换某一个正式数据集。
|
||||
```
|
||||
|
||||
管理后台「补数」对盘后正式数据集同样走 `force_republish_boundary`,不会绕过 A/B 整批边界。
|
||||
|
||||
## 备份
|
||||
|
||||
|
||||
@@ -269,7 +269,7 @@ function renderRelease(data) {
|
||||
|
||||
async function dangerous(kind, dataset) {
|
||||
const date = ($("rel-date") && $("rel-date").value) || "";
|
||||
const ds = dataset || prompt("数据集(daily / valuation / moneyflow / auction / index_daily / reference)", "daily");
|
||||
const ds = dataset || prompt("数据集(daily/valuation/moneyflow/auction/stocks→A组整批;index_daily→B组;或 reference)", "daily");
|
||||
if (!ds) return;
|
||||
const password = prompt("二次确认:输入管理密码");
|
||||
if (!password) return;
|
||||
|
||||
@@ -22,6 +22,18 @@
|
||||
"20:00",
|
||||
"23:10"
|
||||
],
|
||||
"cross_gates": [
|
||||
{
|
||||
"left": "daily",
|
||||
"right": "valuation",
|
||||
"min_key_overlap": 0.98
|
||||
},
|
||||
{
|
||||
"left": "daily",
|
||||
"right": "moneyflow",
|
||||
"min_key_overlap": 0.98
|
||||
}
|
||||
],
|
||||
"field_gates": {
|
||||
"valuation": {
|
||||
"fields": [
|
||||
|
||||
@@ -1,13 +1,13 @@
|
||||
from datahub.adapters.akshare import ADAPTER as akshare
|
||||
from datahub.adapters.eastmoney import ADAPTER as eastmoney
|
||||
from datahub.adapters.eastmoney import EastmoneyAdapter
|
||||
from datahub.adapters.ifind import ADAPTER as ifind
|
||||
from datahub.adapters.tencent import ADAPTER as tencent
|
||||
from datahub.adapters.tencent import TencentAdapter
|
||||
from datahub.adapters.ths import ADAPTER as ths
|
||||
from datahub.adapters.xgb import ADAPTER as xgb
|
||||
|
||||
RESERVED = {
|
||||
"eastmoney": eastmoney,
|
||||
"tencent": tencent,
|
||||
"eastmoney": EastmoneyAdapter(),
|
||||
"tencent": TencentAdapter(),
|
||||
"ths": ths,
|
||||
"xgb": xgb,
|
||||
"akshare": akshare,
|
||||
|
||||
@@ -1,3 +1,250 @@
|
||||
from datahub.adapters.base import ReservedAdapter
|
||||
from __future__ import annotations
|
||||
|
||||
ADAPTER = ReservedAdapter("eastmoney")
|
||||
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"
|
||||
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) -> 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
|
||||
payload = self._get_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",
|
||||
},
|
||||
referer="https://quote.eastmoney.com/",
|
||||
)
|
||||
data = payload.get("data") or {}
|
||||
points = []
|
||||
for raw in data.get("trends") or []:
|
||||
point = _parse_trend(raw)
|
||||
if point:
|
||||
points.append(point)
|
||||
if not points:
|
||||
raise AdapterError("No intraday chart data returned")
|
||||
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 _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
|
||||
return {
|
||||
"time": when.strftime("%H:%M"),
|
||||
"date": when.strftime("%Y-%m-%d"),
|
||||
"open": round4(finite_number(parts[1])),
|
||||
"close": round4(finite_number(parts[2])),
|
||||
"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])),
|
||||
}
|
||||
|
||||
@@ -1,3 +1,99 @@
|
||||
from datahub.adapters.base import ReservedAdapter
|
||||
from __future__ import annotations
|
||||
|
||||
ADAPTER = ReservedAdapter("tencent")
|
||||
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
|
||||
|
||||
@@ -11,8 +11,12 @@ 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,
|
||||
)
|
||||
@@ -31,6 +35,21 @@ 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 = {
|
||||
@@ -42,12 +61,15 @@ 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):
|
||||
@@ -85,6 +107,14 @@ 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)
|
||||
@@ -93,10 +123,67 @@ 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:
|
||||
@@ -104,6 +191,17 @@ 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,
|
||||
@@ -148,12 +246,11 @@ class TushareAdapter(MarketAdapter):
|
||||
try:
|
||||
with urllib.request.urlopen(request, timeout=self.timeout) as response:
|
||||
result = json.loads(response.read().decode("utf-8"))
|
||||
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")
|
||||
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')}"))
|
||||
data = result.get("data") or {}
|
||||
columns = data.get("fields") or []
|
||||
return [dict(zip(columns, item)) for item in data.get("items") 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]
|
||||
|
||||
@@ -6,7 +6,7 @@ from typing import Any
|
||||
from datahub.adapters import RESERVED
|
||||
from datahub.auth import AuthService
|
||||
from datahub.db import HubDB
|
||||
from datahub.pipeline import Pipeline
|
||||
from datahub.pipeline import OFFICIAL_DATASETS, STOCKS_DATASET, Pipeline
|
||||
from datahub.scheduler import Scheduler
|
||||
from datahub.serving import ApiError
|
||||
from datahub.timeutil import isoformat, now_shanghai, session_phase, yyyymmdd
|
||||
@@ -143,8 +143,22 @@ class AdminAPI:
|
||||
self._dangerous(password, confirm, f"{dataset}:{day}")
|
||||
if dataset == "reference":
|
||||
result = self.pipeline.ingest_reference(day)
|
||||
elif dataset in OFFICIAL_DATASETS or dataset == STOCKS_DATASET:
|
||||
# Manual same-day republish must rebuild the full A/B boundary.
|
||||
# Gate failures and mid-switch exceptions both surface as
|
||||
# FAILED_PRECONDITION so the admin API never leaks raw
|
||||
# transaction errors to the client.
|
||||
try:
|
||||
result = self.pipeline.force_republish_boundary(dataset, day)
|
||||
failures = self.pipeline.eod_failures(result)
|
||||
if failures:
|
||||
raise ApiError("FAILED_PRECONDITION", "; ".join(failures))
|
||||
except ApiError:
|
||||
raise
|
||||
except Exception as exc:
|
||||
raise ApiError("FAILED_PRECONDITION", str(exc)) from exc
|
||||
else:
|
||||
result = self.pipeline.run_dataset(dataset, day)
|
||||
raise ApiError("INVALID_ARGUMENT", f"unsupported backfill dataset: {dataset}")
|
||||
self.pipeline.audit(actor, "backfill", f"{dataset}:{day}", json.dumps({"ok": True}))
|
||||
return result
|
||||
|
||||
|
||||
@@ -7,7 +7,7 @@ import json
|
||||
import sys
|
||||
|
||||
from datahub.hub import build_hub
|
||||
from datahub.pipeline import OFFICIAL_DATASETS
|
||||
from datahub.pipeline import EOD_A_DATASETS, OFFICIAL_DATASETS, STOCKS_DATASET
|
||||
from datahub.settings import load_settings
|
||||
from datahub.timeutil import yyyymmdd
|
||||
|
||||
@@ -19,15 +19,15 @@ def main(argv: list[str] | None = None) -> int:
|
||||
history.add_argument("--calendar-start", default=None, help="日历起点,默认配置 calendar_start")
|
||||
history.add_argument("--index-days", type=int, default=None, help="指数回补交易日数量,默认 260")
|
||||
history.add_argument("--force", action="store_true", help="覆盖已发布的指数日期")
|
||||
refresh = sub.add_parser("eod-refresh", help="对指定交易日补跑盘后正式数据(跳过已发布数据集,仍走质量门禁)")
|
||||
refresh = sub.add_parser("eod-refresh", help="对指定交易日补跑盘后正式数据(跳过已完整发布的一致性边界,仍走质量门禁)")
|
||||
refresh.add_argument("--trade-date", default=None, help="交易日 YYYYMMDD,默认今天")
|
||||
refresh.add_argument(
|
||||
"--force", action="store_true",
|
||||
help="对 --dataset 指定的数据集强制重取重发(生成新批次,保留上一批次可回滚)",
|
||||
help="强制重发 --dataset 所属的完整一致性边界(A 组或 B 组),生成新批次并保留上一批次可回滚",
|
||||
)
|
||||
refresh.add_argument(
|
||||
"--dataset", default=None,
|
||||
help="配合 --force 使用:只强制重发该数据集(如 valuation)",
|
||||
help="配合 --force:指定边界内任一成员(如 valuation→整组 A;index_daily→整组 B)",
|
||||
)
|
||||
stocks_refresh = sub.add_parser("stocks-refresh", help="刷新股票主档并发布正式快照(幂等:无变化则跳过)")
|
||||
stocks_refresh.add_argument("--trade-date", default=None, help="交易日 YYYYMMDD,默认今天")
|
||||
@@ -54,26 +54,27 @@ def main(argv: list[str] | None = None) -> int:
|
||||
if args.command == "eod-refresh":
|
||||
day = yyyymmdd(args.trade_date) if args.trade_date else yyyymmdd()
|
||||
if args.force:
|
||||
datasets = tuple(sorted({args.dataset} & OFFICIAL_DATASETS)) if args.dataset else ()
|
||||
if args.dataset and not datasets:
|
||||
parser.error(f"unknown dataset: {args.dataset}")
|
||||
if not datasets:
|
||||
allowed = set(OFFICIAL_DATASETS) | {STOCKS_DATASET}
|
||||
if not args.dataset:
|
||||
parser.error("--force requires --dataset (e.g. --dataset valuation)")
|
||||
result = {}
|
||||
for dataset in datasets:
|
||||
result[dataset] = hub.pipeline.run_dataset(dataset, day)
|
||||
if args.dataset not in allowed:
|
||||
parser.error(f"unknown dataset: {args.dataset}")
|
||||
result = hub.pipeline.force_republish_boundary(args.dataset, day)
|
||||
boundary = "A" if args.dataset in EOD_A_DATASETS or args.dataset == STOCKS_DATASET else "B"
|
||||
else:
|
||||
result = hub.pipeline.run_eod_missing(day)
|
||||
boundary = None
|
||||
hub.pipeline.audit("cli", "eod-refresh", f"eod:{day}", json.dumps(
|
||||
{"force": bool(args.force), "dataset": args.dataset,
|
||||
{"force": bool(args.force), "dataset": args.dataset, "boundary": boundary,
|
||||
**{name: item.get("state") for name, item in result.items() if isinstance(item, dict)}},
|
||||
ensure_ascii=False,
|
||||
))
|
||||
if args.force:
|
||||
payload = {"trade_date": day, "datasets": result}
|
||||
failures = hub.pipeline.eod_failures(result)
|
||||
payload = {"trade_date": day, "boundary": boundary, "datasets": result}
|
||||
json.dump(payload, sys.stdout, ensure_ascii=False, indent=2, default=str)
|
||||
sys.stdout.write("\n")
|
||||
return 0
|
||||
return 0 if not failures else 1
|
||||
missing = hub.pipeline.missing_official_datasets(day)
|
||||
payload = {"trade_date": day, "datasets": result, "missing_after": missing}
|
||||
json.dump(payload, sys.stdout, ensure_ascii=False, indent=2, default=str)
|
||||
|
||||
@@ -0,0 +1,180 @@
|
||||
"""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,9 +7,10 @@ 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
|
||||
|
||||
SCHEMA = """
|
||||
_BASE_SCHEMA = """
|
||||
CREATE TABLE IF NOT EXISTS schema_migrations (
|
||||
version INTEGER PRIMARY KEY,
|
||||
applied_at TEXT NOT NULL
|
||||
@@ -282,6 +283,8 @@ 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"),
|
||||
@@ -289,6 +292,7 @@ DATASET_TABLES = {
|
||||
"auction": ("eod_auction", "staging_auction"),
|
||||
"index_daily": ("eod_index_bars", "staging_index_bars"),
|
||||
"stocks": ("eod_stocks", "staging_stocks"),
|
||||
**EXTENDED_DATASET_TABLES,
|
||||
}
|
||||
|
||||
|
||||
|
||||
@@ -156,6 +156,95 @@ 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
|
||||
@@ -184,6 +273,11 @@ 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,6 +9,11 @@ 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
|
||||
@@ -21,12 +26,16 @@ from datahub.timeutil import add_days, isoformat, now_shanghai, yyyymmdd
|
||||
LOGGER = get_logger()
|
||||
|
||||
HARD_DATASETS = {"daily", "valuation", "index_daily"}
|
||||
SOFT_DATASETS = {"moneyflow", "auction"}
|
||||
OFFICIAL_DATASETS = HARD_DATASETS | SOFT_DATASETS
|
||||
SOFT_DATASETS = {"moneyflow", "auction"} | EXTENDED_SOFT_DATASETS
|
||||
OFFICIAL_DATASETS = HARD_DATASETS | {"moneyflow", "auction"} # A/B retry scope unchanged
|
||||
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 = {
|
||||
@@ -80,6 +89,7 @@ STAGING_INSERT = {
|
||||
r.get("close"), r.get("pct_chg"), r.get("volume"), r.get("amount"),
|
||||
),
|
||||
),
|
||||
**EXTENDED_STAGING_INSERT,
|
||||
}
|
||||
|
||||
EOD_COPY = {
|
||||
@@ -114,6 +124,7 @@ 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,
|
||||
}
|
||||
|
||||
|
||||
@@ -133,6 +144,95 @@ def _staging_row_count(connection: Any, dataset: str, batch_id: str) -> int:
|
||||
return int(row["n"] if row is not None else 0)
|
||||
|
||||
|
||||
def _staging_count_or_raise(connection: Any, dataset: str, trade_date: str, batch_id: str) -> int:
|
||||
rows_out = _staging_row_count(connection, dataset, batch_id)
|
||||
if rows_out <= 0:
|
||||
report = {
|
||||
"rows": 0,
|
||||
"errors": [EMPTY_BATCH_ERROR],
|
||||
"warnings": [],
|
||||
"hard_fail": True,
|
||||
"soft_fail": False,
|
||||
"batch_id": batch_id,
|
||||
"dataset": dataset,
|
||||
"trade_date": trade_date,
|
||||
}
|
||||
LOGGER.warning(
|
||||
"skip official publish for empty batch",
|
||||
extra={
|
||||
"hub": {
|
||||
"dataset": dataset,
|
||||
"trade_date": trade_date,
|
||||
"batch_id": batch_id,
|
||||
"rows_out": rows_out,
|
||||
"reason": "upstream_empty",
|
||||
}
|
||||
},
|
||||
)
|
||||
raise QualityError("empty batch cannot be officially published", report)
|
||||
return rows_out
|
||||
|
||||
|
||||
def _upsert_publication(
|
||||
connection: Any,
|
||||
dataset: str,
|
||||
trade_date: str,
|
||||
batch_id: str,
|
||||
state: str,
|
||||
published_at: str,
|
||||
) -> None:
|
||||
current = connection.execute(
|
||||
"SELECT active_batch FROM publications WHERE dataset = ? AND trade_date = ?",
|
||||
(dataset, trade_date),
|
||||
).fetchone()
|
||||
prev = str(current["active_batch"]) if current else None
|
||||
connection.execute(
|
||||
"""
|
||||
INSERT INTO publications(dataset, trade_date, active_batch, prev_batch, state, published_at)
|
||||
VALUES (?, ?, ?, ?, ?, ?)
|
||||
ON CONFLICT(dataset, trade_date) DO UPDATE SET
|
||||
prev_batch=excluded.prev_batch,
|
||||
active_batch=excluded.active_batch,
|
||||
state=excluded.state,
|
||||
published_at=excluded.published_at
|
||||
""",
|
||||
(dataset, trade_date, batch_id, prev, state, published_at),
|
||||
)
|
||||
|
||||
|
||||
def _record_publication_history(
|
||||
connection: Any,
|
||||
dataset: str,
|
||||
trade_date: str,
|
||||
batch_id: str,
|
||||
published_at: str,
|
||||
quality: dict[str, Any],
|
||||
) -> None:
|
||||
max_gen = connection.execute(
|
||||
"SELECT COALESCE(MAX(generation), 0) AS g FROM publication_history WHERE dataset = ? AND trade_date = ?",
|
||||
(dataset, trade_date),
|
||||
).fetchone()
|
||||
generation = int(max_gen["g"]) + 1
|
||||
connection.execute(
|
||||
"INSERT OR REPLACE INTO publication_history(dataset, trade_date, batch_id, published_at, generation) VALUES (?,?,?,?,?)",
|
||||
(dataset, trade_date, batch_id, published_at, generation),
|
||||
)
|
||||
keep = int(quality.get("publication_generations") or 3)
|
||||
stale = connection.execute(
|
||||
"""
|
||||
SELECT batch_id FROM publication_history
|
||||
WHERE dataset = ? AND trade_date = ?
|
||||
ORDER BY generation DESC
|
||||
""",
|
||||
(dataset, trade_date),
|
||||
).fetchall()
|
||||
for row in stale[keep:]:
|
||||
connection.execute(
|
||||
"DELETE FROM publication_history WHERE dataset = ? AND trade_date = ? AND batch_id = ?",
|
||||
(dataset, trade_date, row["batch_id"]),
|
||||
)
|
||||
|
||||
|
||||
class QualityError(RuntimeError):
|
||||
def __init__(self, message: str, report: dict[str, Any]) -> None:
|
||||
super().__init__(message)
|
||||
@@ -259,11 +359,20 @@ class Pipeline:
|
||||
is published; ``force`` re-publishes unconditionally. New listings,
|
||||
renames (incl. N/C prefix removal) and status changes all flow into
|
||||
the snapshot, which carries batch_id/published_at metadata.
|
||||
|
||||
The ``stock_master`` UPSERT happens inside the same publish
|
||||
transaction as the snapshot switch — fetch / quality-gate / switch
|
||||
failures leave the master on the previous complete values.
|
||||
"""
|
||||
day = yyyymmdd(trade_date or self.clock())
|
||||
rows = self._fetch_dataset(STOCKS_DATASET, day)
|
||||
with self.db.write() as connection:
|
||||
self._upsert_stock_master(connection, rows, isoformat(self.clock()))
|
||||
try:
|
||||
rows = self._fetch_dataset(STOCKS_DATASET, day)
|
||||
except Exception as exc:
|
||||
self.audit(
|
||||
"pipeline", "stocks-refresh", f"{STOCKS_DATASET}:{day}",
|
||||
json.dumps({"state": "failed", "error": str(exc)}, ensure_ascii=False),
|
||||
)
|
||||
raise
|
||||
if not force:
|
||||
active, snapshot = self.published_stock_snapshot(day)
|
||||
if active is not None:
|
||||
@@ -280,7 +389,14 @@ class Pipeline:
|
||||
"batch_id": active,
|
||||
"rows": len(snapshot),
|
||||
}
|
||||
result = self.run_dataset(STOCKS_DATASET, day, prepared_rows=rows)
|
||||
try:
|
||||
result = self.run_dataset(STOCKS_DATASET, day, prepared_rows=rows)
|
||||
except Exception as exc:
|
||||
self.audit(
|
||||
"pipeline", "stocks-refresh", f"{STOCKS_DATASET}:{day}",
|
||||
json.dumps({"state": "failed", "error": str(exc)}, ensure_ascii=False),
|
||||
)
|
||||
raise
|
||||
self.audit(
|
||||
"pipeline", "stocks-refresh", f"{STOCKS_DATASET}:{day}",
|
||||
json.dumps({"batch_id": result["batch_id"], "rows": result["rows"]}, ensure_ascii=False),
|
||||
@@ -548,25 +664,31 @@ class Pipeline:
|
||||
return [dataset for dataset in sorted(OFFICIAL_DATASETS) if dataset not in published]
|
||||
|
||||
def run_eod_missing(self, trade_date: str) -> dict[str, Any]:
|
||||
"""Fetch/publish every official dataset still missing for the date.
|
||||
"""Republish every incomplete EOD consistency group for the date.
|
||||
|
||||
Idempotent: datasets with an existing publication are skipped, so
|
||||
repeats never overwrite the current official batch. Per-dataset
|
||||
failures are collected instead of aborting the remaining datasets.
|
||||
A-group (daily/valuation/moneyflow/auction + stocks) and B-group
|
||||
(index_daily) are separate boundaries. Within a group, either the
|
||||
whole boundary is already published (idempotent skip) or every
|
||||
member is re-staged and switched together — never fill only the
|
||||
missing members on top of older batches from an earlier partial run.
|
||||
"""
|
||||
return self._run_eod_datasets(tuple(sorted(OFFICIAL_DATASETS)), trade_date)
|
||||
|
||||
def run_eod_batch_a(self, trade_date: str) -> dict[str, Any]:
|
||||
return self._run_eod_datasets(EOD_A_DATASETS, trade_date)
|
||||
|
||||
def run_eod_batch_b(self, trade_date: str) -> dict[str, Any]:
|
||||
return self._run_eod_datasets(EOD_B_DATASETS, trade_date)
|
||||
|
||||
def _run_eod_datasets(self, datasets: tuple[str, ...], trade_date: str) -> dict[str, Any]:
|
||||
day = yyyymmdd(trade_date)
|
||||
results: dict[str, Any] = {}
|
||||
results.update(self.run_eod_batch_a(trade_date))
|
||||
results.update(self.run_eod_batch_b(trade_date))
|
||||
return results
|
||||
|
||||
def run_eod_batch_a(self, trade_date: str, force: bool = False) -> dict[str, Any]:
|
||||
return self.run_release_group(EOD_A_DATASETS, trade_date, include_stocks=True, force=force)
|
||||
|
||||
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 self.active_batch(dataset, day) is not None:
|
||||
if not force and self.active_batch(dataset, day):
|
||||
results[dataset] = {
|
||||
"dataset": dataset,
|
||||
"trade_date": day,
|
||||
@@ -575,7 +697,17 @@ class Pipeline:
|
||||
}
|
||||
continue
|
||||
try:
|
||||
results[dataset] = self.run_dataset(dataset, day)
|
||||
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,
|
||||
@@ -583,8 +715,357 @@ class Pipeline:
|
||||
"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.
|
||||
"""
|
||||
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 run_release_group(
|
||||
self,
|
||||
datasets: tuple[str, ...],
|
||||
trade_date: str,
|
||||
include_stocks: bool = False,
|
||||
force: bool = False,
|
||||
) -> dict[str, Any]:
|
||||
"""One post-market publish/republish becomes one atomic visibility flip.
|
||||
|
||||
Consistency boundary: every member (official datasets, plus the daily
|
||||
stocks snapshot when ``include_stocks``) is fetched, staged,
|
||||
field-gated and cross-validated BEFORE any reader can see it. Only
|
||||
when the whole group passes does a single SQLite transaction copy
|
||||
all staging batches to the official tables and flip every
|
||||
``publications`` row at once. Any member failure aborts the group:
|
||||
the previous complete official version keeps serving and the reason
|
||||
is recorded on the batches and in the audit log.
|
||||
|
||||
Skip is all-or-nothing for the boundary unless ``force``: if every
|
||||
official member (and stocks when required) is already published, the
|
||||
group is skipped. If any official member is still missing — or
|
||||
``force`` is set — every official member is re-staged, so a retry or
|
||||
manual republish never mixes old and new batches in one release.
|
||||
"""
|
||||
day = yyyymmdd(trade_date)
|
||||
results: dict[str, Any] = {}
|
||||
staged: dict[str, dict[str, Any]] = {}
|
||||
failure: str | None = None
|
||||
missing_official = [dataset for dataset in datasets if self.active_batch(dataset, day) is None]
|
||||
stocks_missing = include_stocks and self.active_batch(STOCKS_DATASET, day) is None
|
||||
|
||||
if not force and not missing_official and not stocks_missing:
|
||||
for dataset in datasets:
|
||||
results[dataset] = {
|
||||
"dataset": dataset,
|
||||
"trade_date": day,
|
||||
"state": "skipped",
|
||||
"reason": "already_published",
|
||||
}
|
||||
if include_stocks:
|
||||
results[STOCKS_DATASET] = {
|
||||
"dataset": STOCKS_DATASET,
|
||||
"trade_date": day,
|
||||
"state": "skipped",
|
||||
"reason": "already_published",
|
||||
}
|
||||
return results
|
||||
|
||||
# Incomplete or forced boundary → restage every official member together.
|
||||
pending = list(datasets)
|
||||
|
||||
for dataset in pending:
|
||||
if failure is not None:
|
||||
results[dataset] = {
|
||||
"dataset": dataset,
|
||||
"trade_date": day,
|
||||
"state": "aborted",
|
||||
"reason": f"release group aborted: {failure}",
|
||||
}
|
||||
continue
|
||||
try:
|
||||
staged[dataset] = self._stage_and_validate(dataset, day)
|
||||
except Exception as exc:
|
||||
failure = f"{dataset}: {exc}"
|
||||
results[dataset] = {
|
||||
"dataset": dataset,
|
||||
"trade_date": day,
|
||||
"state": "failed",
|
||||
"error": str(exc),
|
||||
}
|
||||
|
||||
# Stocks join the same switch when the official boundary is being
|
||||
# rebuilt (missing or forced), or when only the stocks snapshot is
|
||||
# still missing.
|
||||
rebuild_official = bool(force or missing_official)
|
||||
if include_stocks and failure is None and (rebuild_official or stocks_missing):
|
||||
try:
|
||||
stocks_plan = self._stage_stocks_snapshot(day, force=rebuild_official)
|
||||
except Exception as exc:
|
||||
failure = f"{STOCKS_DATASET}: {exc}"
|
||||
results[STOCKS_DATASET] = {
|
||||
"dataset": STOCKS_DATASET,
|
||||
"trade_date": day,
|
||||
"state": "failed",
|
||||
"error": str(exc),
|
||||
}
|
||||
else:
|
||||
if stocks_plan is not None:
|
||||
staged[STOCKS_DATASET] = stocks_plan
|
||||
|
||||
if failure is None and staged:
|
||||
cross_errors = self._cross_gate_errors(staged)
|
||||
if cross_errors:
|
||||
failure = "; ".join(cross_errors)
|
||||
|
||||
if failure is not None:
|
||||
for dataset, item in staged.items():
|
||||
self._abandon_batch(item["batch_id"], f"release group not switched: {failure}")
|
||||
results[dataset] = {
|
||||
"dataset": dataset,
|
||||
"trade_date": day,
|
||||
"state": "failed",
|
||||
"error": f"release group not switched: {failure}",
|
||||
"batch_id": item["batch_id"],
|
||||
}
|
||||
LOGGER.warning(
|
||||
"release group blocked, previous official version keeps serving",
|
||||
extra={
|
||||
"hub": {
|
||||
"trade_date": day,
|
||||
"datasets": sorted(staged),
|
||||
"reason": failure,
|
||||
"event": "release_group_blocked",
|
||||
}
|
||||
},
|
||||
)
|
||||
self.audit(
|
||||
"pipeline", "release-group", f"eod:{day}",
|
||||
json.dumps({"state": "failed", "reason": failure}, ensure_ascii=False),
|
||||
)
|
||||
return results
|
||||
|
||||
if staged:
|
||||
try:
|
||||
self._switch_release_group(day, staged)
|
||||
except Exception as exc:
|
||||
reason = f"release group switch failed: {exc}"
|
||||
for item in staged.values():
|
||||
self._abandon_batch(item["batch_id"], reason)
|
||||
LOGGER.warning(
|
||||
"release group switch failed, previous official version keeps serving",
|
||||
extra={
|
||||
"hub": {
|
||||
"trade_date": day,
|
||||
"datasets": sorted(staged),
|
||||
"reason": reason,
|
||||
"event": "release_group_switch_failed",
|
||||
}
|
||||
},
|
||||
)
|
||||
self.audit(
|
||||
"pipeline", "release-group", f"eod:{day}",
|
||||
json.dumps(
|
||||
{
|
||||
"state": "failed",
|
||||
"reason": reason,
|
||||
"switched": [],
|
||||
"force": bool(force),
|
||||
},
|
||||
ensure_ascii=False,
|
||||
),
|
||||
)
|
||||
raise
|
||||
for dataset, item in staged.items():
|
||||
results[dataset] = {
|
||||
"dataset": dataset,
|
||||
"trade_date": day,
|
||||
"state": item["state"],
|
||||
"batch_id": item["batch_id"],
|
||||
"rows": item["rows"],
|
||||
}
|
||||
self.audit(
|
||||
"pipeline", "release-group", f"eod:{day}",
|
||||
json.dumps(
|
||||
{"state": "ok", "switched": sorted(staged), "force": bool(force)},
|
||||
ensure_ascii=False,
|
||||
),
|
||||
)
|
||||
return results
|
||||
|
||||
def _stage_and_validate(self, dataset: str, trade_date: str, attempts: int | None = None) -> dict[str, Any]:
|
||||
"""Fetch → stage → quality-gate one member without publishing it."""
|
||||
day = yyyymmdd(trade_date)
|
||||
batch_id = self.next_batch_id(dataset, day)
|
||||
max_attempts = attempts or self.settings.max_publish_attempts
|
||||
rows: list[dict[str, Any]] = []
|
||||
self._set_batch(batch_id, dataset, day, "scheduled", 0)
|
||||
try:
|
||||
self._set_batch(batch_id, dataset, day, "fetching", 1)
|
||||
rows = retry_call(
|
||||
lambda: self._fetch_dataset(dataset, day),
|
||||
attempts=max_attempts,
|
||||
base_delay=0.05,
|
||||
sleeper=lambda _d: time.sleep(_d),
|
||||
)
|
||||
self._stage(dataset, batch_id, rows)
|
||||
self._set_batch(batch_id, dataset, day, "staged", 1, rows_in=len(rows), rows_out=len(rows))
|
||||
self._set_batch(batch_id, dataset, day, "validating", 1, rows_in=len(rows), rows_out=len(rows))
|
||||
report = self.validate(dataset, batch_id, day, rows)
|
||||
if report["hard_fail"]:
|
||||
self._reject_batch(batch_id, dataset, day, rows, report)
|
||||
raise QualityError("integrity gate failed", report)
|
||||
except RetryError as exc:
|
||||
self._set_batch(batch_id, dataset, day, "failed", max_attempts, error=str(exc), finished=True)
|
||||
raise
|
||||
except QualityError as exc:
|
||||
current = self.db.fetchone("SELECT state FROM batches WHERE batch_id = ?", (batch_id,))
|
||||
if current and current["state"] not in {"staged", "failed"}:
|
||||
self._reject_batch(batch_id, dataset, day, rows, exc.report)
|
||||
raise
|
||||
except Exception as exc:
|
||||
self._set_batch(batch_id, dataset, day, "failed", 1, error=str(exc), finished=True)
|
||||
raise
|
||||
self._set_batch(
|
||||
batch_id, dataset, day, "ready", 1, rows_in=len(rows), rows_out=len(rows), quality=report
|
||||
)
|
||||
return {
|
||||
"dataset": dataset,
|
||||
"trade_date": day,
|
||||
"batch_id": batch_id,
|
||||
"rows": len(rows),
|
||||
"quality": report,
|
||||
"state": "degraded" if report["soft_fail"] else "published",
|
||||
}
|
||||
|
||||
def _stage_stocks_snapshot(self, trade_date: str, force: bool = False) -> dict[str, Any] | None:
|
||||
"""Stage the daily stocks snapshot for a release group switch.
|
||||
|
||||
Returns None when the published snapshot is already identical to
|
||||
upstream (idempotent skip) unless ``force`` is set. The stock_master
|
||||
upsert is deferred into the group switch / publish transaction so the
|
||||
master never runs ahead of the published snapshot.
|
||||
"""
|
||||
day = yyyymmdd(trade_date)
|
||||
active, snapshot = self.published_stock_snapshot(day)
|
||||
rows = self._fetch_dataset(STOCKS_DATASET, day)
|
||||
if active is not None and not force:
|
||||
upstream = sorted(
|
||||
tuple(str(row.get(field)) for field in STOCK_SNAPSHOT_FIELDS) for row in rows
|
||||
)
|
||||
published = sorted(tuple(str(row.get(field)) for field in STOCK_SNAPSHOT_FIELDS) for row in snapshot)
|
||||
if upstream == published:
|
||||
return None
|
||||
batch_id = self.next_batch_id(STOCKS_DATASET, day)
|
||||
self._set_batch(batch_id, STOCKS_DATASET, day, "scheduled", 0)
|
||||
self._set_batch(batch_id, STOCKS_DATASET, day, "fetching", 1)
|
||||
self._stage(STOCKS_DATASET, batch_id, rows)
|
||||
self._set_batch(batch_id, STOCKS_DATASET, day, "staged", 1, rows_in=len(rows), rows_out=len(rows))
|
||||
self._set_batch(batch_id, STOCKS_DATASET, day, "validating", 1, rows_in=len(rows), rows_out=len(rows))
|
||||
report = self.validate(STOCKS_DATASET, batch_id, day, rows)
|
||||
if report["hard_fail"]:
|
||||
self._reject_batch(batch_id, STOCKS_DATASET, day, rows, report)
|
||||
raise QualityError("integrity gate failed", report)
|
||||
self._set_batch(
|
||||
batch_id, STOCKS_DATASET, day, "ready", 1, rows_in=len(rows), rows_out=len(rows), quality=report
|
||||
)
|
||||
return {
|
||||
"dataset": STOCKS_DATASET,
|
||||
"trade_date": day,
|
||||
"batch_id": batch_id,
|
||||
"rows": len(rows),
|
||||
"row_values": rows,
|
||||
"quality": report,
|
||||
"state": "degraded" if report["soft_fail"] else "published",
|
||||
}
|
||||
|
||||
def _cross_gate_errors(self, staged: dict[str, dict[str, Any]]) -> list[str]:
|
||||
"""Cross-dataset consistency checks on staged batches (交叉校验)."""
|
||||
errors: list[str] = []
|
||||
gates = self.settings.quality.get("cross_gates") or []
|
||||
for gate in gates if isinstance(gates, list) else []:
|
||||
if not isinstance(gate, dict):
|
||||
continue
|
||||
left = str(gate.get("left") or "")
|
||||
right = str(gate.get("right") or "")
|
||||
if not left or not right or left not in staged or right not in staged:
|
||||
continue
|
||||
floor = float(gate.get("min_key_overlap") or 0.98)
|
||||
left_keys = self._staging_keys(left, staged[left]["batch_id"])
|
||||
right_keys = self._staging_keys(right, staged[right]["batch_id"])
|
||||
denom = max(len(left_keys), len(right_keys))
|
||||
overlap = (len(left_keys & right_keys) / denom) if denom else 1.0
|
||||
if overlap < floor:
|
||||
errors.append(
|
||||
f"cross gate: {left} vs {right} key overlap {overlap:.4f} < {floor}"
|
||||
)
|
||||
return errors
|
||||
|
||||
def _staging_keys(self, dataset: str, batch_id: str) -> set[str]:
|
||||
table = DATASET_TABLES[dataset][1]
|
||||
rows = self.db.fetchall(
|
||||
f"SELECT DISTINCT ts_code FROM {table} WHERE batch_id = ?",
|
||||
(batch_id,),
|
||||
)
|
||||
return {str(row["ts_code"]) for row in rows}
|
||||
|
||||
def _switch_release_group(self, trade_date: str, members: dict[str, dict[str, Any]]) -> None:
|
||||
"""Single transaction: copy every member and flip every publication."""
|
||||
day = yyyymmdd(trade_date)
|
||||
published_at = isoformat(self.clock())
|
||||
with self.db.write() as connection:
|
||||
for dataset, item in members.items():
|
||||
_staging_count_or_raise(connection, dataset, day, item["batch_id"])
|
||||
for dataset, item in members.items():
|
||||
connection.execute(EOD_COPY[dataset], (item["batch_id"],))
|
||||
if dataset == STOCKS_DATASET:
|
||||
self._upsert_stock_master(connection, item["row_values"], published_at)
|
||||
if self.before_commit:
|
||||
self.before_commit()
|
||||
for dataset, item in members.items():
|
||||
_upsert_publication(connection, dataset, day, item["batch_id"], item["state"], published_at)
|
||||
_record_publication_history(
|
||||
connection, dataset, day, item["batch_id"], published_at, self.settings.quality
|
||||
)
|
||||
connection.execute(
|
||||
"UPDATE batches SET state='published', finished_at=? WHERE batch_id=?",
|
||||
(published_at, item["batch_id"]),
|
||||
)
|
||||
|
||||
def _abandon_batch(self, batch_id: str, reason: str) -> None:
|
||||
row = self.db.fetchone("SELECT dataset, trade_date FROM batches WHERE batch_id = ?", (batch_id,))
|
||||
if not row:
|
||||
return
|
||||
self._set_batch(
|
||||
batch_id, str(row["dataset"]), str(row["trade_date"]), "failed", 1,
|
||||
error=reason, finished=True,
|
||||
)
|
||||
|
||||
@staticmethod
|
||||
def eod_failures(results: dict[str, Any]) -> list[str]:
|
||||
return [
|
||||
@@ -602,7 +1083,16 @@ class Pipeline:
|
||||
)
|
||||
listed_n = int((listed or {}).get("n") or 0)
|
||||
row_n = len(rows)
|
||||
keys = [(row.get("ts_code"), row.get("trade_date")) for row in 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]
|
||||
dup = row_n - len(set(keys))
|
||||
if dup:
|
||||
errors.append(f"duplicate keys: {dup}")
|
||||
@@ -623,7 +1113,8 @@ class Pipeline:
|
||||
errors.append(EMPTY_BATCH_ERROR)
|
||||
field_report = self._field_gate(dataset, trade_date, rows, errors)
|
||||
if dataset in SOFT_DATASETS:
|
||||
hard_fail = bool(dup or bad_date or empty)
|
||||
allow_empty = dataset in {"popularity", "dragon_tiger", "moneyflow", "auction"}
|
||||
hard_fail = bool(dup or bad_date or (empty and not allow_empty))
|
||||
else:
|
||||
hard_fail = bool(errors) and (dataset in HARD_DATASETS or dataset == STOCKS_DATASET)
|
||||
report = {
|
||||
@@ -735,77 +1226,24 @@ class Pipeline:
|
||||
return batch_id, stats
|
||||
|
||||
def publish(self, dataset: str, trade_date: str, batch_id: str, state: str = "published") -> None:
|
||||
copy_sql = EOD_COPY[dataset]
|
||||
published_at = isoformat(self.clock())
|
||||
with self.db.write() as connection:
|
||||
rows_out = _staging_row_count(connection, dataset, batch_id)
|
||||
if rows_out <= 0:
|
||||
report = {
|
||||
"rows": 0,
|
||||
"errors": [EMPTY_BATCH_ERROR],
|
||||
"warnings": [],
|
||||
"hard_fail": True,
|
||||
"soft_fail": False,
|
||||
"batch_id": batch_id,
|
||||
"dataset": dataset,
|
||||
"trade_date": trade_date,
|
||||
}
|
||||
LOGGER.warning(
|
||||
"skip official publish for empty batch",
|
||||
extra={
|
||||
"hub": {
|
||||
"dataset": dataset,
|
||||
"trade_date": trade_date,
|
||||
"batch_id": batch_id,
|
||||
"rows_out": rows_out,
|
||||
"reason": "upstream_empty",
|
||||
}
|
||||
},
|
||||
)
|
||||
raise QualityError("empty batch cannot be officially published", report)
|
||||
current = connection.execute(
|
||||
"SELECT active_batch FROM publications WHERE dataset = ? AND trade_date = ?",
|
||||
(dataset, trade_date),
|
||||
).fetchone()
|
||||
prev = str(current["active_batch"]) if current else None
|
||||
connection.execute(copy_sql, (batch_id,))
|
||||
_staging_count_or_raise(connection, dataset, trade_date, batch_id)
|
||||
connection.execute(EOD_COPY[dataset], (batch_id,))
|
||||
if dataset == STOCKS_DATASET:
|
||||
staging = DATASET_TABLES[STOCKS_DATASET][1]
|
||||
stock_rows = [
|
||||
dict(row)
|
||||
for row in connection.execute(
|
||||
f"SELECT * FROM {staging} WHERE batch_id = ?",
|
||||
(batch_id,),
|
||||
).fetchall()
|
||||
]
|
||||
self._upsert_stock_master(connection, stock_rows, published_at)
|
||||
if self.before_commit:
|
||||
self.before_commit()
|
||||
connection.execute(
|
||||
"""
|
||||
INSERT INTO publications(dataset, trade_date, active_batch, prev_batch, state, published_at)
|
||||
VALUES (?, ?, ?, ?, ?, ?)
|
||||
ON CONFLICT(dataset, trade_date) DO UPDATE SET
|
||||
prev_batch=excluded.prev_batch,
|
||||
active_batch=excluded.active_batch,
|
||||
state=excluded.state,
|
||||
published_at=excluded.published_at
|
||||
""",
|
||||
(dataset, trade_date, batch_id, prev, state, published_at),
|
||||
)
|
||||
max_gen = connection.execute(
|
||||
"SELECT COALESCE(MAX(generation), 0) AS g FROM publication_history WHERE dataset = ? AND trade_date = ?",
|
||||
(dataset, trade_date),
|
||||
).fetchone()
|
||||
generation = int(max_gen["g"]) + 1
|
||||
connection.execute(
|
||||
"INSERT OR REPLACE INTO publication_history(dataset, trade_date, batch_id, published_at, generation) VALUES (?,?,?,?,?)",
|
||||
(dataset, trade_date, batch_id, published_at, generation),
|
||||
)
|
||||
keep = int(self.settings.quality.get("publication_generations") or 3)
|
||||
stale = connection.execute(
|
||||
"""
|
||||
SELECT batch_id FROM publication_history
|
||||
WHERE dataset = ? AND trade_date = ?
|
||||
ORDER BY generation DESC
|
||||
""",
|
||||
(dataset, trade_date),
|
||||
).fetchall()
|
||||
for row in stale[keep:]:
|
||||
connection.execute(
|
||||
"DELETE FROM publication_history WHERE dataset = ? AND trade_date = ? AND batch_id = ?",
|
||||
(dataset, trade_date, row["batch_id"]),
|
||||
)
|
||||
_upsert_publication(connection, dataset, trade_date, batch_id, state, published_at)
|
||||
_record_publication_history(connection, dataset, trade_date, batch_id, published_at, self.settings.quality)
|
||||
|
||||
def rollback(self, dataset: str, trade_date: str, actor: str = "admin") -> dict[str, Any]:
|
||||
trade_date = yyyymmdd(trade_date)
|
||||
|
||||
@@ -0,0 +1,188 @@
|
||||
"""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.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)
|
||||
source = "eastmoney:trends2"
|
||||
except Exception as exc:
|
||||
raise RealtimeApiError("SOURCE_UNAVAILABLE", f"intraday unavailable: {exc}") from exc
|
||||
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 _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),
|
||||
)
|
||||
@@ -48,6 +48,10 @@ 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,
|
||||
"stocks_refresh": self._stocks_refresh,
|
||||
"cleanup": self._cleanup,
|
||||
@@ -87,6 +91,10 @@ 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)),
|
||||
]
|
||||
@@ -99,7 +107,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", "stocks_refresh"} and not open_day:
|
||||
if job_id in {"eod_a", "eod_b", "eod_c", "eod_d", "eod_e", "eod_f", "stocks_refresh"} and not open_day:
|
||||
self._fired.add(key)
|
||||
continue
|
||||
self._fired.add(key)
|
||||
@@ -110,7 +118,7 @@ class Scheduler:
|
||||
try:
|
||||
self.run_job(job_id, day)
|
||||
except Exception:
|
||||
if job_id not in {"eod_a", "eod_b", "stocks_refresh"}:
|
||||
if job_id not in {"eod_a", "eod_b", "eod_c", "eod_d", "eod_e", "eod_f", "stocks_refresh"}:
|
||||
raise
|
||||
# Keep the tick alive; evening retries take over.
|
||||
LOGGER.exception("scheduled job %s failed for %s", job_id, day)
|
||||
@@ -317,6 +325,18 @@ 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)
|
||||
|
||||
|
||||
@@ -73,6 +73,20 @@ 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":
|
||||
@@ -197,9 +211,75 @@ 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")
|
||||
try:
|
||||
return fetch_intraday(self.db, code, yyyymmdd(q.get("date") or ""))
|
||||
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")
|
||||
datasets = (
|
||||
"daily", "valuation", "moneyflow", "auction", "index_daily", "stocks",
|
||||
"limit_events", "popularity", "dragon_tiger", "sector_daily",
|
||||
)
|
||||
items = []
|
||||
for dataset in datasets:
|
||||
pub = self.db.fetchone(
|
||||
@@ -244,6 +324,7 @@ 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
|
||||
@@ -264,6 +345,7 @@ 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(
|
||||
@@ -282,6 +364,9 @@ 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))]
|
||||
@@ -317,6 +402,9 @@ 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]
|
||||
|
||||
@@ -45,6 +45,30 @@ 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},
|
||||
],
|
||||
}
|
||||
|
||||
|
||||
@@ -66,4 +90,9 @@ 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]
|
||||
return list(RAW.get(api_name) or [])
|
||||
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
|
||||
|
||||
@@ -0,0 +1,408 @@
|
||||
from __future__ import annotations
|
||||
|
||||
import unittest
|
||||
from pathlib import Path
|
||||
import tempfile
|
||||
|
||||
from datahub.adapters.tushare import TushareAdapter
|
||||
from datahub.crypto import SecretVault
|
||||
from datahub.db import HubDB
|
||||
from datahub.pipeline import Pipeline
|
||||
from datahub.settings import Settings
|
||||
from datahub.serving import V1API
|
||||
from tests.fixtures import TRADE_DATE, fake_transport
|
||||
|
||||
GROUP_A = ("daily", "valuation", "moneyflow", "auction")
|
||||
|
||||
|
||||
class GroupTransport:
|
||||
"""fake_transport with per-API degradation switches for release-group tests."""
|
||||
|
||||
def __init__(self) -> None:
|
||||
self.empty: set[str] = set()
|
||||
self.keep_rows: dict[str, int] = {}
|
||||
self.stocks: list[dict] | None = None
|
||||
self.calls: list[str] = []
|
||||
|
||||
def __call__(self, api_name: str, params: dict, fields: str):
|
||||
self.calls.append(api_name)
|
||||
if api_name in self.empty:
|
||||
return []
|
||||
if api_name == "stock_basic" and self.stocks is not None:
|
||||
return [dict(row) for row in self.stocks]
|
||||
rows = fake_transport(api_name, params, fields)
|
||||
keep = self.keep_rows.get(api_name)
|
||||
if keep is not None:
|
||||
return rows[:keep]
|
||||
return rows
|
||||
|
||||
|
||||
def make_pipe(transport: GroupTransport, quality_extra: dict | None = None):
|
||||
tmp = tempfile.TemporaryDirectory()
|
||||
db = HubDB(Path(tmp.name) / "hub.db")
|
||||
adapter = TushareAdapter("test-token", transport=transport)
|
||||
quality = {
|
||||
"daily_row_ratio": 0.98,
|
||||
"null_rate_max": 0.01,
|
||||
"max_publish_attempts": 2,
|
||||
"publication_generations": 3,
|
||||
}
|
||||
if quality_extra:
|
||||
quality.update(quality_extra)
|
||||
settings = Settings(
|
||||
encryption_key=SecretVault.generate_key(),
|
||||
api_token="t" * 32,
|
||||
admin_password="admin-pass",
|
||||
tushare_token="test-token",
|
||||
db_path=db.path,
|
||||
quality=quality,
|
||||
scheduler_enabled=False,
|
||||
)
|
||||
pipe = Pipeline(db, adapter, settings)
|
||||
pipe._tmp = tmp
|
||||
return pipe, db
|
||||
|
||||
|
||||
def publications_map(db: HubDB, day: str) -> dict[str, str]:
|
||||
rows = db.fetchall("SELECT dataset, active_batch FROM publications WHERE trade_date = ?", (day,))
|
||||
return {str(row["dataset"]): str(row["active_batch"]) for row in rows}
|
||||
|
||||
|
||||
class ReleaseGroupSwitchTests(unittest.TestCase):
|
||||
def setUp(self) -> None:
|
||||
self.transport = GroupTransport()
|
||||
self.pipe, self.db = make_pipe(self.transport)
|
||||
self.pipe.ingest_reference(TRADE_DATE)
|
||||
|
||||
def test_whole_group_switches_in_one_publish_instant(self) -> None:
|
||||
results = self.pipe.run_eod_batch_a(TRADE_DATE)
|
||||
self.assertEqual(set(results), {*GROUP_A, "stocks"})
|
||||
self.assertEqual({item["state"] for item in results.values()}, {"published"})
|
||||
pubs = self.db.fetchall("SELECT * FROM publications WHERE trade_date = ?", (TRADE_DATE,))
|
||||
self.assertEqual(len(pubs), 5)
|
||||
self.assertEqual(len({row["published_at"] for row in pubs}), 1)
|
||||
# official rows copied and serving resolves the new batches
|
||||
api = V1API(self.db, self.pipe, self.pipe.settings)
|
||||
payload = api.handle("/v1/bars/daily", {"date": [TRADE_DATE], "code": ["600000.SH"]})
|
||||
self.assertEqual(payload["meta"]["batch_id"], results["daily"]["batch_id"])
|
||||
stocks = api.handle("/v1/stocks", {})
|
||||
self.assertEqual(stocks["meta"]["batch_id"], results["stocks"]["batch_id"])
|
||||
|
||||
def test_any_member_failure_blocks_entire_group(self) -> None:
|
||||
self.transport.empty = {"daily_basic"} # valuation upstream returns nothing
|
||||
results = self.pipe.run_eod_batch_a(TRADE_DATE)
|
||||
self.assertEqual(results["valuation"]["state"], "failed")
|
||||
self.assertEqual(results["moneyflow"]["state"], "aborted")
|
||||
self.assertEqual(results["auction"]["state"], "aborted")
|
||||
self.assertEqual(results["daily"]["state"], "failed") # staged fine, then abandoned
|
||||
# nothing became visible, and the reason is recorded
|
||||
self.assertEqual(publications_map(self.db, TRADE_DATE), {})
|
||||
abandoned = self.db.fetchall(
|
||||
"SELECT * FROM batches WHERE trade_date = ? AND state = 'failed'",
|
||||
(TRADE_DATE,),
|
||||
)
|
||||
self.assertTrue(any("release group not switched" in str(row["error"] or "") for row in abandoned))
|
||||
audit = self.db.fetchone(
|
||||
"SELECT * FROM audit_log WHERE action = 'release-group' ORDER BY id DESC"
|
||||
)
|
||||
self.assertIn("valuation", str(audit["detail"]))
|
||||
# still missing → evening retries keep trying
|
||||
self.assertIn("daily", self.pipe.missing_official_datasets(TRADE_DATE))
|
||||
|
||||
def test_failure_keeps_previous_complete_version_serving(self) -> None:
|
||||
first = self.pipe.run_dataset("daily", TRADE_DATE)
|
||||
self.transport.empty = {"daily_basic"}
|
||||
results = self.pipe.run_eod_missing(TRADE_DATE)
|
||||
# incomplete A-group restages daily with the others; valuation fails → no A switch
|
||||
self.assertEqual(results["daily"]["state"], "failed")
|
||||
self.assertEqual(results["valuation"]["state"], "failed")
|
||||
# the already-published daily batch is untouched and keeps serving
|
||||
self.assertEqual(self.pipe.active_batch("daily", TRADE_DATE), first["batch_id"])
|
||||
pubs = publications_map(self.db, TRADE_DATE)
|
||||
self.assertEqual(pubs["daily"], first["batch_id"])
|
||||
self.assertNotIn("valuation", pubs)
|
||||
self.assertNotIn("moneyflow", pubs)
|
||||
self.assertNotIn("auction", pubs)
|
||||
# B-group is an independent boundary and may still publish
|
||||
self.assertEqual(results["index_daily"]["state"], "published")
|
||||
payload = V1API(self.db, self.pipe, self.pipe.settings).handle(
|
||||
"/v1/bars/daily", {"date": [TRADE_DATE], "code": ["600000.SH"]}
|
||||
)
|
||||
self.assertEqual(payload["meta"]["batch_id"], first["batch_id"])
|
||||
|
||||
def test_partial_group_retry_does_not_mix_batches(self) -> None:
|
||||
"""Already-published A members must be restaged with missing ones."""
|
||||
first_daily = self.pipe.run_dataset("daily", TRADE_DATE)
|
||||
first_moneyflow = self.pipe.run_dataset("moneyflow", TRADE_DATE)
|
||||
results = self.pipe.run_eod_missing(TRADE_DATE)
|
||||
# A-group switched as one boundary; B-group (index) also published
|
||||
for name in (*GROUP_A, "stocks"):
|
||||
self.assertEqual(results[name]["state"], "published", name)
|
||||
self.assertEqual(results["index_daily"]["state"], "published")
|
||||
pubs = self.db.fetchall(
|
||||
"SELECT dataset, active_batch, published_at FROM publications WHERE trade_date = ?",
|
||||
(TRADE_DATE,),
|
||||
)
|
||||
by_ds = {str(row["dataset"]): row for row in pubs}
|
||||
# old partial batches replaced — no cross-batch mix of the first wave
|
||||
self.assertNotEqual(by_ds["daily"]["active_batch"], first_daily["batch_id"])
|
||||
self.assertNotEqual(by_ds["moneyflow"]["active_batch"], first_moneyflow["batch_id"])
|
||||
a_times = {by_ds[name]["published_at"] for name in (*GROUP_A, "stocks")}
|
||||
self.assertEqual(len(a_times), 1)
|
||||
# serving resolves the new complete A-group batches
|
||||
api = V1API(self.db, self.pipe, self.pipe.settings)
|
||||
daily = api.handle("/v1/bars/daily", {"date": [TRADE_DATE], "code": ["600000.SH"]})
|
||||
self.assertEqual(daily["meta"]["batch_id"], results["daily"]["batch_id"])
|
||||
self.assertEqual(daily["meta"]["batch_id"], by_ds["daily"]["active_batch"])
|
||||
|
||||
def test_reads_during_switch_see_old_state_until_commit(self) -> None:
|
||||
snapshots: list[dict] = []
|
||||
|
||||
def watcher() -> None:
|
||||
with self.db.connect() as connection:
|
||||
rows = connection.execute(
|
||||
"SELECT dataset, active_batch FROM publications WHERE trade_date = ?",
|
||||
(TRADE_DATE,),
|
||||
).fetchall()
|
||||
snapshots.append({str(row["dataset"]): row["active_batch"] for row in rows})
|
||||
|
||||
self.pipe.before_commit = watcher
|
||||
self.pipe.run_eod_batch_a(TRADE_DATE)
|
||||
# inside the switch transaction the group was still invisible
|
||||
self.assertEqual(snapshots[0], {})
|
||||
after = publications_map(self.db, TRADE_DATE)
|
||||
self.assertEqual(set(after), {*GROUP_A, "stocks"})
|
||||
|
||||
def test_switch_crash_rolls_back_whole_group(self) -> None:
|
||||
def explode() -> None:
|
||||
raise RuntimeError("killed mid-switch")
|
||||
|
||||
self.pipe.before_commit = explode
|
||||
with self.assertRaises(RuntimeError):
|
||||
self.pipe.run_eod_batch_a(TRADE_DATE)
|
||||
self.assertEqual(publications_map(self.db, TRADE_DATE), {})
|
||||
for table in ("eod_bars", "eod_valuation", "eod_moneyflow", "eod_auction", "eod_stocks"):
|
||||
rows = self.db.fetchall(f"SELECT * FROM {table} WHERE trade_date = ?", (TRADE_DATE,))
|
||||
self.assertEqual(rows, [], table)
|
||||
audit = self.db.fetchone(
|
||||
"SELECT * FROM audit_log WHERE action = 'release-group' ORDER BY id DESC"
|
||||
)
|
||||
self.assertIsNotNone(audit)
|
||||
detail = str(audit["detail"])
|
||||
self.assertIn("killed mid-switch", detail)
|
||||
self.assertIn("failed", detail)
|
||||
|
||||
def test_duplicate_runs_are_idempotent(self) -> None:
|
||||
self.pipe.run_eod_batch_a(TRADE_DATE)
|
||||
self.pipe.run_eod_batch_b(TRADE_DATE)
|
||||
batches_before = {
|
||||
str(row["batch_id"])
|
||||
for row in self.db.fetchall("SELECT batch_id FROM batches WHERE trade_date = ?", (TRADE_DATE,))
|
||||
}
|
||||
calls_before = len(self.transport.calls)
|
||||
again = self.pipe.run_eod_missing(TRADE_DATE)
|
||||
self.assertEqual({item["state"] for item in again.values()}, {"skipped"})
|
||||
self.assertEqual({item["reason"] for item in again.values()}, {"already_published"})
|
||||
batches_after = {
|
||||
str(row["batch_id"])
|
||||
for row in self.db.fetchall("SELECT batch_id FROM batches WHERE trade_date = ?", (TRADE_DATE,))
|
||||
}
|
||||
self.assertEqual(batches_after, batches_before)
|
||||
self.assertEqual(len(self.transport.calls), calls_before)
|
||||
self.assertEqual(self.pipe.missing_official_datasets(TRADE_DATE), [])
|
||||
|
||||
def test_cross_gate_failure_blocks_switch(self) -> None:
|
||||
transport = GroupTransport()
|
||||
pipe, db = make_pipe(
|
||||
transport,
|
||||
quality_extra={"cross_gates": [
|
||||
{"left": "daily", "right": "moneyflow", "min_key_overlap": 1.0},
|
||||
]},
|
||||
)
|
||||
pipe.ingest_reference(TRADE_DATE)
|
||||
transport.keep_rows["moneyflow"] = 1 # moneyflow covers only half the market
|
||||
results = pipe.run_eod_batch_a(TRADE_DATE)
|
||||
self.assertEqual(results["moneyflow"]["state"], "failed")
|
||||
self.assertIn("cross gate", str(results["moneyflow"]["error"]))
|
||||
self.assertEqual(publications_map(db, TRADE_DATE), {})
|
||||
|
||||
def test_stocks_master_and_snapshot_switch_together_or_not_at_all(self) -> None:
|
||||
original = [
|
||||
{"ts_code": "600000.SH", "symbol": "600000", "name": "浦发银行", "area": "上海",
|
||||
"industry": "银行", "market": "主板", "list_status": "L", "list_date": "19991110"},
|
||||
{"ts_code": "920071.BJ", "symbol": "920071", "name": "N金钛", "area": "辽宁",
|
||||
"industry": "小金属", "market": "北交所", "list_status": "L", "list_date": "20240901"},
|
||||
]
|
||||
renamed = [dict(original[0]), {**original[1], "name": "金钛股份"}]
|
||||
self.transport.stocks = renamed
|
||||
self.pipe.run_eod_batch_a(TRADE_DATE)
|
||||
master = self.db.fetchone("SELECT name FROM stock_master WHERE ts_code = '920071.BJ'")
|
||||
self.assertEqual(master["name"], "金钛股份")
|
||||
stocks_pub = self.db.fetchone(
|
||||
"SELECT active_batch FROM publications WHERE dataset = 'stocks' AND trade_date = ?",
|
||||
(TRADE_DATE,),
|
||||
)
|
||||
self.assertIsNotNone(stocks_pub)
|
||||
|
||||
# failure path: rename staged but the group is blocked → master stays untouched
|
||||
transport = GroupTransport()
|
||||
transport.stocks = original
|
||||
pipe, db = make_pipe(
|
||||
transport,
|
||||
quality_extra={"cross_gates": [
|
||||
{"left": "daily", "right": "moneyflow", "min_key_overlap": 1.0},
|
||||
]},
|
||||
)
|
||||
pipe.ingest_reference(TRADE_DATE) # master seeded with "N金钛"
|
||||
transport.stocks = renamed
|
||||
transport.keep_rows["moneyflow"] = 1
|
||||
results = pipe.run_eod_batch_a(TRADE_DATE)
|
||||
self.assertEqual(results["stocks"]["state"], "failed")
|
||||
master = db.fetchone("SELECT name FROM stock_master WHERE ts_code = '920071.BJ'")
|
||||
self.assertEqual(master["name"], "N金钛") # rename not applied
|
||||
stocks_pub = db.fetchone(
|
||||
"SELECT active_batch FROM publications WHERE dataset = 'stocks' AND trade_date = ?",
|
||||
(TRADE_DATE,),
|
||||
)
|
||||
self.assertIsNone(stocks_pub)
|
||||
|
||||
|
||||
class StocksRefreshAtomicTests(unittest.TestCase):
|
||||
def setUp(self) -> None:
|
||||
self.transport = GroupTransport()
|
||||
self.pipe, self.db = make_pipe(self.transport)
|
||||
self.pipe.ingest_reference(TRADE_DATE)
|
||||
self.transport.stocks = [
|
||||
{"ts_code": "600000.SH", "symbol": "600000", "name": "浦发银行", "area": "上海",
|
||||
"industry": "银行", "market": "主板", "list_status": "L", "list_date": "19991110"},
|
||||
{"ts_code": "920071.BJ", "symbol": "920071", "name": "N金钛", "area": "辽宁",
|
||||
"industry": "小金属", "market": "北交所", "list_status": "L", "list_date": "20240901"},
|
||||
]
|
||||
first = self.pipe.refresh_stocks(TRADE_DATE)
|
||||
self.assertEqual(first["state"], "published")
|
||||
self.first_batch = first["batch_id"]
|
||||
|
||||
def test_refresh_keeps_master_when_snapshot_publish_fails(self) -> None:
|
||||
self.transport.stocks = [
|
||||
{"ts_code": "600000.SH", "symbol": "600000", "name": "浦发银行", "area": "上海",
|
||||
"industry": "银行", "market": "主板", "list_status": "L", "list_date": "19991110"},
|
||||
{"ts_code": "920071.BJ", "symbol": "920071", "name": "金钛股份", "area": "辽宁",
|
||||
"industry": "小金属", "market": "北交所", "list_status": "L", "list_date": "20240901"},
|
||||
]
|
||||
|
||||
def explode() -> None:
|
||||
raise RuntimeError("snapshot switch killed")
|
||||
|
||||
self.pipe.before_commit = explode
|
||||
with self.assertRaises(RuntimeError):
|
||||
self.pipe.refresh_stocks(TRADE_DATE)
|
||||
master = self.db.fetchone("SELECT name FROM stock_master WHERE ts_code = '920071.BJ'")
|
||||
self.assertEqual(master["name"], "N金钛") # rename not applied
|
||||
self.assertEqual(self.pipe.active_batch("stocks", TRADE_DATE), self.first_batch)
|
||||
audit = self.db.fetchone(
|
||||
"SELECT * FROM audit_log WHERE action = 'stocks-refresh' ORDER BY id DESC"
|
||||
)
|
||||
self.assertIn("failed", str(audit["detail"]))
|
||||
self.assertIn("snapshot switch killed", str(audit["detail"]))
|
||||
|
||||
def test_refresh_keeps_master_when_quality_gate_rejects(self) -> None:
|
||||
self.transport.stocks = [] # empty → hard fail before publish
|
||||
with self.assertRaises(Exception):
|
||||
self.pipe.refresh_stocks(TRADE_DATE)
|
||||
master = self.db.fetchone("SELECT name FROM stock_master WHERE ts_code = '920071.BJ'")
|
||||
self.assertEqual(master["name"], "N金钛")
|
||||
self.assertEqual(self.pipe.active_batch("stocks", TRADE_DATE), self.first_batch)
|
||||
audit = self.db.fetchone(
|
||||
"SELECT * FROM audit_log WHERE action = 'stocks-refresh' ORDER BY id DESC"
|
||||
)
|
||||
self.assertIn("failed", str(audit["detail"]))
|
||||
|
||||
|
||||
class ForceBoundaryEntryTests(unittest.TestCase):
|
||||
"""CLI force / admin backfill must rebuild the full A/B boundary."""
|
||||
|
||||
def setUp(self) -> None:
|
||||
self.transport = GroupTransport()
|
||||
self.pipe, self.db = make_pipe(self.transport)
|
||||
self.pipe.ingest_reference(TRADE_DATE)
|
||||
self.first = self.pipe.run_eod_batch_a(TRADE_DATE)
|
||||
self.pipe.run_eod_batch_b(TRADE_DATE)
|
||||
|
||||
def test_force_republish_valuation_rebuilds_whole_a_group(self) -> None:
|
||||
before = publications_map(self.db, TRADE_DATE)
|
||||
results = self.pipe.force_republish_boundary("valuation", TRADE_DATE)
|
||||
self.assertEqual({item["state"] for item in results.values()}, {"published"})
|
||||
after = publications_map(self.db, TRADE_DATE)
|
||||
for name in (*GROUP_A, "stocks"):
|
||||
self.assertNotEqual(after[name], before[name], name)
|
||||
self.assertEqual(after[name], results[name]["batch_id"], name)
|
||||
# B-group left alone
|
||||
self.assertEqual(after["index_daily"], before["index_daily"])
|
||||
pubs = self.db.fetchall(
|
||||
"SELECT dataset, published_at FROM publications WHERE trade_date = ?",
|
||||
(TRADE_DATE,),
|
||||
)
|
||||
a_times = {row["published_at"] for row in pubs if row["dataset"] in {*GROUP_A, "stocks"}}
|
||||
self.assertEqual(len(a_times), 1)
|
||||
|
||||
def test_force_republish_index_rebuilds_only_b_group(self) -> None:
|
||||
before = publications_map(self.db, TRADE_DATE)
|
||||
results = self.pipe.force_republish_boundary("index_daily", TRADE_DATE)
|
||||
self.assertEqual(results["index_daily"]["state"], "published")
|
||||
after = publications_map(self.db, TRADE_DATE)
|
||||
self.assertNotEqual(after["index_daily"], before["index_daily"])
|
||||
for name in GROUP_A:
|
||||
self.assertEqual(after[name], before[name], name)
|
||||
|
||||
def test_admin_backfill_official_dataset_uses_boundary(self) -> None:
|
||||
from datahub.admin_api import AdminAPI
|
||||
from datahub.auth import AuthService
|
||||
from datahub.crypto import SecretVault
|
||||
from datahub.scheduler import Scheduler
|
||||
from datahub.serving import ApiError
|
||||
|
||||
vault = SecretVault(self.pipe.settings.encryption_key)
|
||||
auth = AuthService(self.db, vault, self.pipe.settings.api_token, "StartPass1")
|
||||
admin = AdminAPI(self.db, self.pipe, Scheduler(self.db, self.pipe), auth)
|
||||
before = publications_map(self.db, TRADE_DATE)
|
||||
result = admin.backfill("moneyflow", TRADE_DATE, "StartPass1", f"moneyflow:{TRADE_DATE}", "tester")
|
||||
self.assertEqual(result["moneyflow"]["state"], "published")
|
||||
after = publications_map(self.db, TRADE_DATE)
|
||||
for name in (*GROUP_A, "stocks"):
|
||||
self.assertNotEqual(after[name], before[name], name)
|
||||
# bad password / wrong confirm still rejected
|
||||
with self.assertRaises(ApiError):
|
||||
admin.backfill("daily", TRADE_DATE, "wrong", f"daily:{TRADE_DATE}", "tester")
|
||||
|
||||
def test_admin_backfill_switch_crash_is_failed_precondition(self) -> None:
|
||||
from datahub.admin_api import AdminAPI
|
||||
from datahub.auth import AuthService
|
||||
from datahub.crypto import SecretVault
|
||||
from datahub.scheduler import Scheduler
|
||||
from datahub.serving import ApiError
|
||||
|
||||
vault = SecretVault(self.pipe.settings.encryption_key)
|
||||
auth = AuthService(self.db, vault, self.pipe.settings.api_token, "StartPass1")
|
||||
admin = AdminAPI(self.db, self.pipe, Scheduler(self.db, self.pipe), auth)
|
||||
before = publications_map(self.db, TRADE_DATE)
|
||||
|
||||
def explode() -> None:
|
||||
raise RuntimeError("killed mid-switch")
|
||||
|
||||
self.pipe.before_commit = explode
|
||||
with self.assertRaises(ApiError) as ctx:
|
||||
admin.backfill("valuation", TRADE_DATE, "StartPass1", f"valuation:{TRADE_DATE}", "tester")
|
||||
self.assertEqual(ctx.exception.code, "FAILED_PRECONDITION")
|
||||
self.assertIn("killed mid-switch", ctx.exception.message)
|
||||
# previous complete A/B versions keep serving
|
||||
self.assertEqual(publications_map(self.db, TRADE_DATE), before)
|
||||
audit = self.db.fetchone(
|
||||
"SELECT * FROM audit_log WHERE action = 'release-group' ORDER BY id DESC"
|
||||
)
|
||||
self.assertIsNotNone(audit)
|
||||
self.assertIn("failed", str(audit["detail"]))
|
||||
self.assertIn("killed mid-switch", str(audit["detail"]))
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
unittest.main()
|
||||
@@ -0,0 +1,65 @@
|
||||
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,11 +19,17 @@ class LayoutTests(unittest.TestCase):
|
||||
def test_reserved_adapters_present(self) -> None:
|
||||
from datahub.adapters import RESERVED
|
||||
|
||||
for name in ("eastmoney", "tencent", "ths", "xgb", "akshare", "ifind"):
|
||||
for name in ("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__":
|
||||
|
||||
@@ -228,25 +228,34 @@ class GateRetryInterplayTests(unittest.TestCase):
|
||||
|
||||
|
||||
class ForceRepublishTests(unittest.TestCase):
|
||||
def test_run_dataset_over_published_keeps_prev_for_rollback(self) -> None:
|
||||
def test_force_boundary_republish_keeps_prev_for_rollback(self) -> None:
|
||||
transport = ValuationTransport()
|
||||
pipe, db = make_pipe(transport)
|
||||
pipe.ingest_reference(TRADE_DATE)
|
||||
first = pipe.run_dataset("valuation", TRADE_DATE)
|
||||
first = pipe.run_eod_batch_a(TRADE_DATE)
|
||||
first_val = first["valuation"]["batch_id"]
|
||||
first_daily = first["daily"]["batch_id"]
|
||||
transport.mode = "vr_all_null"
|
||||
with self.assertRaises(QualityError):
|
||||
pipe.run_dataset("valuation", TRADE_DATE) # gate holds: bad re-publish refused
|
||||
blocked = pipe.force_republish_boundary("valuation", TRADE_DATE)
|
||||
self.assertEqual(blocked["valuation"]["state"], "failed")
|
||||
self.assertEqual(pipe.active_batch("valuation", TRADE_DATE), first_val)
|
||||
self.assertEqual(pipe.active_batch("daily", TRADE_DATE), first_daily)
|
||||
transport.mode = "ok"
|
||||
second = pipe.run_dataset("valuation", TRADE_DATE) # CLI --force path
|
||||
self.assertNotEqual(first["batch_id"], second["batch_id"])
|
||||
pub = db.fetchone(
|
||||
"SELECT * FROM publications WHERE dataset='valuation' AND trade_date=?",
|
||||
second = pipe.force_republish_boundary("valuation", TRADE_DATE)
|
||||
self.assertEqual(second["valuation"]["state"], "published")
|
||||
self.assertNotEqual(second["valuation"]["batch_id"], first_val)
|
||||
self.assertNotEqual(second["daily"]["batch_id"], first_daily)
|
||||
pubs = db.fetchall(
|
||||
"SELECT dataset, active_batch, prev_batch, published_at FROM publications WHERE trade_date=?",
|
||||
(TRADE_DATE,),
|
||||
)
|
||||
self.assertEqual(pub["active_batch"], second["batch_id"])
|
||||
self.assertEqual(pub["prev_batch"], first["batch_id"])
|
||||
by_ds = {str(row["dataset"]): row for row in pubs}
|
||||
a_times = {by_ds[name]["published_at"] for name in ("daily", "valuation", "moneyflow", "auction", "stocks")}
|
||||
self.assertEqual(len(a_times), 1)
|
||||
self.assertEqual(by_ds["valuation"]["active_batch"], second["valuation"]["batch_id"])
|
||||
self.assertEqual(by_ds["valuation"]["prev_batch"], first_val)
|
||||
rolled = pipe.rollback("valuation", TRADE_DATE, actor="cli")
|
||||
self.assertEqual(rolled["active_batch"], first["batch_id"])
|
||||
self.assertEqual(rolled["active_batch"], first_val)
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
|
||||
Reference in New Issue
Block a user