diff --git a/app/service/market_quote_sync_service.py b/app/service/market_quote_sync_service.py index 5926504..f82a299 100644 --- a/app/service/market_quote_sync_service.py +++ b/app/service/market_quote_sync_service.py @@ -61,9 +61,16 @@ class MarketQuoteSyncService: try: quotes = await asyncio.to_thread(loader, codes) error_type = None + rejected: list[str] = [] except Exception as exc: - quotes = {} + # 数据源对**整批**调用抛错的常见原因是:其中某个代码它不接受 + # (本库混入了一只非本公司产品,而 loader 侧有白名单校验)。 + # 整批失败会连带**其余正常产品一起拿不到行情** —— 实测后果是 + # `fin_market_price` 长期不更新、所有客户下单 503,而错误信息只说 + # "某产品行情已过期",完全看不出根因。所以这里降级为逐只调用: + # 保留能取到的,把被拒的代码记下来(进审计与返回值)。 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} for code, quote in usable.items(): selected.setdefault(code, (source, quote)) @@ -76,7 +83,12 @@ class MarketQuoteSyncService: source_results.append({ "source": source, "priority": priority, "status": status, "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, }) if len(selected) == len(codes): @@ -101,6 +113,27 @@ class MarketQuoteSyncService: "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 async def _record_source_run( session: Any, run_no: str, item: dict[str, object], now: datetime