mirror of
https://github.com/ZhuLinsen/daily_stock_analysis.git
synced 2026-10-06 14:33:11 +08:00
feat: 增加 A 股指数多数据源 fallback 路由 (#2258)
* feat: add A-share index fallback routing * docs: clarify index priority configuration scope
This commit is contained in:
+7
-3
@@ -892,6 +892,10 @@ ADMIN_AUTH_ENABLED=false
|
||||
# ===========================================
|
||||
# Lower number = higher priority (tried first)
|
||||
# Default priorities: efinance(0) > akshare(1) > tushare/pytdx(2) > baostock(3) > yfinance(4) > tencent(5)
|
||||
# 上述优先级仅控制普通 A 股日 K 通用链路。
|
||||
# IndexRegistry 已登记的沪深指数(当前为 sh000016、sh000688、sz399001、sz399006、sh000300)
|
||||
# 固定按 Tencent > AkShare > TickFlow > YFinance 尝试,不读取这些 *_PRIORITY 配置;未配置或不可用的来源会跳过。
|
||||
# 只有显式市场输入(如 sh000016 或 000016.SH)会触发指数链;裸 000016 仍按股票处理。
|
||||
# For US stocks, set YFINANCE_PRIORITY=0 to use Yahoo Finance first
|
||||
|
||||
# EFINANCE_PRIORITY=0 # EastMoney (China) - default: 0
|
||||
@@ -899,7 +903,7 @@ ADMIN_AUTH_ENABLED=false
|
||||
# # when eastmoney hosts are unreachable. Default: 30
|
||||
# AKSHARE_PRIORITY=1 # AkShare (China) - default: 1
|
||||
# TUSHARE_PRIORITY=2 # Tushare Pro (China) - default: 2
|
||||
# TICKFLOW_PRIORITY=2 # TickFlow(A 股)- 默认:2;可选,需配置 TICKFLOW_API_KEY
|
||||
# TICKFLOW_PRIORITY=2 # TickFlow(普通 A 股日 K)- 默认:2;可选,需配置 TICKFLOW_API_KEY
|
||||
# PYTDX_PRIORITY=2 # Tongdaxin (China) - default: 2
|
||||
#
|
||||
# Pytdx custom server (for intranet/deploy): use custom host instead of built-in public servers
|
||||
@@ -908,8 +912,8 @@ ADMIN_AUTH_ENABLED=false
|
||||
# Or multiple servers: PYTDX_SERVERS=ip1:port1,ip2:port2
|
||||
# BAOSTOCK_PRIORITY=3 # Baostock (China) - default: 3
|
||||
# YFINANCE_PRIORITY=4 # Yahoo Finance (Global) - default: 4
|
||||
# TENCENT_PRIORITY=5 # Tencent direct daily K-line (China) - last-resort
|
||||
# # A-share fallback. Default: 5
|
||||
# TENCENT_PRIORITY=5 # Tencent direct daily K-line - ordinary A-share last-resort
|
||||
# # fallback. Registered indices ignore this value. Default: 5
|
||||
|
||||
# Example: Prioritize Yahoo Finance for US stocks
|
||||
# YFINANCE_PRIORITY=0
|
||||
|
||||
@@ -286,10 +286,10 @@ const settingsHelpZhCN: SettingsHelpMap = {
|
||||
},
|
||||
'settings.data_source.TICKFLOW_PRIORITY': {
|
||||
title: 'TickFlow 日 K 优先级',
|
||||
summary: '控制 TickFlow 在 A 股日 K 数据源回退链中的位置。',
|
||||
usage: '填写整数;数字越小越早尝试,默认 2。未配置 TICKFLOW_API_KEY 时该优先级不会生效。',
|
||||
valueNotes: ['该设置只影响日 K 等通用数据源回退链,不控制实时行情源顺序。'],
|
||||
impact: ['影响 A 股日 K 获取的数据源尝试顺序;实时行情仍由 REALTIME_SOURCE_PRIORITY 单独决定。'],
|
||||
summary: '控制 TickFlow 在普通 A 股日 K 数据源回退链中的位置。',
|
||||
usage: '填写整数;数字越小越早尝试,默认 2。未配置 TICKFLOW_API_KEY 时该优先级不会生效;已登记指数固定按 Tencent → AkShare → TickFlow → YFinance 降级,不读取本配置。',
|
||||
valueNotes: ['该设置只影响普通 A 股日 K 等通用数据源回退链,不控制已登记指数或实时行情源顺序。'],
|
||||
impact: ['影响普通 A 股日 K 获取的数据源尝试顺序;已登记指数和实时行情均使用各自独立顺序。'],
|
||||
notes: ['如果希望优先使用 TickFlow 日 K,可以适当调低该值;如果希望实时行情优先使用 TickFlow,请在 REALTIME_SOURCE_PRIORITY 中显式加入 tickflow。'],
|
||||
},
|
||||
'settings.data_source.TICKFLOW_KLINE_ADJUST': {
|
||||
@@ -1492,10 +1492,10 @@ const settingsHelpEnUS: SettingsHelpMap = {
|
||||
},
|
||||
'settings.data_source.TICKFLOW_PRIORITY': {
|
||||
title: 'TickFlow Daily K-line Priority',
|
||||
summary: 'Controls where TickFlow sits in the A-share daily K-line provider fallback chain.',
|
||||
usage: 'Use an integer. Lower numbers are tried earlier. The default is 2. This has no effect unless TICKFLOW_API_KEY is configured.',
|
||||
valueNotes: ['This setting only affects the daily K-line/general data-source fallback chain; it does not control realtime quote provider order.'],
|
||||
impact: ['Affects provider order for A-share daily K-line fetching. Realtime quotes are still controlled separately by REALTIME_SOURCE_PRIORITY.'],
|
||||
summary: 'Controls where TickFlow sits in the generic A-share daily K-line provider fallback chain.',
|
||||
usage: 'Use an integer. Lower numbers are tried earlier. The default is 2. This has no effect unless TICKFLOW_API_KEY is configured. Registered indices use the fixed Tencent → AkShare → TickFlow → YFinance chain and ignore this setting.',
|
||||
valueNotes: ['This setting only affects generic A-share daily K-lines; it does not control registered-index or realtime-quote provider order.'],
|
||||
impact: ['Affects provider order for generic A-share daily K-line fetching. Registered indices and realtime quotes use separate ordering.'],
|
||||
notes: ['Lower this value only if you want TickFlow daily K-lines to be tried earlier. Add tickflow to REALTIME_SOURCE_PRIORITY when you want TickFlow realtime quotes in the realtime fallback chain.'],
|
||||
},
|
||||
'settings.data_source.TICKFLOW_KLINE_ADJUST': {
|
||||
|
||||
@@ -232,7 +232,7 @@ const fieldDescriptionMap: Record<string, string> = {
|
||||
NEWS_MAX_AGE_DAYS: '新闻最大时效上限。实际窗口 = min(策略档位天数, NEWS_MAX_AGE_DAYS)。例如 ultra_short + 7 仍为 1 天。',
|
||||
REALTIME_SOURCE_PRIORITY: '按逗号分隔填写数据源调用优先级。',
|
||||
TICKFLOW_API_KEY: '用于接入 TickFlow 数据服务的 API 密钥。',
|
||||
TICKFLOW_PRIORITY: '控制 TickFlow 在 A 股日 K 数据源回退链中的尝试顺序;不控制实时行情,实时行情顺序由 REALTIME_SOURCE_PRIORITY 决定。',
|
||||
TICKFLOW_PRIORITY: '控制 TickFlow 在普通 A 股日 K 回退链中的尝试顺序;已登记指数使用 Tencent → AkShare → TickFlow → YFinance 固定链,不读取本配置;实时行情顺序由 REALTIME_SOURCE_PRIORITY 决定。',
|
||||
TICKFLOW_KLINE_ADJUST: '控制 TickFlow 日 K 线的复权口径,默认 none 保持未复权技术指标基线。',
|
||||
TICKFLOW_BATCH_DAILY_ENABLED: '批量分析前是否使用 TickFlow 批量日 K 接口预热缓存;权限不足时会继续回退。',
|
||||
TICKFLOW_BATCH_SIZE: '控制 TickFlow 日 K 和实时行情批量请求的单批最大标的数。',
|
||||
|
||||
@@ -45,6 +45,7 @@ from tenacity import (
|
||||
|
||||
from src.patches.eastmoney_patch import eastmoney_patch
|
||||
from src.config import get_config
|
||||
from src.services.stock_list_parser import ParseStatus, parse_analysis_target
|
||||
from .base import BaseFetcher, DataFetchError, RateLimitError, STANDARD_COLUMNS, is_bse_code, is_st_stock, is_kc_cy_stock, normalize_stock_code
|
||||
from .realtime_types import (
|
||||
UnifiedRealtimeQuote, ChipDistribution, RealtimeSource,
|
||||
@@ -467,6 +468,7 @@ class AkshareFetcher(BaseFetcher):
|
||||
从 Akshare 获取原始数据
|
||||
|
||||
根据代码类型自动选择 API:
|
||||
- 已登记 A 股指数:使用 ak.stock_zh_index_daily_em()
|
||||
- 美股:不支持,抛出异常由 YfinanceFetcher 处理(Issue #311)
|
||||
- 港股:使用 ak.stock_hk_hist()
|
||||
- ETF 基金:使用 ak.fund_etf_hist_em()
|
||||
@@ -479,6 +481,10 @@ class AkshareFetcher(BaseFetcher):
|
||||
4. 调用对应的 akshare API
|
||||
5. 处理返回数据
|
||||
"""
|
||||
target = parse_analysis_target(stock_code)
|
||||
if target.asset_type == ParseStatus.INDEX:
|
||||
return self._fetch_index_data(target.canonical_id, start_date, end_date)
|
||||
|
||||
# 根据代码类型选择不同的获取方法
|
||||
if _is_us_code(stock_code):
|
||||
# 美股:akshare 的 stock_us_daily 接口复权存在已知问题(参见 Issue #311)
|
||||
@@ -492,6 +498,85 @@ class AkshareFetcher(BaseFetcher):
|
||||
return self._fetch_etf_data(stock_code, start_date, end_date)
|
||||
else:
|
||||
return self._fetch_stock_data(stock_code, start_date, end_date)
|
||||
|
||||
def _fetch_index_data(
|
||||
self, stock_code: str, start_date: str, end_date: str
|
||||
) -> pd.DataFrame:
|
||||
"""Fetch a registry-recognized A-share index from Eastmoney."""
|
||||
import akshare as ak
|
||||
|
||||
self._set_random_user_agent()
|
||||
self._enforce_rate_limit()
|
||||
logger.info(
|
||||
"[API调用] ak.stock_zh_index_daily_em(symbol=%s, start_date=%s, end_date=%s)",
|
||||
stock_code,
|
||||
start_date.replace("-", ""),
|
||||
end_date.replace("-", ""),
|
||||
)
|
||||
try:
|
||||
df = _akshare_call_with_timeout(
|
||||
ak.stock_zh_index_daily_em,
|
||||
timeout=self._history_call_timeout,
|
||||
call_name="ak.stock_zh_index_daily_em",
|
||||
symbol=stock_code,
|
||||
start_date=start_date.replace("-", ""),
|
||||
end_date=end_date.replace("-", ""),
|
||||
)
|
||||
except (ConnectionError, TimeoutError):
|
||||
raise
|
||||
except RateLimitError:
|
||||
raise
|
||||
except Exception as exc:
|
||||
error_msg = str(exc).lower()
|
||||
if any(
|
||||
keyword in error_msg
|
||||
for keyword in ("banned", "blocked", "频率", "rate", "限制")
|
||||
):
|
||||
raise RateLimitError(f"Akshare 指数接口可能被限流: {exc}") from exc
|
||||
raise DataFetchError(f"Akshare 获取指数数据失败: {exc}") from exc
|
||||
|
||||
if df is None:
|
||||
return pd.DataFrame()
|
||||
if not isinstance(df, pd.DataFrame):
|
||||
raise DataFetchError(
|
||||
f"Akshare 指数接口返回无效类型: {type(df).__name__}"
|
||||
)
|
||||
if df.empty:
|
||||
return df.copy()
|
||||
|
||||
required_columns = {"date", "open", "high", "low", "close", "volume"}
|
||||
missing_columns = sorted(required_columns - set(df.columns))
|
||||
if missing_columns:
|
||||
raise DataFetchError(
|
||||
"Akshare 指数数据缺少必需列: " + ", ".join(missing_columns)
|
||||
)
|
||||
|
||||
result = df.copy()
|
||||
parsed_dates = pd.to_datetime(result["date"], errors="coerce", format="mixed")
|
||||
if parsed_dates.isna().any():
|
||||
invalid_dates = result.loc[parsed_dates.isna(), "date"].astype(str).tolist()
|
||||
raise DataFetchError(
|
||||
"Akshare 指数数据包含无法解析的 date: "
|
||||
+ ", ".join(invalid_dates[:3])
|
||||
)
|
||||
result["_index_sort_date"] = parsed_dates
|
||||
result = (
|
||||
result.sort_values("_index_sort_date", kind="stable")
|
||||
.drop(columns="_index_sort_date")
|
||||
.reset_index(drop=True)
|
||||
)
|
||||
if "pct_chg" not in result.columns:
|
||||
close = pd.to_numeric(result["close"], errors="coerce")
|
||||
pct_chg = close.pct_change(fill_method=None) * 100
|
||||
pct_chg = pct_chg.replace(
|
||||
[float("inf"), float("-inf")], float("nan")
|
||||
)
|
||||
if not pct_chg.empty and pd.isna(pct_chg.iloc[0]):
|
||||
pct_chg.iloc[0] = 0.0
|
||||
result["pct_chg"] = pct_chg
|
||||
if "amount" not in result.columns:
|
||||
result["amount"] = pd.NA
|
||||
return result
|
||||
|
||||
def _fetch_stock_data(self, stock_code: str, start_date: str, end_date: str) -> pd.DataFrame:
|
||||
"""
|
||||
@@ -910,7 +995,7 @@ class AkshareFetcher(BaseFetcher):
|
||||
|
||||
# 重命名列
|
||||
df = df.rename(columns=column_mapping)
|
||||
|
||||
|
||||
# 添加股票代码列
|
||||
df['code'] = stock_code
|
||||
|
||||
@@ -920,6 +1005,56 @@ class AkshareFetcher(BaseFetcher):
|
||||
df = df[existing_cols]
|
||||
|
||||
return df
|
||||
|
||||
def get_stock_name(self, stock_code: str) -> Optional[str]:
|
||||
target = parse_analysis_target(stock_code)
|
||||
if target.asset_type != ParseStatus.INDEX or target.matched_index is None:
|
||||
return None
|
||||
|
||||
import akshare as ak
|
||||
|
||||
table_name = (
|
||||
"上证系列指数"
|
||||
if target.matched_index.exchange.upper() == "SH"
|
||||
else "深证系列指数"
|
||||
)
|
||||
try:
|
||||
self._set_random_user_agent()
|
||||
self._enforce_rate_limit()
|
||||
df = _akshare_call_with_timeout(
|
||||
ak.stock_zh_index_spot_em,
|
||||
timeout=self._history_call_timeout,
|
||||
call_name="ak.stock_zh_index_spot_em",
|
||||
symbol=table_name,
|
||||
)
|
||||
if df is None or df.empty or not {"代码", "名称"}.issubset(df.columns):
|
||||
return None
|
||||
|
||||
bare_code = target.matched_index.bare_code
|
||||
|
||||
def normalize_index_code(value: Any) -> str:
|
||||
text = str(value).strip()
|
||||
if text.endswith(".0"):
|
||||
text = text[:-2]
|
||||
return text.zfill(6) if text.isdigit() else text
|
||||
|
||||
matches = df[df["代码"].map(normalize_index_code) == bare_code]
|
||||
if matches.empty:
|
||||
return None
|
||||
for name_value in matches["名称"]:
|
||||
if pd.isna(name_value):
|
||||
continue
|
||||
name = str(name_value).strip()
|
||||
if name:
|
||||
return name
|
||||
return None
|
||||
except Exception:
|
||||
logger.debug(
|
||||
"Akshare index name lookup failed for %s",
|
||||
target.canonical_id,
|
||||
exc_info=True,
|
||||
)
|
||||
return None
|
||||
|
||||
def get_realtime_quote(self, stock_code: str, source: str = "em") -> Optional[UnifiedRealtimeQuote]:
|
||||
"""
|
||||
|
||||
+390
-6
@@ -28,6 +28,7 @@ from src.data.stock_index_loader import get_index_stock_name
|
||||
from src.data.stock_mapping import STOCK_NAME_MAP, is_meaningful_stock_name
|
||||
from src.services.market_symbol_utils import is_suffix_market_symbol
|
||||
from src.services.run_diagnostics import record_provider_run, record_provider_run_started
|
||||
from src.services.stock_list_parser import AnalysisTarget, ParseStatus, parse_analysis_target
|
||||
from .fundamental_adapter import AkshareFundamentalAdapter
|
||||
from .yfinance_fundamental_adapter import YfinanceFundamentalAdapter
|
||||
from .realtime_types import CircuitBreaker
|
||||
@@ -629,6 +630,18 @@ class DataFetcherManager:
|
||||
"AlphaVantageFetcher": {"us"},
|
||||
}
|
||||
_daily_source_health = CircuitBreaker(failure_threshold=3, cooldown_seconds=300.0)
|
||||
_CN_INDEX_DAILY_SOURCE_ORDER = (
|
||||
"TencentFetcher",
|
||||
"AkshareFetcher",
|
||||
"TickFlowFetcher",
|
||||
"YfinanceFetcher",
|
||||
)
|
||||
_CN_INDEX_NAME_SOURCE_ORDER = (
|
||||
"TencentFetcher",
|
||||
"AkshareFetcher",
|
||||
"TickFlowFetcher",
|
||||
)
|
||||
_CN_INDEX_BARE_CODE_CONFLICTS = frozenset({"000001", "000016", "000688"})
|
||||
_CONCEPT_RANKINGS_CACHE_TTL_SECONDS = 300.0
|
||||
_CONCEPT_RANKINGS_EMPTY_CACHE_TTL_SECONDS = 30.0
|
||||
_concept_rankings_cache_lock = RLock()
|
||||
@@ -847,6 +860,353 @@ class DataFetcherManager:
|
||||
self._stock_name_cache[stock_code] = name
|
||||
return name
|
||||
|
||||
def _discard_cached_stock_name(self, stock_code: str) -> None:
|
||||
self._ensure_concurrency_guards()
|
||||
with self._stock_name_cache_lock:
|
||||
self._stock_name_cache.pop(stock_code, None)
|
||||
|
||||
@classmethod
|
||||
def _warn_bare_index_conflict(cls, target: AnalysisTarget) -> None:
|
||||
if target.asset_type != ParseStatus.STOCK or target.normalized_prefix is not None:
|
||||
return
|
||||
bare_code = target.normalized_code or (target.raw_input or "").strip()
|
||||
if (
|
||||
target.matched_index is None
|
||||
and bare_code not in cls._CN_INDEX_BARE_CODE_CONFLICTS
|
||||
):
|
||||
return
|
||||
logger.warning(
|
||||
"[指数路由] 裸代码 %s 存在股票/指数歧义,按股票路由处理",
|
||||
bare_code,
|
||||
)
|
||||
|
||||
@staticmethod
|
||||
def _cn_index_provider_symbol(target: AnalysisTarget, fetcher_name: str) -> str:
|
||||
entry = target.matched_index
|
||||
if entry is None:
|
||||
return ""
|
||||
exchange = entry.exchange.upper()
|
||||
if exchange not in {"SH", "SZ"}:
|
||||
return ""
|
||||
if fetcher_name in {"TencentFetcher", "AkshareFetcher"}:
|
||||
return f"{exchange.lower()}{entry.bare_code}"
|
||||
if fetcher_name == "TickFlowFetcher":
|
||||
return f"{entry.bare_code}.{exchange}"
|
||||
if fetcher_name == "YfinanceFetcher":
|
||||
suffix = "SS" if exchange == "SH" else "SZ"
|
||||
return f"{entry.bare_code}.{suffix}"
|
||||
return ""
|
||||
|
||||
@classmethod
|
||||
def _is_meaningful_cn_index_name(
|
||||
cls, name: Optional[str], target: AnalysisTarget
|
||||
) -> bool:
|
||||
if not is_meaningful_stock_name(name, target.canonical_id):
|
||||
return False
|
||||
|
||||
aliases = {target.canonical_id}
|
||||
entry = target.matched_index
|
||||
if entry is not None:
|
||||
if normalize_stock_code(str(name).strip()) == entry.bare_code:
|
||||
return False
|
||||
aliases.add(entry.bare_code)
|
||||
aliases.update(entry.aliases)
|
||||
display_code = (target.display_code or "").strip()
|
||||
registry_name = (entry.display_name or "").strip()
|
||||
# The parser currently uses display_code for the human label too.
|
||||
if display_code and (
|
||||
display_code != registry_name
|
||||
or normalize_stock_code(display_code) == entry.bare_code
|
||||
):
|
||||
aliases.add(display_code)
|
||||
for source_name in cls._CN_INDEX_DAILY_SOURCE_ORDER:
|
||||
aliases.add(cls._cn_index_provider_symbol(target, source_name))
|
||||
|
||||
candidate = str(name).strip().upper()
|
||||
code_aliases = {
|
||||
str(alias).strip().upper()
|
||||
for alias in aliases
|
||||
if alias is not None and str(alias).strip()
|
||||
}
|
||||
return candidate not in code_aliases
|
||||
|
||||
def _get_cn_index_daily_data(
|
||||
self,
|
||||
target: AnalysisTarget,
|
||||
start_date: Optional[str],
|
||||
end_date: Optional[str],
|
||||
days: int,
|
||||
) -> Tuple[pd.DataFrame, str]:
|
||||
source_order = self._CN_INDEX_DAILY_SOURCE_ORDER
|
||||
fetchers_by_name = {
|
||||
fetcher.name: fetcher for fetcher in self._get_fetchers_snapshot()
|
||||
}
|
||||
errors: List[str] = []
|
||||
request_start = time.time()
|
||||
|
||||
for index, source_name in enumerate(source_order):
|
||||
fallback_to = source_order[index + 1] if index + 1 < len(source_order) else None
|
||||
fetcher = fetchers_by_name.get(source_name)
|
||||
provider_symbol = self._cn_index_provider_symbol(target, source_name)
|
||||
|
||||
if not provider_symbol:
|
||||
reason = (
|
||||
"unsupported index provider symbol: "
|
||||
f"{target.canonical_id} -> {source_name}"
|
||||
)
|
||||
record_provider_run(
|
||||
data_type="daily_data",
|
||||
provider=source_name,
|
||||
operation="get_daily_data",
|
||||
success=False,
|
||||
latency_ms=0,
|
||||
error_type="unsupported",
|
||||
error_message=reason,
|
||||
fallback_to=fallback_to,
|
||||
record_count=0,
|
||||
)
|
||||
logger.warning(
|
||||
"[指数数据源不支持 %d/%d] [%s] %s: %s",
|
||||
index + 1,
|
||||
len(source_order),
|
||||
source_name,
|
||||
target.canonical_id,
|
||||
reason,
|
||||
)
|
||||
errors.append(f"[{source_name}] {reason}")
|
||||
continue
|
||||
|
||||
if fetcher is None or not self._is_fetcher_available(
|
||||
fetcher, capability="daily_data"
|
||||
):
|
||||
reason = "数据源未配置或暂不可用"
|
||||
record_provider_run(
|
||||
data_type="daily_data",
|
||||
provider=source_name,
|
||||
operation="get_daily_data",
|
||||
success=False,
|
||||
latency_ms=0,
|
||||
error_type="unavailable",
|
||||
error_message=reason,
|
||||
fallback_to=fallback_to,
|
||||
record_count=0,
|
||||
)
|
||||
logger.warning(
|
||||
"[指数数据源失败 %d/%d] [%s] %s: %s",
|
||||
index + 1,
|
||||
len(source_order),
|
||||
source_name,
|
||||
target.canonical_id,
|
||||
reason,
|
||||
)
|
||||
errors.append(f"[{source_name}] {reason}")
|
||||
continue
|
||||
|
||||
if not self._is_daily_source_available(fetcher, "cn_index"):
|
||||
reason = self._daily_source_unavailable_error(fetcher)
|
||||
record_provider_run(
|
||||
data_type="daily_data",
|
||||
provider=source_name,
|
||||
operation="get_daily_data",
|
||||
success=False,
|
||||
latency_ms=0,
|
||||
error_type="CircuitOpen",
|
||||
error_message=reason,
|
||||
fallback_to=fallback_to,
|
||||
record_count=0,
|
||||
)
|
||||
logger.warning(
|
||||
"[指数数据源失败 %d/%d] [%s] %s: %s",
|
||||
index + 1,
|
||||
len(source_order),
|
||||
source_name,
|
||||
target.canonical_id,
|
||||
reason,
|
||||
)
|
||||
errors.append(reason)
|
||||
continue
|
||||
|
||||
attempt_start = time.time()
|
||||
try:
|
||||
logger.info(
|
||||
"[指数数据源尝试 %d/%d] [%s] %s -> %s",
|
||||
index + 1,
|
||||
len(source_order),
|
||||
source_name,
|
||||
target.canonical_id,
|
||||
provider_symbol,
|
||||
)
|
||||
record_provider_run_started(
|
||||
data_type="daily_data",
|
||||
provider=source_name,
|
||||
operation="get_daily_data",
|
||||
)
|
||||
df = self._call_fetcher_method(
|
||||
fetcher,
|
||||
"get_daily_data",
|
||||
stock_code=provider_symbol,
|
||||
start_date=start_date,
|
||||
end_date=end_date,
|
||||
days=days,
|
||||
)
|
||||
duration_ms = int((time.time() - attempt_start) * 1000)
|
||||
if df is not None and not df.empty:
|
||||
record_provider_run(
|
||||
data_type="daily_data",
|
||||
provider=source_name,
|
||||
operation="get_daily_data",
|
||||
success=True,
|
||||
latency_ms=duration_ms,
|
||||
record_count=len(df),
|
||||
)
|
||||
self._record_daily_source_success(fetcher, "cn_index")
|
||||
logger.info(
|
||||
"[指数数据源完成] %s 使用 [%s] 获取成功: rows=%d, elapsed=%.2fs",
|
||||
target.canonical_id,
|
||||
source_name,
|
||||
len(df),
|
||||
time.time() - request_start,
|
||||
)
|
||||
return df, source_name
|
||||
|
||||
reason = "empty result"
|
||||
record_provider_run(
|
||||
data_type="daily_data",
|
||||
provider=source_name,
|
||||
operation="get_daily_data",
|
||||
success=False,
|
||||
latency_ms=duration_ms,
|
||||
error_type="empty",
|
||||
error_message=reason,
|
||||
fallback_to=fallback_to,
|
||||
record_count=0,
|
||||
)
|
||||
if df is not None and df.empty:
|
||||
self._record_daily_source_success(fetcher, "cn_index")
|
||||
else:
|
||||
self._record_daily_source_failure(fetcher, "cn_index", reason)
|
||||
logger.warning(
|
||||
"[指数数据源失败 %d/%d] [%s] %s: %s",
|
||||
index + 1,
|
||||
len(source_order),
|
||||
source_name,
|
||||
target.canonical_id,
|
||||
reason,
|
||||
)
|
||||
errors.append(f"[{source_name}] {reason}")
|
||||
except Exception as exc:
|
||||
error_type, error_reason = summarize_exception(exc)
|
||||
duration_ms = int((time.time() - attempt_start) * 1000)
|
||||
record_provider_run(
|
||||
data_type="daily_data",
|
||||
provider=source_name,
|
||||
operation="get_daily_data",
|
||||
success=False,
|
||||
latency_ms=duration_ms,
|
||||
error_type=error_type,
|
||||
error_message=error_reason,
|
||||
fallback_to=fallback_to,
|
||||
record_count=0,
|
||||
)
|
||||
self._record_daily_source_failure(
|
||||
fetcher, "cn_index", error_reason
|
||||
)
|
||||
logger.warning(
|
||||
"[指数数据源失败 %d/%d] [%s] %s: error_type=%s, reason=%s",
|
||||
index + 1,
|
||||
len(source_order),
|
||||
source_name,
|
||||
target.canonical_id,
|
||||
error_type,
|
||||
error_reason,
|
||||
)
|
||||
errors.append(f"[{source_name}] ({error_type}) {error_reason}")
|
||||
|
||||
logger.warning(
|
||||
"[指数数据源终止] %s 所有指数日线数据源均失败: elapsed=%.2fs; %s",
|
||||
target.canonical_id,
|
||||
time.time() - request_start,
|
||||
"; ".join(errors) or "暂无可用数据源",
|
||||
)
|
||||
return pd.DataFrame(columns=STANDARD_COLUMNS), ""
|
||||
|
||||
def _get_cn_index_name(self, target: AnalysisTarget) -> str:
|
||||
cache_key = target.canonical_id
|
||||
entry = target.matched_index
|
||||
display_name = entry.display_name if entry is not None else ""
|
||||
if self._is_meaningful_cn_index_name(display_name, target):
|
||||
return self._cache_stock_name(cache_key, display_name) or display_name
|
||||
|
||||
cached_name = self._get_cached_stock_name(cache_key)
|
||||
if cached_name is not None:
|
||||
if self._is_meaningful_cn_index_name(cached_name, target):
|
||||
return cached_name
|
||||
self._discard_cached_stock_name(cache_key)
|
||||
|
||||
fetchers_by_name = {
|
||||
fetcher.name: fetcher for fetcher in self._get_fetchers_snapshot()
|
||||
}
|
||||
for source_name in self._CN_INDEX_NAME_SOURCE_ORDER:
|
||||
provider_symbol = self._cn_index_provider_symbol(target, source_name)
|
||||
if not provider_symbol:
|
||||
logger.warning(
|
||||
"[指数名称] [%s] 不支持 provider symbol: %s",
|
||||
source_name,
|
||||
cache_key,
|
||||
)
|
||||
continue
|
||||
fetcher = fetchers_by_name.get(source_name)
|
||||
if (
|
||||
fetcher is None
|
||||
or not hasattr(fetcher, "get_stock_name")
|
||||
or not self._is_fetcher_available(fetcher, capability="stock_name")
|
||||
):
|
||||
logger.warning(
|
||||
"[指数名称] [%s] 未配置或暂不可用: %s",
|
||||
source_name,
|
||||
cache_key,
|
||||
)
|
||||
continue
|
||||
try:
|
||||
name = self._call_fetcher_method(
|
||||
fetcher, "get_stock_name", provider_symbol
|
||||
)
|
||||
except Exception as exc:
|
||||
logger.warning(
|
||||
"[指数名称] [%s] 获取 %s 失败: %s",
|
||||
source_name,
|
||||
cache_key,
|
||||
exc,
|
||||
)
|
||||
continue
|
||||
if self._is_meaningful_cn_index_name(name, target):
|
||||
self._cache_stock_name(cache_key, name)
|
||||
logger.info(
|
||||
"[指数名称] 从 %s 获取: %s -> %s",
|
||||
source_name,
|
||||
cache_key,
|
||||
name,
|
||||
)
|
||||
return name
|
||||
logger.warning(
|
||||
"[指数名称] [%s] 获取 %s 返回空结果",
|
||||
source_name,
|
||||
cache_key,
|
||||
)
|
||||
|
||||
static_name = STOCK_NAME_MAP.get(cache_key)
|
||||
if static_name and self._is_meaningful_cn_index_name(static_name, target):
|
||||
self._cache_stock_name(cache_key, static_name)
|
||||
logger.info(
|
||||
"[指数名称] 从静态映射获取: %s -> %s",
|
||||
cache_key,
|
||||
static_name,
|
||||
)
|
||||
return static_name
|
||||
|
||||
logger.warning("[指数名称] 所有数据源都无法获取 %s 的名称", cache_key)
|
||||
return cache_key
|
||||
|
||||
def _get_tickflow_fetcher(self):
|
||||
"""Lazily create a TickFlow fetcher for market-review-only calls."""
|
||||
from src.config import get_config
|
||||
@@ -1255,10 +1615,10 @@ class DataFetcherManager:
|
||||
|
||||
故障切换策略:
|
||||
1. 美股指数/美股股票直接路由到 YfinanceFetcher
|
||||
2. 其他代码从最高优先级数据源开始尝试
|
||||
3. 捕获异常后自动切换到下一个
|
||||
4. 记录每个数据源的失败原因
|
||||
5. 所有数据源失败后抛出详细异常
|
||||
2. 当前注册表已识别的 A 股指数使用固定指数数据源链
|
||||
3. 其他代码从最高优先级数据源开始尝试
|
||||
4. 捕获异常后自动切换到下一个并记录失败原因
|
||||
5. 指数全源失败返回标准空结果;非指数全源失败抛出详细异常
|
||||
|
||||
Args:
|
||||
stock_code: 股票代码
|
||||
@@ -1270,10 +1630,21 @@ class DataFetcherManager:
|
||||
Tuple[DataFrame, str]: (数据, 成功的数据源名称)
|
||||
|
||||
Raises:
|
||||
DataFetchError: 所有数据源都失败时抛出
|
||||
DataFetchError: 非指数代码的所有数据源都失败时抛出
|
||||
"""
|
||||
from .us_index_mapping import is_us_index_code, is_us_stock_code
|
||||
|
||||
raw_stock_code = (stock_code or "").strip()
|
||||
target = parse_analysis_target(raw_stock_code)
|
||||
self._warn_bare_index_conflict(target)
|
||||
if target.asset_type == ParseStatus.INDEX:
|
||||
return self._get_cn_index_daily_data(
|
||||
target,
|
||||
start_date=start_date,
|
||||
end_date=end_date,
|
||||
days=days,
|
||||
)
|
||||
|
||||
# Normalize code (strip SH/SZ prefix etc.)
|
||||
stock_code = normalize_stock_code(stock_code)
|
||||
|
||||
@@ -2252,6 +2623,11 @@ class DataFetcherManager:
|
||||
股票中文名称,所有数据源都失败则返回 None
|
||||
"""
|
||||
raw_stock_code = (stock_code or "").strip()
|
||||
target = parse_analysis_target(raw_stock_code)
|
||||
self._warn_bare_index_conflict(target)
|
||||
if target.asset_type == ParseStatus.INDEX:
|
||||
return self._get_cn_index_name(target)
|
||||
|
||||
# Normalize code (strip SH/SZ prefix etc.)
|
||||
stock_code = normalize_stock_code(stock_code)
|
||||
static_name = STOCK_NAME_MAP.get(stock_code)
|
||||
@@ -2382,7 +2758,15 @@ class DataFetcherManager:
|
||||
"""
|
||||
if not stock_codes:
|
||||
return
|
||||
stock_codes = [normalize_stock_code(c) for c in stock_codes]
|
||||
normalized_codes: List[str] = []
|
||||
for code in stock_codes:
|
||||
target = parse_analysis_target(code)
|
||||
normalized_codes.append(
|
||||
target.canonical_id
|
||||
if target.asset_type == ParseStatus.INDEX
|
||||
else normalize_stock_code(code)
|
||||
)
|
||||
stock_codes = normalized_codes
|
||||
if use_bulk:
|
||||
self.batch_get_stock_names(stock_codes)
|
||||
return
|
||||
|
||||
@@ -34,14 +34,14 @@ class TencentFetcher(BaseFetcher):
|
||||
allow_empty_daily_data = True
|
||||
|
||||
_KLINE_ENDPOINT = "https://web.ifzq.gtimg.cn/appstock/app/fqkline/get"
|
||||
_QUOTE_ENDPOINT = "https://qt.gtimg.cn/q"
|
||||
_HTTP_TIMEOUT_SECONDS = 8
|
||||
|
||||
def __init__(self) -> None:
|
||||
self.priority = _read_tencent_priority()
|
||||
|
||||
def _fetch_raw_data(self, stock_code: str, start_date: str, end_date: str) -> pd.DataFrame:
|
||||
code = normalize_stock_code(stock_code)
|
||||
symbol = _to_tencent_symbol(code)
|
||||
symbol = _to_tencent_symbol(stock_code)
|
||||
if not symbol:
|
||||
raise DataFetchError(f"TencentFetcher unsupported stock code: {stock_code}")
|
||||
|
||||
@@ -103,11 +103,48 @@ class TencentFetcher(BaseFetcher):
|
||||
normalized = normalized[["date", "open", "high", "low", "close", "volume", "amount", "pct_chg"]]
|
||||
return normalized
|
||||
|
||||
def get_stock_name(self, stock_code: str) -> Optional[str]:
|
||||
symbol = _to_tencent_symbol(stock_code)
|
||||
if not symbol:
|
||||
return None
|
||||
try:
|
||||
response = requests.get(
|
||||
f"{self._QUOTE_ENDPOINT}={symbol}",
|
||||
headers={"Referer": "https://finance.qq.com", "User-Agent": "Mozilla/5.0"},
|
||||
timeout=self._HTTP_TIMEOUT_SECONDS,
|
||||
)
|
||||
response.raise_for_status()
|
||||
response.encoding = "gbk"
|
||||
content = response.text.strip()
|
||||
data_start = content.find('"')
|
||||
data_end = content.rfind('"')
|
||||
if data_start == -1 or data_end <= data_start:
|
||||
return None
|
||||
fields = content[data_start + 1:data_end].split("~")
|
||||
if len(fields) <= 2 or fields[2].strip() != symbol[2:]:
|
||||
return None
|
||||
name = fields[1].strip()
|
||||
return name or None
|
||||
except Exception:
|
||||
logger.debug(
|
||||
"TencentFetcher stock name lookup failed for %s",
|
||||
stock_code,
|
||||
exc_info=True,
|
||||
)
|
||||
return None
|
||||
|
||||
|
||||
def _to_tencent_symbol(stock_code: str) -> str:
|
||||
raw_code = (stock_code or "").strip().upper()
|
||||
code = normalize_stock_code(stock_code)
|
||||
if not code or not code.isdigit() or len(code) != 6:
|
||||
return ""
|
||||
if raw_code.startswith(("SH", "SS")) or raw_code.endswith((".SH", ".SS")):
|
||||
return f"sh{code}"
|
||||
if raw_code.startswith("SZ") or raw_code.endswith(".SZ"):
|
||||
return f"sz{code}"
|
||||
if raw_code.startswith("BJ") or raw_code.endswith(".BJ"):
|
||||
return f"bj{code}"
|
||||
if is_bse_code(code):
|
||||
return f"bj{code}"
|
||||
if code.startswith(("6", "5", "9")):
|
||||
|
||||
@@ -27,6 +27,7 @@ and this project adheres to [Semantic Versioning](https://semver.org/).
|
||||
- [修复] 按最新 review 复核收敛 3 处正确性问题(OR-COM-7f3d3f5b / 3d6b61f8 / a1e8b0c2):`BaseAgent._filtered_registry()` 携带源 registry 的类别超时映射(工具子集仍生效类别上限,不再绕过 #1890 类别超时);并行批次 >5 时排队调用的 per-tool 超时自 worker 实际开始起算(不再提交即烧预算导致对未启动调用的假超时);`get_tool_registry()` 缓存命中快路径在锁内读取一致对(消除与 `reset_tool_registry()` 竞态返回 `None` 或错配 registry)。新增对应回归测试。
|
||||
- [新功能] 股票名称解析引擎重构增强:新增 `resolver_name_to_code_list()` 公开 API,返回按市场排序(A 股→港股→美股)的 `Stock` 候选列表(最多 5 个),新增 `US_stock_code_match()` 匹配美股 ticker(1~5 位字母且仅限本地库已存在代码,避免 hello/open 等英文词误判为股票);AkShare 全量 A 股数据经幂等 `extend_AkShare()` 合并进全局 `stockDB`(30 分钟缓存 + 失败 5 分钟退避 + Future 单飞:TTL 过期 stale-while-revalidate 零等待、冷启动等待上界由拉取超时推导(拉取经子进程封顶 25s)、worker 先清账唤醒等待者再做日志/落盘(finally 兜底 BaseException)、成功拉取落盘 `data/cache` 跨重启复用,非中文输入跳过网络扩展),匹配策略升级为「精确→子串(≥2 汉字)→拼音子串(≥5 字母)→difflib 模糊(0.8,单字误写 0.7 兜底)」;`resolve_name_to_code()` 保持既有本地优先语义(本地精确命中零网络,调用方离线低延迟契约不变),跨市场候选能力由 `resolver_name_to_code_list()` 独立提供;解析全链路线程安全(`stockDB` 读写加锁、名称/拼音索引随库变更自动失效),新增 40 个单元测试覆盖精确/跨市场排序/子串/拼音/模糊/幂等扩展/失败退避/多候选场景。
|
||||
- [改进] `StockDaily` 表新增可空 `canonical_id` 列并支持双写(Expand-Contract PR2,issue #2207):自愈式迁移幂等加列 + 普通索引 `ix_stock_daily_canonical_id`,存量行与 `save_daily_data` 未显式传参时均经 index-aware 推导(裸指数码命中注册表时统一到指数 `canonical_id`,避免同一指数按输入形态分裂到不同桶——例如裸 `000300` 与显式 `sh000300` 现在都收敛到 `sh000300`,而非裸码被推导为 `sz000300`),推导失败写 NULL 降级;读路径仍用 `code` 列,`(code, date)` 唯一约束保留不变。显式登记契约漂移:PRD Glossary/FR-1/DD-3 与架构 AD-1/AD-7 中 canonical_id 的点分格式描述(`000016.SH`)已被 Phase 1 已合入代码的前缀格式(`sh000016`)取代,本变更遵循代码,PRD/架构文档的同步修正留待后续 PR 统一收敛。
|
||||
- [新功能] 当前注册表已识别的 5 个沪深 A 股指数以显式市场输入按 canonical 身份路由:名称优先使用注册表并以 Tencent、AkShare、TickFlow 兜底,日线固定使用 Tencent、AkShare、TickFlow、yfinance 多源降级链,不读取普通 A 股日 K 的 `*_PRIORITY` 配置;裸代码仍按股票处理,并隔离同码股票名称缓存。
|
||||
|
||||
## [3.30.0] - 2026-08-09
|
||||
|
||||
|
||||
@@ -8,8 +8,9 @@
|
||||
|
||||
如果遇到“数据源失败”,通常不是系统只能用一个源,而是免费源被限流、上游接口临时变更、网络抖动或当前市场/标的不支持。DSA 已经内置多数据源 fallback,会按场景自动尝试下一个源;如果你希望更稳定,建议至少配置一个 token 型稳定源:
|
||||
|
||||
- A 股个股与选股:优先配置 `TUSHARE_TOKEN`,并保留 AkShare / Efinance / Tencent / Baostock / YFinance 兜底。
|
||||
- A 股大盘复盘:配置 `TICKFLOW_API_KEY` 后,指数和市场宽度会优先尝试 TickFlow,失败后回退现有免费源。
|
||||
- A 股个股与选股:优先配置 `TUSHARE_TOKEN`,并保留 AkShare / Efinance / Tencent / TickFlow / Baostock / YFinance 兜底;普通个股日线按 priority 配置排序。
|
||||
- 已登记 A 股指数:固定按 Tencent → AkShare → TickFlow → YFinance 降级,不读取普通日 K 的 `*_PRIORITY` 配置。
|
||||
- A 股大盘复盘:配置 `TICKFLOW_API_KEY` 后,复盘聚合所需的指数和市场宽度会优先尝试 TickFlow,失败后回退现有免费源;这与单标的指数日线的 Tencent-first 固定链是不同入口。
|
||||
- 港股 / 美股:配置 `LONGBRIDGE_*` 后优先使用 Longbridge,YFinance、Finnhub、AlphaVantage 继续兜底。
|
||||
- 热点题材:选股的热点实现参考 AlphaSift,默认走 EastMoney provider,并使用本地 last-good cache 降低实时接口失败影响。
|
||||
|
||||
@@ -17,9 +18,10 @@
|
||||
|
||||
| 场景 | 已接入源 | 默认使用方式 | 失败处理 |
|
||||
| --- | --- | --- | --- |
|
||||
| A 股日线 / 技术面 | Efinance、Tencent、AkShare、Tushare、Pytdx、Baostock、YFinance | `DataFetcherManager` 按优先级尝试;配置 `TUSHARE_TOKEN` 后 Tushare 自动进入候选源 | 单源失败后尝试下一个源;连续失败会短期熔断该源 |
|
||||
| A 股个股日线 / 技术面 | Efinance、Tencent、AkShare、Tushare、TickFlow、Pytdx、Baostock、YFinance | `DataFetcherManager` 按 priority 尝试;配置 `TUSHARE_TOKEN` 后 Tushare 自动进入候选源 | 单源失败后尝试下一个源;连续失败会短期熔断该源 |
|
||||
| 已登记 A 股指数日线 / 技术面 | Tencent、AkShare、TickFlow、YFinance | 当前 5 个 `IndexRegistry` 标的固定按 Tencent → AkShare → TickFlow → YFinance 尝试,不读取普通日 K 的 `*_PRIORITY` 配置 | 未配置、熔断、异常或空结果均继续下一源;全部失败返回空结果并记录汇总告警 |
|
||||
| A 股实时行情 | Tencent、AkShare Sina、Efinance、AkShare EM、Tushare | `REALTIME_SOURCE_PRIORITY` 控制顺序,默认偏向 Tencent / Sina 这类轻量源 | 失败源记录 `fallback_from`,成功源继续返回 |
|
||||
| A 股大盘复盘 | TickFlow、AkShare、Tushare、Efinance | 配置 `TICKFLOW_API_KEY` 后,主指数和市场宽度优先尝试 TickFlow | TickFlow 权限不足或失败时回退 AkShare / Tushare / Efinance 链路 |
|
||||
| A 股大盘复盘 | TickFlow、AkShare、Tushare、Efinance | 配置 `TICKFLOW_API_KEY` 后,复盘聚合的主指数和市场宽度优先尝试 TickFlow;不等同于单标的指数日线链 | TickFlow 权限不足或失败时回退 AkShare / Tushare / Efinance 链路 |
|
||||
| 选股快照 | Tushare、Sina、Efinance、AkShare EM、EastMoney Datacenter | 有 `TUSHARE_TOKEN` 时自动把 `tushare` 放入快照优先级;否则使用免费源链路 | 选股引擎维护 source health;状态接口透出 snapshot/daily health |
|
||||
| 选股日线补特征 | `DataFetcherManager` | 选股引擎优先复用现有日线与缓存链路 | 现有链路失败后才回到引擎自身的日线源 |
|
||||
| 选股热点题材 | EastMoney provider、参考 AlphaSift 的 hotspot 实现、last-good cache | 未指定 provider 时默认使用 EastMoney provider | 实时失败时回退热点缓存;无缓存时返回稳定空态和可读错误 |
|
||||
@@ -39,7 +41,8 @@ flowchart TD
|
||||
D --> C[本地 stock_daily 缓存]
|
||||
C -->|命中且新鲜| COK[复用缓存]
|
||||
C -->|缺失或过期| DM{市场}
|
||||
DM -->|A 股| CN[Tushare if token -> Efinance/Tencent -> AkShare -> Pytdx -> Baostock -> YFinance]
|
||||
DM -->|A 股个股/未登记标的| CN[按 priority 动态排序: Efinance/AkShare/Tushare/TickFlow/Pytdx/Baostock/YFinance/Tencent]
|
||||
DM -->|已登记沪深指数| CNI[Tencent -> AkShare -> TickFlow -> YFinance]
|
||||
DM -->|港股| HK[Longbridge if configured -> AkShare/Tushare -> YFinance]
|
||||
DM -->|美股| US[Longbridge/YFinance -> Finnhub/AlphaVantage -> Stooq]
|
||||
|
||||
@@ -57,6 +60,7 @@ flowchart TD
|
||||
TF -->|no or failed| MF[AkShare/Tushare/Efinance fallback]
|
||||
|
||||
CN --> QL[质量标记: source/fallback/stale/fetch_failed]
|
||||
CNI --> QL
|
||||
HK --> QL
|
||||
US --> QL
|
||||
RS --> QL
|
||||
@@ -106,7 +110,7 @@ flowchart TD
|
||||
TS -->|yes| SP1[tushare -> sina -> efinance -> akshare_em -> em_datacenter]
|
||||
TS -->|no| SP2[sina -> efinance -> akshare_em -> em_datacenter]
|
||||
ENV --> DAILY[DSA provider context]
|
||||
DAILY --> DFM[DataFetcherManager: Tushare/Efinance/Tencent/AkShare/Pytdx/Baostock/YFinance]
|
||||
DAILY --> DFM[DataFetcherManager: Tushare/Efinance/Tencent/AkShare/TickFlow/Pytdx/Baostock/YFinance]
|
||||
DFM --> RESULT[候选股 + source_errors/warnings/llm_parse_errors]
|
||||
|
||||
API --> HOT{hotspots,与 screen 并行}
|
||||
@@ -132,7 +136,7 @@ ENABLE_EASTMONEY_PATCH=true
|
||||
|
||||
### A 股稳定模式
|
||||
|
||||
适合经常跑选股、批量分析或对外服务。Tushare 用于增强 A 股日线与快照稳定性;TickFlow 可增强 A 股日 K、实时行情和大盘复盘(实时行情需显式加入 `REALTIME_SOURCE_PRIORITY`);免费源继续作为兜底。
|
||||
适合经常跑选股、批量分析或对外服务。Tushare 用于增强普通 A 股日线与快照稳定性;TickFlow 可增强普通 A 股日 K、已登记指数固定 fallback、实时行情和大盘复盘(实时行情需显式加入 `REALTIME_SOURCE_PRIORITY`);免费源继续作为兜底。
|
||||
|
||||
```env
|
||||
TUSHARE_TOKEN=your_tushare_token
|
||||
@@ -176,7 +180,9 @@ LONGBRIDGE_ACCESS_TOKEN=your_access_token
|
||||
| --- | --- |
|
||||
| 单个源失败但 fallback 成功 | 本次使用了降级数据源,分析仍可继续;报告中会标记实际成功源。 |
|
||||
| 多个源失败但有缓存 | 实时源不可用,本次使用上一次成功缓存;结论会降低置信度。 |
|
||||
| 全部源失败且无缓存 | 当前数据不可用,请稍后重试,或配置 Tushare / TickFlow / Longbridge 等 token 型数据源。 |
|
||||
| 全部源失败且无缓存 | 当前数据不可用,请稍后重试。普通 A 股可检查 Tushare/TickFlow;已登记指数应检查 Tencent/AkShare 连通性或配置 TickFlow;港美股可检查 Longbridge。 |
|
||||
|
||||
普通 A 股日线使用 `cn` 健康度命名空间,已登记指数固定链使用独立的 `cn_index` 命名空间;任一链的连续失败不会熔断另一条链。指数源返回空结果时会继续 fallback 并记录诊断,四源全部失败则返回空结果和空来源;普通股票保留既有最终异常语义。
|
||||
|
||||
### 新闻面证据缺失的报告标注(已实现)
|
||||
|
||||
|
||||
+5
-3
@@ -418,8 +418,8 @@ daily_stock_analysis/
|
||||
| `TUSHARE_TOKEN` | Tushare Pro Token | - | 可选 |
|
||||
| `TUSHARE_HTTP_URL` | Tushare Pro HTTP 接入地址;留空时使用官方端点 `http://api.tushare.pro`,仅在需通过公司内网代理、跨境网络或自建镜像时填 `http://` 或 `https://` 开头的完整地址 | `http://api.tushare.pro` | 可选 |
|
||||
| `TICKFLOW_API_KEY` | TickFlow API Key;可选,用于 A 股日 K、实时行情、股票列表/名称与大盘复盘增强;失败或权限不足时自动回退。 | - | 可选 |
|
||||
| `TICKFLOW_PRIORITY` | TickFlow 日 K 数据源优先级;数字越小越早尝试,默认 `2`;未配置 API Key 时不启用;不影响实时行情,实时行情顺序由 `REALTIME_SOURCE_PRIORITY` 控制。 | `2` | 可选 |
|
||||
| `TENCENT_PRIORITY` | 腾讯直连 A 股日 K 数据源优先级;数字越小越早尝试,默认 `5`,作为 Efinance、AkShare、Tushare、TickFlow、PyTDX、Baostock 和 YFinance 之后的最终兜底;不影响实时行情。 | `5` | 可选 |
|
||||
| `TICKFLOW_PRIORITY` | TickFlow 普通 A 股日 K 数据源优先级;数字越小越早尝试,默认 `2`;未配置 API Key 时不启用;已登记指数使用独立固定链,不读取本变量;不影响实时行情。 | `2` | 可选 |
|
||||
| `TENCENT_PRIORITY` | 腾讯直连普通 A 股日 K 数据源优先级;数字越小越早尝试,默认 `5`,作为通用链路最终兜底;已登记指数使用独立固定链,不读取本变量;不影响实时行情。 | `5` | 可选 |
|
||||
| `TICKFLOW_KLINE_ADJUST` | TickFlow 日 K 复权模式:`none`、`forward`、`backward`、`forward_additive`、`backward_additive`。 | `none` | 可选 |
|
||||
| `TICKFLOW_BATCH_DAILY_ENABLED` | 是否启用 TickFlow 批量日 K 预取;权限不足会短期缓存失败状态,并继续走常规回退。 | `true` | 可选 |
|
||||
| `TICKFLOW_BATCH_SIZE` | TickFlow 日 K 与实时行情批量请求的单批最大标的数。 | `100` | 可选 |
|
||||
@@ -448,7 +448,9 @@ daily_stock_analysis/
|
||||
> - 日股/韩股:当前仅走 Yfinance 基础路径获取日线与实时行情;`institution`、`capital_flow`、`dragon_tiger`、`boards` 等依赖 A 股专属源/离岸完整版的能力会降级为 `not_supported`(详见 [市场支持与边界](market-support.md));
|
||||
> - 台股:在美股/港股 offshore 基础路径之外,`institution` 区块额外展示三大法人原始买卖超净额(TWSE T86 / TPEx,默认开启、fail-open,取不到数据时维持 `not_supported`);`capital_flow`、`dragon_tiger`、`boards` 仍为 `not_supported`;
|
||||
> - 任何异常走 fail-open,仅记录错误,不影响技术面/新闻/筹码主链路。
|
||||
> - 配置 `TICKFLOW_API_KEY` 后,TickFlow 会作为可选 A 股日 K 数据源和大盘复盘增强源实例化;`TICKFLOW_PRIORITY` 只影响日 K/通用数据源回退链。实时行情优先级由 `REALTIME_SOURCE_PRIORITY` 单独控制,只有显式包含 `tickflow` 时才会使用 TickFlow 实时行情。`REALTIME_SOURCE_PRIORITY` 中排在 `tickflow` 前面的数据源会先被尝试。
|
||||
> - 配置 `TICKFLOW_API_KEY` 后,TickFlow 会作为可选 A 股日 K 数据源和大盘复盘增强源实例化;`TICKFLOW_PRIORITY` 只影响普通 A 股日 K/通用数据源回退链。实时行情优先级由 `REALTIME_SOURCE_PRIORITY` 单独控制,只有显式包含 `tickflow` 时才会使用 TickFlow 实时行情。`REALTIME_SOURCE_PRIORITY` 中排在 `tickflow` 前面的数据源会先被尝试。
|
||||
> - 当前 `IndexRegistry` 已登记的 5 个沪深指数为 `sh000016`(上证50)、`sh000688`(科创50)、`sz399001`(深证成指)、`sz399006`(创业板指)和 `sh000300`(沪深300)。使用显式市场输入(也接受 `000016.SH` 等交易所后缀形式)时,它们不参与通用 priority 排序,固定按 Tencent → AkShare → TickFlow → YFinance 降级;未配置或不可用的来源会跳过。裸 `000016` 等代码仍按股票处理,不触发指数链。该固定链不读取 `EFINANCE_PRIORITY`、`AKSHARE_PRIORITY`、`TUSHARE_PRIORITY`、`TICKFLOW_PRIORITY`、`PYTDX_PRIORITY`、`BAOSTOCK_PRIORITY`、`YFINANCE_PRIORITY` 或 `TENCENT_PRIORITY`,不影响普通股票和实时行情的既有顺序。
|
||||
> - 已登记指数名称优先来自本地注册表;只有注册名称无效时才按 Tencent → AkShare → TickFlow 查询,名称链不使用 YFinance。指数日线四源全部失败时返回空结果并记录汇总告警;普通股票仍保持既有的 `DataFetchError` 最终失败契约。
|
||||
> - TickFlow 日 K 默认 `TICKFLOW_KLINE_ADJUST=none`;日线 `volume` 从手统一转为股,`amount` 保持元口径。
|
||||
> - TickFlow 日 K 区间请求会显式传入 `start_time` / `end_time` / `count`;官方 quickstart 明确说明时间范围查询仍受 `count` 限制。若返回非空但行数打满 `count` 且首个返回交易日晚于请求起始交易日,系统会判定为疑似截断,不写入缓存并让 manager 继续回退。
|
||||
> - 批量分析时,`prefetch_daily_klines()` 会在逐股 `get_daily_data()` 之前预热进程内缓存,不改变对外调用路径。
|
||||
|
||||
@@ -359,8 +359,8 @@ For the notification baseline, diagnostics, and deployment notes, see [Notificat
|
||||
| `TUSHARE_TOKEN` | Tushare Pro Token | - | Optional |
|
||||
| `TUSHARE_HTTP_URL` | Tushare Pro HTTP endpoint; defaults to `http://api.tushare.pro` when unset/empty. Set only when routing through a corporate proxy, cross-border network, or a self-hosted mirror (must start with `http://` or `https://`). | `http://api.tushare.pro` | Optional |
|
||||
| `TICKFLOW_API_KEY` | TickFlow API key; enables optional A-share daily K-lines, realtime quotes, stock list/name lookup, and CN market review enhancement. Permission failures fall back to existing providers. | - | Optional |
|
||||
| `TICKFLOW_PRIORITY` | TickFlow daily K-line provider priority; lower values are tried earlier. No effect unless `TICKFLOW_API_KEY` is configured. Does not affect realtime quotes, which are ordered by `REALTIME_SOURCE_PRIORITY`. | `2` | Optional |
|
||||
| `TENCENT_PRIORITY` | Tencent direct A-share daily K-line provider priority; lower values are tried earlier. Defaults to `5` as the last fallback after Efinance, AkShare, Tushare, TickFlow, PyTDX, Baostock, and YFinance. Does not affect realtime quotes. | `5` | Optional |
|
||||
| `TICKFLOW_PRIORITY` | TickFlow priority for the generic A-share daily K-line route; lower values are tried earlier. No effect unless `TICKFLOW_API_KEY` is configured. Registered indices use a separate fixed chain and ignore this variable. Realtime quotes are ordered by `REALTIME_SOURCE_PRIORITY`. | `2` | Optional |
|
||||
| `TENCENT_PRIORITY` | Tencent direct priority for the generic A-share daily K-line route; lower values are tried earlier and `5` is the default last fallback. Registered indices use a separate fixed chain and ignore this variable. Does not affect realtime quotes. | `5` | Optional |
|
||||
| `TICKFLOW_KLINE_ADJUST` | TickFlow daily K-line adjustment mode: `none`, `forward`, `backward`, `forward_additive`, or `backward_additive`. | `none` | Optional |
|
||||
| `TICKFLOW_BATCH_DAILY_ENABLED` | Enable TickFlow batch daily K-line prefetch when the current plan supports it; permission failures are negative-cached and fall back to per-stock providers. | `true` | Optional |
|
||||
| `TICKFLOW_BATCH_SIZE` | Maximum symbols per TickFlow batch request for daily K-lines and realtime quotes. | `100` | Optional |
|
||||
@@ -420,7 +420,9 @@ For the notification baseline, diagnostics, and deployment notes, see [Notificat
|
||||
| `SAVE_CONTEXT_SNAPSHOT` | Persist analysis-history `context_snapshot`. When false, new history records do not save enhanced_context, market_phase_summary, AnalysisContextPack overview, or diagnostic snapshots, but current-run prompt summaries remain enabled | `true` |
|
||||
|
||||
> Behavior notes:
|
||||
> - When `TICKFLOW_API_KEY` is configured, TickFlow is instantiated as an optional A-share daily K-line data source and CN market-review enhancer. `TICKFLOW_PRIORITY` only affects the daily K-line/general provider fallback chain. Realtime quote priority is controlled separately by `REALTIME_SOURCE_PRIORITY`; TickFlow realtime quotes are used only when that list explicitly includes `tickflow`, and any source listed before `tickflow` is tried first.
|
||||
> - When `TICKFLOW_API_KEY` is configured, TickFlow is instantiated as an optional A-share daily K-line data source and CN market-review enhancer. `TICKFLOW_PRIORITY` only affects the generic A-share daily K-line/provider fallback chain. Realtime quote priority is controlled separately by `REALTIME_SOURCE_PRIORITY`; TickFlow realtime quotes are used only when that list explicitly includes `tickflow`, and any source listed before `tickflow` is tried first.
|
||||
> - The five SH/SZ indices currently registered in `IndexRegistry` are `sh000016` (SSE 50), `sh000688` (STAR 50), `sz399001` (SZSE Component), `sz399006` (ChiNext), and `sh000300` (CSI 300). Explicit-market inputs (exchange-suffix forms such as `000016.SH` are also accepted) bypass generic priority sorting and use the fixed Tencent → AkShare → TickFlow → YFinance fallback chain; unconfigured or unavailable providers are skipped. Bare `000016`-style inputs remain stocks and do not enter the index chain. This fixed chain ignores `EFINANCE_PRIORITY`, `AKSHARE_PRIORITY`, `TUSHARE_PRIORITY`, `TICKFLOW_PRIORITY`, `PYTDX_PRIORITY`, `BAOSTOCK_PRIORITY`, `YFINANCE_PRIORITY`, and `TENCENT_PRIORITY`. Existing stock and realtime-quote ordering is unchanged.
|
||||
> - Registered index names normally come from the local registry. Only an invalid registry name triggers the Tencent → AkShare → TickFlow name fallback; YFinance is not part of the name chain. If all four daily providers fail, an index request returns an empty result with a summary warning, while ordinary stocks retain the existing final `DataFetchError` contract.
|
||||
> - TickFlow daily K-lines default to `TICKFLOW_KLINE_ADJUST=none`; daily `volume` is converted from lots to shares, while `amount` remains in yuan.
|
||||
> - TickFlow daily K-line range requests pass explicit `start_time` / `end_time` / `count`. Because the official quickstart documents that time-range queries are still limited by `count`, non-empty count-capped responses whose first returned trading date is later than the requested start trading date are rejected before normalization or cache writes, allowing manager fallback to continue.
|
||||
> - Batch analysis can warm the per-process TickFlow daily K-line cache through `prefetch_daily_klines()` before per-stock `get_daily_data()` calls. Only validated frames are cached; batch permission failures are negative-cached and degrade to single-stock requests or existing providers.
|
||||
|
||||
@@ -794,7 +794,7 @@ _FIELD_DEFINITIONS: Dict[str, Dict[str, Any]] = {
|
||||
},
|
||||
"TICKFLOW_PRIORITY": {
|
||||
"title": "TickFlow Daily K-line Priority",
|
||||
"description": "Priority for TickFlow daily K-line fetcher. Lower numbers are tried earlier; realtime quote order is controlled separately by REALTIME_SOURCE_PRIORITY.",
|
||||
"description": "Priority for TickFlow in the generic A-share daily K-line route. Registered indices use the fixed Tencent -> AkShare -> TickFlow -> YFinance chain and ignore this setting; realtime quote order is controlled separately by REALTIME_SOURCE_PRIORITY.",
|
||||
"category": "data_source",
|
||||
"data_type": "integer",
|
||||
"ui_control": "number",
|
||||
|
||||
@@ -0,0 +1,565 @@
|
||||
import logging
|
||||
from typing import cast
|
||||
from unittest.mock import patch
|
||||
|
||||
import pandas as pd
|
||||
import pytest
|
||||
|
||||
from data_provider.base import BaseFetcher, DataFetcherManager, STANDARD_COLUMNS
|
||||
from src.services.stock_list_parser import (
|
||||
AnalysisTarget,
|
||||
IndexEntry,
|
||||
ParseStatus,
|
||||
parse_analysis_target,
|
||||
)
|
||||
|
||||
|
||||
def _daily_frame() -> pd.DataFrame:
|
||||
return pd.DataFrame(
|
||||
{
|
||||
"date": [pd.Timestamp("2026-08-21")],
|
||||
"open": [100.0],
|
||||
"high": [102.0],
|
||||
"low": [99.0],
|
||||
"close": [101.0],
|
||||
"volume": [1000],
|
||||
"amount": [101000.0],
|
||||
"pct_chg": [1.0],
|
||||
}
|
||||
)
|
||||
|
||||
|
||||
class _FakeFetcher(BaseFetcher):
|
||||
def __init__(
|
||||
self,
|
||||
name: str,
|
||||
*,
|
||||
priority: int,
|
||||
daily_result=None,
|
||||
name_result=None,
|
||||
available: bool = True,
|
||||
) -> None:
|
||||
self.name = name
|
||||
self.priority = priority
|
||||
self.daily_result = daily_result
|
||||
self.name_result = name_result
|
||||
self.available = available
|
||||
self.daily_calls: list[str] = []
|
||||
self.name_calls: list[str] = []
|
||||
|
||||
def _fetch_raw_data(self, stock_code, start_date, end_date):
|
||||
raise NotImplementedError
|
||||
|
||||
def _normalize_data(self, df, stock_code):
|
||||
raise NotImplementedError
|
||||
|
||||
def is_available_for_request(self, _capability: str) -> bool:
|
||||
return self.available
|
||||
|
||||
def get_daily_data(self, stock_code, start_date=None, end_date=None, days=30):
|
||||
self.daily_calls.append(stock_code)
|
||||
if isinstance(self.daily_result, Exception):
|
||||
raise self.daily_result
|
||||
if self.daily_result is None:
|
||||
return pd.DataFrame()
|
||||
return self.daily_result.copy()
|
||||
|
||||
def get_stock_name(self, stock_code):
|
||||
self.name_calls.append(stock_code)
|
||||
if isinstance(self.name_result, Exception):
|
||||
raise self.name_result
|
||||
return self.name_result
|
||||
|
||||
|
||||
def _manager_without_fetchers() -> DataFetcherManager:
|
||||
manager = DataFetcherManager.__new__(DataFetcherManager)
|
||||
manager._fetchers = []
|
||||
manager._ensure_concurrency_guards()
|
||||
return manager
|
||||
|
||||
|
||||
@pytest.fixture(autouse=True)
|
||||
def _reset_daily_source_health():
|
||||
DataFetcherManager.reset_daily_source_health()
|
||||
yield
|
||||
DataFetcherManager.reset_daily_source_health()
|
||||
|
||||
|
||||
def test_index_daily_route_uses_fixed_order_symbols_and_diagnostics() -> None:
|
||||
tencent = _FakeFetcher(
|
||||
"TencentFetcher", priority=9, daily_result=RuntimeError("tencent failed")
|
||||
)
|
||||
akshare = _FakeFetcher("AkshareFetcher", priority=8, daily_result=pd.DataFrame())
|
||||
tickflow = _FakeFetcher(
|
||||
"TickFlowFetcher", priority=7, daily_result=_daily_frame(), available=False
|
||||
)
|
||||
yfinance = _FakeFetcher("YfinanceFetcher", priority=0, daily_result=_daily_frame())
|
||||
manager = DataFetcherManager(fetchers=[yfinance, tickflow, akshare, tencent])
|
||||
|
||||
with patch("data_provider.base.record_provider_run") as record_run:
|
||||
df, source = manager.get_daily_data(
|
||||
"sh000016", start_date="2026-08-01", end_date="2026-08-21"
|
||||
)
|
||||
|
||||
assert not df.empty
|
||||
assert source == "YfinanceFetcher"
|
||||
assert tencent.daily_calls == ["sh000016"]
|
||||
assert akshare.daily_calls == ["sh000016"]
|
||||
assert tickflow.daily_calls == []
|
||||
assert yfinance.daily_calls == ["000016.SS"]
|
||||
assert [item.kwargs["provider"] for item in record_run.call_args_list] == [
|
||||
"TencentFetcher",
|
||||
"AkshareFetcher",
|
||||
"TickFlowFetcher",
|
||||
"YfinanceFetcher",
|
||||
]
|
||||
assert [item.kwargs.get("fallback_to") for item in record_run.call_args_list] == [
|
||||
"AkshareFetcher",
|
||||
"TickFlowFetcher",
|
||||
"YfinanceFetcher",
|
||||
None,
|
||||
]
|
||||
assert [
|
||||
(
|
||||
item.kwargs["success"],
|
||||
item.kwargs.get("error_type"),
|
||||
item.kwargs.get("error_message"),
|
||||
item.kwargs.get("record_count"),
|
||||
)
|
||||
for item in record_run.call_args_list
|
||||
] == [
|
||||
(False, "RuntimeError", "tencent failed", 0),
|
||||
(False, "empty", "empty result", 0),
|
||||
(False, "unavailable", "数据源未配置或暂不可用", 0),
|
||||
(True, None, None, 1),
|
||||
]
|
||||
|
||||
|
||||
def test_index_daily_route_short_circuits_after_tencent_success() -> None:
|
||||
tencent = _FakeFetcher("TencentFetcher", priority=9, daily_result=_daily_frame())
|
||||
akshare = _FakeFetcher("AkshareFetcher", priority=0, daily_result=_daily_frame())
|
||||
tickflow = _FakeFetcher("TickFlowFetcher", priority=1, daily_result=_daily_frame())
|
||||
yfinance = _FakeFetcher("YfinanceFetcher", priority=2, daily_result=_daily_frame())
|
||||
manager = DataFetcherManager(fetchers=[akshare, tickflow, yfinance, tencent])
|
||||
|
||||
with patch("data_provider.base.record_provider_run") as record_run:
|
||||
df, source = manager.get_daily_data("sz399001")
|
||||
|
||||
assert not df.empty
|
||||
assert source == "TencentFetcher"
|
||||
assert tencent.daily_calls == ["sz399001"]
|
||||
assert akshare.daily_calls == []
|
||||
assert tickflow.daily_calls == []
|
||||
assert yfinance.daily_calls == []
|
||||
assert record_run.call_args.kwargs.get("fallback_to") is None
|
||||
|
||||
|
||||
def test_index_daily_route_short_circuits_after_akshare_success() -> None:
|
||||
tencent = _FakeFetcher("TencentFetcher", priority=0, daily_result=pd.DataFrame())
|
||||
akshare = _FakeFetcher("AkshareFetcher", priority=1, daily_result=_daily_frame())
|
||||
tickflow = _FakeFetcher("TickFlowFetcher", priority=2, daily_result=_daily_frame())
|
||||
yfinance = _FakeFetcher("YfinanceFetcher", priority=3, daily_result=_daily_frame())
|
||||
manager = DataFetcherManager(fetchers=[tencent, akshare, tickflow, yfinance])
|
||||
|
||||
_, source = manager.get_daily_data("sh000016")
|
||||
|
||||
assert source == "AkshareFetcher"
|
||||
assert tencent.daily_calls == ["sh000016"]
|
||||
assert akshare.daily_calls == ["sh000016"]
|
||||
assert tickflow.daily_calls == []
|
||||
assert yfinance.daily_calls == []
|
||||
|
||||
|
||||
def test_index_daily_route_short_circuits_after_tickflow_success() -> None:
|
||||
tencent = _FakeFetcher("TencentFetcher", priority=0, daily_result=pd.DataFrame())
|
||||
akshare = _FakeFetcher("AkshareFetcher", priority=1, daily_result=pd.DataFrame())
|
||||
tickflow = _FakeFetcher("TickFlowFetcher", priority=2, daily_result=_daily_frame())
|
||||
yfinance = _FakeFetcher("YfinanceFetcher", priority=3, daily_result=_daily_frame())
|
||||
manager = DataFetcherManager(fetchers=[tencent, akshare, tickflow, yfinance])
|
||||
|
||||
_, source = manager.get_daily_data("sh000016")
|
||||
|
||||
assert source == "TickFlowFetcher"
|
||||
assert tencent.daily_calls == ["sh000016"]
|
||||
assert akshare.daily_calls == ["sh000016"]
|
||||
assert tickflow.daily_calls == ["000016.SH"]
|
||||
assert yfinance.daily_calls == []
|
||||
|
||||
|
||||
def test_index_provider_symbols_cover_shanghai_and_shenzhen() -> None:
|
||||
sh_target = parse_analysis_target("sh000016")
|
||||
sz_target = parse_analysis_target("sz399001")
|
||||
|
||||
assert DataFetcherManager._cn_index_provider_symbol(
|
||||
sh_target, "TencentFetcher"
|
||||
) == "sh000016"
|
||||
assert DataFetcherManager._cn_index_provider_symbol(
|
||||
sh_target, "TickFlowFetcher"
|
||||
) == "000016.SH"
|
||||
assert DataFetcherManager._cn_index_provider_symbol(
|
||||
sh_target, "YfinanceFetcher"
|
||||
) == "000016.SS"
|
||||
assert DataFetcherManager._cn_index_provider_symbol(
|
||||
sz_target, "AkshareFetcher"
|
||||
) == "sz399001"
|
||||
assert DataFetcherManager._cn_index_provider_symbol(
|
||||
sz_target, "TickFlowFetcher"
|
||||
) == "399001.SZ"
|
||||
assert DataFetcherManager._cn_index_provider_symbol(
|
||||
sz_target, "YfinanceFetcher"
|
||||
) == "399001.SZ"
|
||||
|
||||
|
||||
def test_index_daily_health_is_isolated_from_regular_cn_route() -> None:
|
||||
tencent = _FakeFetcher("TencentFetcher", priority=0, daily_result=_daily_frame())
|
||||
manager = DataFetcherManager(fetchers=[tencent])
|
||||
regular_cn_key = manager._daily_health_key(tencent, "cn")
|
||||
for _ in range(3):
|
||||
manager._daily_source_health.record_failure(regular_cn_key, error="failed")
|
||||
|
||||
_, source = manager.get_daily_data("sh000016")
|
||||
|
||||
assert source == "TencentFetcher"
|
||||
assert tencent.daily_calls == ["sh000016"]
|
||||
|
||||
|
||||
def test_regular_cn_daily_health_is_isolated_from_index_route() -> None:
|
||||
tencent = _FakeFetcher(
|
||||
"TencentFetcher", priority=0, daily_result=RuntimeError("index failed")
|
||||
)
|
||||
manager = DataFetcherManager(fetchers=[tencent])
|
||||
index_key = manager._daily_health_key(tencent, "cn_index")
|
||||
|
||||
for _ in range(3):
|
||||
result, source = manager.get_daily_data("sh000016")
|
||||
assert result.empty
|
||||
assert source == ""
|
||||
|
||||
assert not manager._daily_source_health.is_available(index_key)
|
||||
assert manager._daily_source_health.is_available(
|
||||
manager._daily_health_key(tencent, "cn")
|
||||
)
|
||||
|
||||
tencent.daily_result = _daily_frame()
|
||||
_, source = manager.get_daily_data("600519")
|
||||
|
||||
assert source == "TencentFetcher"
|
||||
assert tencent.daily_calls == ["sh000016", "sh000016", "sh000016", "600519"]
|
||||
|
||||
|
||||
def test_index_daily_skips_unsupported_provider_symbols_without_health_failure() -> None:
|
||||
entry = IndexEntry(
|
||||
bare_code="000016",
|
||||
exchange="HK",
|
||||
canonical_id="hk000016",
|
||||
display_name="",
|
||||
)
|
||||
target = AnalysisTarget(
|
||||
raw_input="hk000016",
|
||||
asset_type=ParseStatus.INDEX,
|
||||
canonical_id="hk000016",
|
||||
display_code="000016.HK",
|
||||
exchange="HK",
|
||||
matched_index=entry,
|
||||
)
|
||||
fetchers = [
|
||||
_FakeFetcher("TencentFetcher", priority=0, daily_result=_daily_frame()),
|
||||
_FakeFetcher("AkshareFetcher", priority=1, daily_result=_daily_frame()),
|
||||
_FakeFetcher("TickFlowFetcher", priority=2, daily_result=_daily_frame()),
|
||||
_FakeFetcher("YfinanceFetcher", priority=3, daily_result=_daily_frame()),
|
||||
]
|
||||
manager = DataFetcherManager(fetchers=cast(list[BaseFetcher], fetchers))
|
||||
|
||||
with patch("data_provider.base.parse_analysis_target", return_value=target), patch(
|
||||
"data_provider.base.record_provider_run"
|
||||
) as record_run, patch.object(
|
||||
DataFetcherManager, "_record_daily_source_failure"
|
||||
) as record_health_failure:
|
||||
df, source = manager.get_daily_data("hk000016")
|
||||
|
||||
assert df.empty
|
||||
assert source == ""
|
||||
assert all(fetcher.daily_calls == [] for fetcher in fetchers)
|
||||
assert [item.kwargs["error_type"] for item in record_run.call_args_list] == [
|
||||
"unsupported",
|
||||
"unsupported",
|
||||
"unsupported",
|
||||
"unsupported",
|
||||
]
|
||||
assert all(
|
||||
"unsupported index provider symbol" in item.kwargs["error_message"]
|
||||
for item in record_run.call_args_list
|
||||
)
|
||||
assert all(item.kwargs["record_count"] == 0 for item in record_run.call_args_list)
|
||||
record_health_failure.assert_not_called()
|
||||
|
||||
|
||||
def test_index_daily_route_returns_standard_empty_result_when_all_sources_fail(
|
||||
caplog,
|
||||
) -> None:
|
||||
fetchers = [
|
||||
_FakeFetcher("TencentFetcher", priority=0, daily_result=RuntimeError("boom")),
|
||||
_FakeFetcher("AkshareFetcher", priority=1, daily_result=pd.DataFrame()),
|
||||
_FakeFetcher("TickFlowFetcher", priority=2, daily_result=RuntimeError("offline")),
|
||||
_FakeFetcher("YfinanceFetcher", priority=3, daily_result=pd.DataFrame()),
|
||||
]
|
||||
manager = DataFetcherManager(fetchers=cast(list[BaseFetcher], fetchers))
|
||||
|
||||
with caplog.at_level(logging.WARNING):
|
||||
df, source = manager.get_daily_data("sh000688")
|
||||
|
||||
assert df.empty
|
||||
assert list(df.columns) == STANDARD_COLUMNS
|
||||
assert source == ""
|
||||
assert [fetcher.daily_calls for fetcher in fetchers] == [
|
||||
["sh000688"],
|
||||
["sh000688"],
|
||||
["000688.SH"],
|
||||
["000688.SS"],
|
||||
]
|
||||
assert "sh000688" in caplog.text
|
||||
assert "所有指数日线数据源" in caplog.text
|
||||
|
||||
|
||||
@pytest.mark.parametrize("stock_code", ["000016", "000001", "000688"])
|
||||
def test_bare_conflict_codes_keep_stock_route_and_warn(stock_code, caplog) -> None:
|
||||
efinance = _FakeFetcher("EfinanceFetcher", priority=0, daily_result=_daily_frame())
|
||||
tencent = _FakeFetcher("TencentFetcher", priority=1, daily_result=_daily_frame())
|
||||
manager = DataFetcherManager(fetchers=[efinance, tencent])
|
||||
|
||||
with caplog.at_level(logging.WARNING):
|
||||
_, source = manager.get_daily_data(stock_code)
|
||||
|
||||
assert source == "EfinanceFetcher"
|
||||
assert efinance.daily_calls == [stock_code]
|
||||
assert tencent.daily_calls == []
|
||||
assert "裸代码" in caplog.text
|
||||
assert "股票路由" in caplog.text
|
||||
|
||||
|
||||
@pytest.mark.parametrize(
|
||||
("stock_code", "stock_name"),
|
||||
[
|
||||
("000016", "深康佳A"),
|
||||
("000001", "平安银行"),
|
||||
("000688", "国城矿业"),
|
||||
],
|
||||
)
|
||||
def test_bare_conflict_get_stock_name_keeps_stock_name_and_warns(
|
||||
stock_code, stock_name, caplog
|
||||
) -> None:
|
||||
manager = _manager_without_fetchers()
|
||||
|
||||
with patch.dict(
|
||||
"data_provider.base.STOCK_NAME_MAP", {stock_code: stock_name}, clear=True
|
||||
), caplog.at_level(logging.WARNING):
|
||||
name = manager.get_stock_name(stock_code, allow_realtime=False)
|
||||
|
||||
assert name == stock_name
|
||||
assert manager._stock_name_cache == {stock_code: stock_name}
|
||||
assert "裸代码" in caplog.text
|
||||
assert "股票路由" in caplog.text
|
||||
|
||||
|
||||
def test_non_index_prefixed_stock_keeps_existing_normalized_route() -> None:
|
||||
efinance = _FakeFetcher("EfinanceFetcher", priority=0, daily_result=_daily_frame())
|
||||
tencent = _FakeFetcher("TencentFetcher", priority=1, daily_result=_daily_frame())
|
||||
manager = DataFetcherManager(fetchers=[efinance, tencent])
|
||||
|
||||
_, source = manager.get_daily_data("sh600519")
|
||||
|
||||
assert source == "EfinanceFetcher"
|
||||
assert efinance.daily_calls == ["600519"]
|
||||
assert tencent.daily_calls == []
|
||||
|
||||
|
||||
def test_index_name_uses_registry_name_and_canonical_cache_without_network() -> None:
|
||||
tencent = _FakeFetcher("TencentFetcher", priority=2, name_result="错误名称")
|
||||
akshare = _FakeFetcher("AkshareFetcher", priority=1, name_result="错误名称")
|
||||
tickflow = _FakeFetcher("TickFlowFetcher", priority=0, name_result="错误名称")
|
||||
manager = DataFetcherManager(fetchers=[tickflow, akshare, tencent])
|
||||
manager._stock_name_cache["sh000016"] = "旧错误名称"
|
||||
|
||||
name = manager.get_stock_name("sh000016", allow_realtime=False)
|
||||
|
||||
assert name == "上证50"
|
||||
assert manager._stock_name_cache == {"sh000016": "上证50"}
|
||||
assert tencent.name_calls == []
|
||||
assert akshare.name_calls == []
|
||||
assert tickflow.name_calls == []
|
||||
|
||||
|
||||
def test_index_name_fallback_uses_fixed_order_and_provider_symbols() -> None:
|
||||
entry = IndexEntry(
|
||||
bare_code="000016",
|
||||
exchange="SH",
|
||||
canonical_id="sh000016",
|
||||
display_name="",
|
||||
)
|
||||
target = AnalysisTarget(
|
||||
raw_input="sh000016",
|
||||
asset_type=ParseStatus.INDEX,
|
||||
canonical_id="sh000016",
|
||||
display_code="",
|
||||
exchange="SH",
|
||||
matched_index=entry,
|
||||
)
|
||||
tencent = _FakeFetcher("TencentFetcher", priority=9, name_result=None)
|
||||
akshare = _FakeFetcher("AkshareFetcher", priority=8, name_result="")
|
||||
tickflow = _FakeFetcher("TickFlowFetcher", priority=0, name_result="网络上证50")
|
||||
manager = DataFetcherManager(fetchers=[tickflow, akshare, tencent])
|
||||
|
||||
with patch("data_provider.base.parse_analysis_target", return_value=target):
|
||||
name = manager.get_stock_name("sh000016", allow_realtime=False)
|
||||
|
||||
assert name == "网络上证50"
|
||||
assert tencent.name_calls == ["sh000016"]
|
||||
assert akshare.name_calls == ["sh000016"]
|
||||
assert tickflow.name_calls == ["000016.SH"]
|
||||
assert manager._stock_name_cache == {"sh000016": "网络上证50"}
|
||||
|
||||
|
||||
def test_index_name_static_fallback_uses_only_canonical_key() -> None:
|
||||
entry = IndexEntry(
|
||||
bare_code="000016",
|
||||
exchange="SH",
|
||||
canonical_id="sh000016",
|
||||
display_name="",
|
||||
)
|
||||
target = AnalysisTarget(
|
||||
raw_input="sh000016",
|
||||
asset_type=ParseStatus.INDEX,
|
||||
canonical_id="sh000016",
|
||||
display_code="",
|
||||
exchange="SH",
|
||||
matched_index=entry,
|
||||
)
|
||||
tencent = _FakeFetcher("TencentFetcher", priority=2, name_result=None)
|
||||
akshare = _FakeFetcher("AkshareFetcher", priority=1, name_result="")
|
||||
tickflow = _FakeFetcher("TickFlowFetcher", priority=0, name_result=None)
|
||||
manager = DataFetcherManager(fetchers=[tickflow, akshare, tencent])
|
||||
|
||||
with patch.dict(
|
||||
"data_provider.base.STOCK_NAME_MAP",
|
||||
{"sh000016": "静态上证50", "000016": "深康佳A"},
|
||||
clear=True,
|
||||
):
|
||||
with patch("data_provider.base.parse_analysis_target", return_value=target):
|
||||
index_name = manager.get_stock_name("sh000016", allow_realtime=False)
|
||||
stock_name = manager.get_stock_name("000016", allow_realtime=False)
|
||||
|
||||
assert index_name == "静态上证50"
|
||||
assert stock_name == "深康佳A"
|
||||
assert tencent.name_calls == ["sh000016"]
|
||||
assert akshare.name_calls == ["sh000016"]
|
||||
assert tickflow.name_calls == ["000016.SH"]
|
||||
assert manager._stock_name_cache == {
|
||||
"sh000016": "静态上证50",
|
||||
"000016": "深康佳A",
|
||||
}
|
||||
|
||||
|
||||
@pytest.mark.parametrize(
|
||||
("stock_code", "aliases"),
|
||||
[
|
||||
(
|
||||
"sh000016",
|
||||
[
|
||||
"sh000016",
|
||||
"000016",
|
||||
"000016.SH",
|
||||
"000016.SS",
|
||||
"000016.SZ",
|
||||
"SSE50",
|
||||
],
|
||||
),
|
||||
("sz399001", ["sz399001", "399001", "399001.SZ"]),
|
||||
],
|
||||
)
|
||||
def test_index_name_validation_rejects_all_code_aliases(stock_code, aliases) -> None:
|
||||
target = parse_analysis_target(stock_code)
|
||||
|
||||
assert all(
|
||||
not DataFetcherManager._is_meaningful_cn_index_name(alias, target)
|
||||
for alias in aliases
|
||||
)
|
||||
assert DataFetcherManager._is_meaningful_cn_index_name("有效指数名称", target)
|
||||
|
||||
|
||||
def test_index_name_rejects_code_aliases_from_registry_cache_network_and_static() -> None:
|
||||
entry = IndexEntry(
|
||||
bare_code="000016",
|
||||
exchange="SH",
|
||||
canonical_id="sh000016",
|
||||
display_name="000016.SH",
|
||||
)
|
||||
target = AnalysisTarget(
|
||||
raw_input="sh000016",
|
||||
asset_type=ParseStatus.INDEX,
|
||||
canonical_id="sh000016",
|
||||
display_code="SSE50",
|
||||
exchange="SH",
|
||||
matched_index=entry,
|
||||
)
|
||||
tencent = _FakeFetcher("TencentFetcher", priority=0, name_result="sh000016")
|
||||
akshare = _FakeFetcher("AkshareFetcher", priority=1, name_result="000016.SS")
|
||||
tickflow = _FakeFetcher("TickFlowFetcher", priority=2, name_result="SSE50")
|
||||
manager = DataFetcherManager(fetchers=[tencent, akshare, tickflow])
|
||||
manager._stock_name_cache["sh000016"] = "000016"
|
||||
|
||||
with patch.dict(
|
||||
"data_provider.base.STOCK_NAME_MAP", {"sh000016": "000016"}, clear=True
|
||||
), patch("data_provider.base.parse_analysis_target", return_value=target):
|
||||
name = manager.get_stock_name("sh000016", allow_realtime=False)
|
||||
|
||||
assert name == "sh000016"
|
||||
assert manager._stock_name_cache == {}
|
||||
assert tencent.name_calls == ["sh000016"]
|
||||
assert akshare.name_calls == ["sh000016"]
|
||||
assert tickflow.name_calls == ["000016.SH"]
|
||||
|
||||
|
||||
def test_index_name_returns_canonical_id_when_registry_and_network_fail() -> None:
|
||||
entry = IndexEntry(
|
||||
bare_code="399001",
|
||||
exchange="SZ",
|
||||
canonical_id="sz399001",
|
||||
display_name="",
|
||||
)
|
||||
target = AnalysisTarget(
|
||||
raw_input="sz399001",
|
||||
asset_type=ParseStatus.INDEX,
|
||||
canonical_id="sz399001",
|
||||
display_code="",
|
||||
exchange="SZ",
|
||||
matched_index=entry,
|
||||
)
|
||||
fetchers = [
|
||||
_FakeFetcher("TencentFetcher", priority=0, name_result=RuntimeError("failed")),
|
||||
_FakeFetcher("AkshareFetcher", priority=1, name_result=None),
|
||||
_FakeFetcher("TickFlowFetcher", priority=2, name_result=""),
|
||||
]
|
||||
manager = DataFetcherManager(fetchers=cast(list[BaseFetcher], fetchers))
|
||||
|
||||
with patch("data_provider.base.parse_analysis_target", return_value=target):
|
||||
name = manager.get_stock_name("sz399001", allow_realtime=False)
|
||||
|
||||
assert name == "sz399001"
|
||||
assert manager._stock_name_cache == {}
|
||||
|
||||
|
||||
def test_index_and_same_bare_stock_names_use_isolated_cache_keys() -> None:
|
||||
manager = _manager_without_fetchers()
|
||||
|
||||
with patch.dict(
|
||||
"data_provider.base.STOCK_NAME_MAP", {"000016": "深康佳A"}, clear=True
|
||||
):
|
||||
stock_name = manager.get_stock_name("000016", allow_realtime=False)
|
||||
index_name = manager.get_stock_name("sh000016", allow_realtime=False)
|
||||
|
||||
assert stock_name == "深康佳A"
|
||||
assert index_name == "上证50"
|
||||
assert manager._stock_name_cache == {
|
||||
"000016": "深康佳A",
|
||||
"sh000016": "上证50",
|
||||
}
|
||||
@@ -0,0 +1,308 @@
|
||||
import sys
|
||||
import types
|
||||
from types import SimpleNamespace
|
||||
from unittest.mock import MagicMock, patch
|
||||
|
||||
import pandas as pd
|
||||
import pytest
|
||||
|
||||
from data_provider.akshare_fetcher import AkshareFetcher
|
||||
from data_provider.base import DataFetchError, RateLimitError, STANDARD_COLUMNS
|
||||
|
||||
|
||||
def _make_fetcher() -> AkshareFetcher:
|
||||
with patch(
|
||||
"data_provider.akshare_fetcher.get_config",
|
||||
return_value=SimpleNamespace(enable_eastmoney_patch=False),
|
||||
):
|
||||
return AkshareFetcher(sleep_min=0, sleep_max=0)
|
||||
|
||||
|
||||
def _call_akshare_inline(func, *args, timeout=None, call_name="", **kwargs):
|
||||
return func(*args, **kwargs)
|
||||
|
||||
|
||||
def test_akshare_recognized_index_uses_index_daily_api_and_standard_columns() -> None:
|
||||
raw = pd.DataFrame(
|
||||
{
|
||||
"date": ["2026-08-20", "2026-08-21"],
|
||||
"open": [100.0, 101.0],
|
||||
"close": [101.0, 102.0],
|
||||
"high": [102.0, 103.0],
|
||||
"low": [99.0, 100.0],
|
||||
"volume": [1000, 1100],
|
||||
"amount": [101000.0, 112200.0],
|
||||
}
|
||||
)
|
||||
index_daily = MagicMock(return_value=raw)
|
||||
fake_akshare = types.SimpleNamespace(stock_zh_index_daily_em=index_daily)
|
||||
fetcher = _make_fetcher()
|
||||
|
||||
with patch.dict(sys.modules, {"akshare": fake_akshare}), patch.object(
|
||||
fetcher, "_set_random_user_agent"
|
||||
), patch.object(fetcher, "_enforce_rate_limit"), patch(
|
||||
"data_provider.akshare_fetcher._akshare_call_with_timeout",
|
||||
side_effect=_call_akshare_inline,
|
||||
) as timeout_call:
|
||||
df = fetcher.get_daily_data(
|
||||
"sh000016", start_date="2026-08-01", end_date="2026-08-21"
|
||||
)
|
||||
|
||||
index_daily.assert_called_once_with(
|
||||
symbol="sh000016", start_date="20260801", end_date="20260821"
|
||||
)
|
||||
timeout_call.assert_called_once_with(
|
||||
index_daily,
|
||||
timeout=fetcher._history_call_timeout,
|
||||
call_name="ak.stock_zh_index_daily_em",
|
||||
symbol="sh000016",
|
||||
start_date="20260801",
|
||||
end_date="20260821",
|
||||
)
|
||||
assert set(STANDARD_COLUMNS).issubset(df.columns)
|
||||
assert set(["code", "ma5", "ma10", "ma20"]).issubset(df.columns)
|
||||
assert df["code"].tolist() == ["sh000016", "sh000016"]
|
||||
assert df["pct_chg"].round(2).tolist() == [0.0, 0.99]
|
||||
|
||||
|
||||
def test_akshare_index_daily_sorts_and_handles_missing_amount_and_infinite_pct() -> None:
|
||||
raw = pd.DataFrame(
|
||||
{
|
||||
"date": ["2026-08-21", "2026-08-20"],
|
||||
"open": [10.0, 0.0],
|
||||
"close": [10.0, 0.0],
|
||||
"high": [11.0, 0.0],
|
||||
"low": [9.0, 0.0],
|
||||
"volume": [1000, 900],
|
||||
}
|
||||
)
|
||||
fake_akshare = types.SimpleNamespace(
|
||||
stock_zh_index_daily_em=MagicMock(return_value=raw)
|
||||
)
|
||||
fetcher = _make_fetcher()
|
||||
|
||||
with patch.dict(sys.modules, {"akshare": fake_akshare}), patch(
|
||||
"data_provider.akshare_fetcher._akshare_call_with_timeout",
|
||||
side_effect=_call_akshare_inline,
|
||||
), patch.object(fetcher, "_set_random_user_agent"), patch.object(
|
||||
fetcher, "_enforce_rate_limit"
|
||||
):
|
||||
result = fetcher._fetch_index_data(
|
||||
"sh000016", "2026-08-01", "2026-08-21"
|
||||
)
|
||||
|
||||
assert result["date"].tolist() == ["2026-08-20", "2026-08-21"]
|
||||
assert "amount" in result.columns
|
||||
assert result["amount"].isna().all()
|
||||
assert result["pct_chg"].iloc[0] == 0.0
|
||||
assert pd.isna(result["pct_chg"].iloc[1])
|
||||
|
||||
|
||||
def test_akshare_index_daily_rejects_invalid_required_schema() -> None:
|
||||
raw = pd.DataFrame(
|
||||
{
|
||||
"date": ["2026-08-21"],
|
||||
"open": [10.0],
|
||||
"close": [10.0],
|
||||
"high": [11.0],
|
||||
"low": [9.0],
|
||||
}
|
||||
)
|
||||
fake_akshare = types.SimpleNamespace(
|
||||
stock_zh_index_daily_em=MagicMock(return_value=raw)
|
||||
)
|
||||
fetcher = _make_fetcher()
|
||||
|
||||
with patch.dict(sys.modules, {"akshare": fake_akshare}), patch(
|
||||
"data_provider.akshare_fetcher._akshare_call_with_timeout",
|
||||
side_effect=_call_akshare_inline,
|
||||
), patch.object(fetcher, "_set_random_user_agent"), patch.object(
|
||||
fetcher, "_enforce_rate_limit"
|
||||
), pytest.raises(DataFetchError, match="volume"):
|
||||
fetcher._fetch_index_data("sh000016", "2026-08-01", "2026-08-21")
|
||||
|
||||
|
||||
def test_akshare_index_daily_rejects_unparseable_dates() -> None:
|
||||
raw = pd.DataFrame(
|
||||
{
|
||||
"date": ["2026-08-20", "not-a-date"],
|
||||
"open": [10.0, 10.1],
|
||||
"close": [10.1, 10.2],
|
||||
"high": [10.2, 10.3],
|
||||
"low": [9.9, 10.0],
|
||||
"volume": [1000, 1100],
|
||||
}
|
||||
)
|
||||
fake_akshare = types.SimpleNamespace(
|
||||
stock_zh_index_daily_em=MagicMock(return_value=raw)
|
||||
)
|
||||
fetcher = _make_fetcher()
|
||||
|
||||
with patch.dict(sys.modules, {"akshare": fake_akshare}), patch(
|
||||
"data_provider.akshare_fetcher._akshare_call_with_timeout",
|
||||
side_effect=_call_akshare_inline,
|
||||
), patch.object(fetcher, "_set_random_user_agent"), patch.object(
|
||||
fetcher, "_enforce_rate_limit"
|
||||
), pytest.raises(DataFetchError, match="date"):
|
||||
fetcher._fetch_index_data("sh000016", "2026-08-01", "2026-08-21")
|
||||
|
||||
|
||||
@pytest.mark.parametrize("error_type", [ConnectionError, TimeoutError])
|
||||
def test_akshare_index_daily_preserves_retryable_builtin_errors(error_type) -> None:
|
||||
fake_akshare = types.SimpleNamespace(stock_zh_index_daily_em=object())
|
||||
fetcher = _make_fetcher()
|
||||
|
||||
with patch.dict(sys.modules, {"akshare": fake_akshare}), patch(
|
||||
"data_provider.akshare_fetcher._akshare_call_with_timeout",
|
||||
side_effect=error_type("retry me"),
|
||||
), patch.object(fetcher, "_set_random_user_agent"), patch.object(
|
||||
fetcher, "_enforce_rate_limit"
|
||||
), pytest.raises(error_type, match="retry me"):
|
||||
fetcher._fetch_index_data("sh000016", "2026-08-01", "2026-08-21")
|
||||
|
||||
|
||||
@pytest.mark.parametrize(
|
||||
("error", "expected_type"),
|
||||
[
|
||||
(RuntimeError("rate limit reached"), RateLimitError),
|
||||
(ValueError("bad payload"), DataFetchError),
|
||||
],
|
||||
)
|
||||
def test_akshare_index_daily_classifies_non_retryable_errors(
|
||||
error, expected_type
|
||||
) -> None:
|
||||
fake_akshare = types.SimpleNamespace(stock_zh_index_daily_em=object())
|
||||
fetcher = _make_fetcher()
|
||||
|
||||
with patch.dict(sys.modules, {"akshare": fake_akshare}), patch(
|
||||
"data_provider.akshare_fetcher._akshare_call_with_timeout",
|
||||
side_effect=error,
|
||||
), patch.object(fetcher, "_set_random_user_agent"), patch.object(
|
||||
fetcher, "_enforce_rate_limit"
|
||||
), pytest.raises(expected_type):
|
||||
fetcher._fetch_index_data("sh000016", "2026-08-01", "2026-08-21")
|
||||
|
||||
|
||||
def test_akshare_generic_normalization_does_not_fill_index_only_columns() -> None:
|
||||
fetcher = _make_fetcher()
|
||||
raw = pd.DataFrame(
|
||||
{
|
||||
"日期": ["2026-08-21"],
|
||||
"开盘": [10.0],
|
||||
"收盘": [10.2],
|
||||
"最高": [10.5],
|
||||
"最低": [9.8],
|
||||
"成交量": [1000],
|
||||
}
|
||||
)
|
||||
|
||||
result = fetcher._normalize_data(raw, "600519")
|
||||
|
||||
assert "amount" not in result.columns
|
||||
assert "pct_chg" not in result.columns
|
||||
|
||||
|
||||
def test_akshare_bare_conflict_code_stays_on_stock_api() -> None:
|
||||
fetcher = _make_fetcher()
|
||||
stock_result = pd.DataFrame({"日期": ["2026-08-21"]})
|
||||
|
||||
with patch.object(
|
||||
fetcher, "_fetch_stock_data", return_value=stock_result
|
||||
) as fetch_stock, patch.object(
|
||||
fetcher, "_fetch_index_data", return_value=pd.DataFrame()
|
||||
) as fetch_index:
|
||||
result = fetcher._fetch_raw_data("000016", "2026-08-01", "2026-08-21")
|
||||
|
||||
assert result is stock_result
|
||||
fetch_stock.assert_called_once_with("000016", "2026-08-01", "2026-08-21")
|
||||
fetch_index.assert_not_called()
|
||||
|
||||
|
||||
def test_akshare_index_name_filters_exchange_spot_table() -> None:
|
||||
spot = MagicMock(
|
||||
return_value=pd.DataFrame(
|
||||
{
|
||||
"代码": ["000001", "000016"],
|
||||
"名称": ["上证指数", "上证50"],
|
||||
}
|
||||
)
|
||||
)
|
||||
fake_akshare = types.SimpleNamespace(stock_zh_index_spot_em=spot)
|
||||
fetcher = _make_fetcher()
|
||||
|
||||
with patch.dict(sys.modules, {"akshare": fake_akshare}), patch.object(
|
||||
fetcher, "_set_random_user_agent"
|
||||
), patch.object(fetcher, "_enforce_rate_limit"), patch(
|
||||
"data_provider.akshare_fetcher._akshare_call_with_timeout",
|
||||
side_effect=_call_akshare_inline,
|
||||
) as timeout_call:
|
||||
name = fetcher.get_stock_name("sh000016")
|
||||
|
||||
assert name == "上证50"
|
||||
spot.assert_called_once_with(symbol="上证系列指数")
|
||||
timeout_call.assert_called_once_with(
|
||||
spot,
|
||||
timeout=fetcher._history_call_timeout,
|
||||
call_name="ak.stock_zh_index_spot_em",
|
||||
symbol="上证系列指数",
|
||||
)
|
||||
|
||||
|
||||
def test_akshare_index_name_uses_shenzhen_spot_table() -> None:
|
||||
spot = MagicMock(
|
||||
return_value=pd.DataFrame({"代码": [399001], "名称": ["深证成指"]})
|
||||
)
|
||||
fake_akshare = types.SimpleNamespace(stock_zh_index_spot_em=spot)
|
||||
fetcher = _make_fetcher()
|
||||
|
||||
with patch.dict(sys.modules, {"akshare": fake_akshare}), patch.object(
|
||||
fetcher, "_set_random_user_agent"
|
||||
), patch.object(fetcher, "_enforce_rate_limit"), patch(
|
||||
"data_provider.akshare_fetcher._akshare_call_with_timeout",
|
||||
side_effect=_call_akshare_inline,
|
||||
):
|
||||
name = fetcher.get_stock_name("sz399001")
|
||||
|
||||
assert name == "深证成指"
|
||||
spot.assert_called_once_with(symbol="深证系列指数")
|
||||
|
||||
|
||||
def test_akshare_index_name_rejects_missing_name_value() -> None:
|
||||
spot = MagicMock(
|
||||
return_value=pd.DataFrame({"代码": ["000016"], "名称": [pd.NA]})
|
||||
)
|
||||
fake_akshare = types.SimpleNamespace(stock_zh_index_spot_em=spot)
|
||||
fetcher = _make_fetcher()
|
||||
|
||||
with patch.dict(sys.modules, {"akshare": fake_akshare}), patch.object(
|
||||
fetcher, "_set_random_user_agent"
|
||||
), patch.object(fetcher, "_enforce_rate_limit"), patch(
|
||||
"data_provider.akshare_fetcher._akshare_call_with_timeout",
|
||||
side_effect=_call_akshare_inline,
|
||||
):
|
||||
name = fetcher.get_stock_name("sh000016")
|
||||
|
||||
assert name is None
|
||||
|
||||
|
||||
def test_akshare_index_name_uses_first_non_blank_duplicate() -> None:
|
||||
spot = MagicMock(
|
||||
return_value=pd.DataFrame(
|
||||
{
|
||||
"代码": ["000016", "000016", "000016"],
|
||||
"名称": [pd.NA, " ", "上证50"],
|
||||
}
|
||||
)
|
||||
)
|
||||
fake_akshare = types.SimpleNamespace(stock_zh_index_spot_em=spot)
|
||||
fetcher = _make_fetcher()
|
||||
|
||||
with patch.dict(sys.modules, {"akshare": fake_akshare}), patch.object(
|
||||
fetcher, "_set_random_user_agent"
|
||||
), patch.object(fetcher, "_enforce_rate_limit"), patch(
|
||||
"data_provider.akshare_fetcher._akshare_call_with_timeout",
|
||||
side_effect=_call_akshare_inline,
|
||||
):
|
||||
name = fetcher.get_stock_name("sh000016")
|
||||
|
||||
assert name == "上证50"
|
||||
@@ -70,12 +70,15 @@ class TestPrefetchStockNames(unittest.TestCase):
|
||||
manager = DataFetcherManager.__new__(DataFetcherManager)
|
||||
manager.get_stock_name = MagicMock(return_value="")
|
||||
|
||||
DataFetcherManager.prefetch_stock_names(manager, ["SH600519", "000001"], use_bulk=False)
|
||||
DataFetcherManager.prefetch_stock_names(
|
||||
manager, ["SH600519", "000001", "sh000016"], use_bulk=False
|
||||
)
|
||||
|
||||
manager.get_stock_name.assert_has_calls(
|
||||
[
|
||||
call("600519", allow_realtime=False),
|
||||
call("000001", allow_realtime=False),
|
||||
call("sh000016", allow_realtime=False),
|
||||
]
|
||||
)
|
||||
|
||||
|
||||
@@ -7,9 +7,10 @@ import os
|
||||
import subprocess
|
||||
import sys
|
||||
from pathlib import Path
|
||||
from unittest.mock import patch
|
||||
from unittest.mock import MagicMock, patch
|
||||
|
||||
import pandas as pd
|
||||
from requests import Response
|
||||
|
||||
from data_provider.tencent_fetcher import TencentFetcher, _to_tencent_symbol
|
||||
|
||||
@@ -107,6 +108,66 @@ def test_tencent_symbol_conversion_supports_a_share_markets() -> None:
|
||||
assert _to_tencent_symbol("600519") == "sh600519"
|
||||
assert _to_tencent_symbol("000001") == "sz000001"
|
||||
assert _to_tencent_symbol("920748") == "bj920748"
|
||||
assert _to_tencent_symbol("sh000016") == "sh000016"
|
||||
assert _to_tencent_symbol("000016.SH") == "sh000016"
|
||||
assert _to_tencent_symbol("sz399001") == "sz399001"
|
||||
|
||||
|
||||
def test_tencent_fetcher_preserves_explicit_index_market_for_daily_request() -> None:
|
||||
payload = {
|
||||
"data": {
|
||||
"sh000016": {
|
||||
"qfqday": [
|
||||
["2026-08-21", "100", "101", "102", "99", "1000", "101000"]
|
||||
]
|
||||
}
|
||||
}
|
||||
}
|
||||
response = MagicMock()
|
||||
response.json.return_value = payload
|
||||
|
||||
with patch("data_provider.tencent_fetcher.requests.get", return_value=response) as request:
|
||||
df = TencentFetcher().get_daily_data(
|
||||
"sh000016", start_date="2026-08-01", end_date="2026-08-21"
|
||||
)
|
||||
|
||||
assert not df.empty
|
||||
assert request.call_args.kwargs["params"]["param"].startswith("sh000016,day,")
|
||||
|
||||
|
||||
def test_tencent_fetcher_get_stock_name_uses_lightweight_quote_request() -> None:
|
||||
response = MagicMock()
|
||||
response.text = 'v_sh000016="1~上证50~000016~0";'
|
||||
|
||||
with patch("data_provider.tencent_fetcher.requests.get", return_value=response) as request:
|
||||
name = TencentFetcher().get_stock_name("sh000016")
|
||||
|
||||
assert name == "上证50"
|
||||
assert request.call_args.args[0] == "https://qt.gtimg.cn/q=sh000016"
|
||||
response.raise_for_status.assert_called_once_with()
|
||||
|
||||
|
||||
def test_tencent_fetcher_get_stock_name_rejects_mismatched_response_code() -> None:
|
||||
response = MagicMock()
|
||||
response.text = 'v_sh000016="1~沪深300~000300~0";'
|
||||
|
||||
with patch("data_provider.tencent_fetcher.requests.get", return_value=response):
|
||||
name = TencentFetcher().get_stock_name("sh000016")
|
||||
|
||||
assert name is None
|
||||
|
||||
|
||||
def test_tencent_fetcher_get_stock_name_decodes_real_gbk_response() -> None:
|
||||
response = Response()
|
||||
response.status_code = 200
|
||||
response.url = "https://qt.gtimg.cn/q=sh000016"
|
||||
response.encoding = "ISO-8859-1"
|
||||
response._content = 'v_sh000016="1~上证50~000016~0";'.encode("gbk")
|
||||
|
||||
with patch("data_provider.tencent_fetcher.requests.get", return_value=response):
|
||||
name = TencentFetcher().get_stock_name("sh000016")
|
||||
|
||||
assert name == "上证50"
|
||||
|
||||
|
||||
def test_tencent_fetcher_parses_qfq_daily_response() -> None:
|
||||
|
||||
Reference in New Issue
Block a user