fix(sync): 行情同步改为逐只容错,避免一只坏代码拖垮全量
全流程验收时发现的核心故障:**所有客户都无法下单**(任意产品都返回
503 FUND_QUOTE_UNAVAILABLE),而错误信息只说"某产品行情已过期",看不出根因。
根因链:
1) MarketQuoteSyncService.sync() 把库里 fund_manager='南方基金' 的 20 只产品
打包交给 hq.py 的 loader;
2) 而 hq.py 的白名单校验是 ny(code not in SOUTHERN_FUND_CODES) -> raise,
**一个代码不在名单里整批就抛错**;
3) 库里混着 510300 —— 它的真实管理人是华泰柏瑞(前面查费率时确认过),
不在南方基金白名单里 -> 两个数据源全部 ValueError -> quotes=0;
4) 于是 in_market_price 长期不更新(当时只有 2 行、时间停在 08:19),
下单的 MAX_QUOTE_AGE 校验全线失败。
修法:整批调用失败时**降级为逐只调用**,保留能取到的行情,把被拒代码记进
source_results(并拼进 error_type,让它在 dvisor_market_quote_source_run
审计行里可见,不必为此改表)。被拒代码通常意味着 und_manager 标注有问题,
值得人工核对。
不动 hq.py 的白名单:那是刻意的防误用设计(不许拿任意代码查南方行情),
该容错的是编排层。
验证:重跑 tools/sync_advisor_market_quotes.py 后不再是全批失败;
配合重跑 tools/seed_sim_account_demo 刷新 fin_market_price,
下单恢复 —— 510500、510300 均 HTTP=201 已成交。
⚠️ 仍未解决(数据供给缺口,需要产品/业务侧决定):
in_market_price 目前**只有演示种子脚本写**,没有生产同步链路。
后果是 20 只产品里只有 2 只(510300/510500)有行情、其余 18 只下单报
"缺少场内行情",而且行情一旦过期就没有任何机制刷新。
This commit is contained in:
@@ -61,9 +61,16 @@ class MarketQuoteSyncService:
|
|||||||
try:
|
try:
|
||||||
quotes = await asyncio.to_thread(loader, codes)
|
quotes = await asyncio.to_thread(loader, codes)
|
||||||
error_type = None
|
error_type = None
|
||||||
|
rejected: list[str] = []
|
||||||
except Exception as exc:
|
except Exception as exc:
|
||||||
quotes = {}
|
# 数据源对**整批**调用抛错的常见原因是:其中某个代码它不接受
|
||||||
|
# (本库混入了一只非本公司产品,而 loader 侧有白名单校验)。
|
||||||
|
# 整批失败会连带**其余正常产品一起拿不到行情** —— 实测后果是
|
||||||
|
# `fin_market_price` 长期不更新、所有客户下单 503,而错误信息只说
|
||||||
|
# "某产品行情已过期",完全看不出根因。所以这里降级为逐只调用:
|
||||||
|
# 保留能取到的,把被拒的代码记下来(进审计与返回值)。
|
||||||
error_type = type(exc).__name__
|
error_type = type(exc).__name__
|
||||||
|
quotes, rejected = await self._load_one_by_one(loader, codes)
|
||||||
usable = {code: value for code, value in quotes.items() if code in product_ids}
|
usable = {code: value for code, value in quotes.items() if code in product_ids}
|
||||||
for code, quote in usable.items():
|
for code, quote in usable.items():
|
||||||
selected.setdefault(code, (source, quote))
|
selected.setdefault(code, (source, quote))
|
||||||
@@ -76,7 +83,12 @@ class MarketQuoteSyncService:
|
|||||||
source_results.append({
|
source_results.append({
|
||||||
"source": source, "priority": priority, "status": status,
|
"source": source, "priority": priority, "status": status,
|
||||||
"requested_count": len(codes), "quote_count": len(usable),
|
"requested_count": len(codes), "quote_count": len(usable),
|
||||||
"error_type": error_type,
|
# 逐只降级后被数据源拒绝的代码:拼进 error_type 是为了让它出现在
|
||||||
|
# `advisor_market_quote_source_run` 审计行里(该表没有独立列,也不想为它改表)。
|
||||||
|
"error_type": (
|
||||||
|
f"{error_type}(rejected={len(rejected)})" if rejected else error_type
|
||||||
|
),
|
||||||
|
"rejected_codes": rejected,
|
||||||
"started_at": started, "completed_at": completed,
|
"started_at": started, "completed_at": completed,
|
||||||
})
|
})
|
||||||
if len(selected) == len(codes):
|
if len(selected) == len(codes):
|
||||||
@@ -101,6 +113,27 @@ class MarketQuoteSyncService:
|
|||||||
"degraded": len(selected) < len(codes),
|
"degraded": len(selected) < len(codes),
|
||||||
}
|
}
|
||||||
|
|
||||||
|
@staticmethod
|
||||||
|
async def _load_one_by_one(
|
||||||
|
loader: QuoteLoader, codes: list[str]
|
||||||
|
) -> tuple[dict[str, dict[str, str | None]], list[str]]:
|
||||||
|
"""逐只调用 loader,返回(成功合并的行情, 被数据源拒绝的代码)。
|
||||||
|
|
||||||
|
只在整批调用失败后走这条路:代价是 N 次请求,换来的是"一只坏代码不拖垮全量"。
|
||||||
|
被拒的代码通常是**不该出现在本同步范围内的产品**(例如混进来的非本公司产品),
|
||||||
|
记下来供人工核对 `fin_product.fund_manager` 的标注是否正确。
|
||||||
|
"""
|
||||||
|
merged: dict[str, dict[str, str | None]] = {}
|
||||||
|
rejected: list[str] = []
|
||||||
|
for code in codes:
|
||||||
|
try:
|
||||||
|
one = await asyncio.to_thread(loader, [code])
|
||||||
|
except Exception: # noqa: BLE001 - 单只失败不影响其余
|
||||||
|
rejected.append(code)
|
||||||
|
continue
|
||||||
|
merged.update(one)
|
||||||
|
return merged, rejected
|
||||||
|
|
||||||
@staticmethod
|
@staticmethod
|
||||||
async def _record_source_run(
|
async def _record_source_run(
|
||||||
session: Any, run_no: str, item: dict[str, object], now: datetime
|
session: Any, run_no: str, item: dict[str, object], now: datetime
|
||||||
|
|||||||
Reference in New Issue
Block a user