跳转至

data_master — 主接口

StockDataMaster 是整个系统对外暴露的单例入口,封装多数据源、缓存、健康检测的统一访问。

data_master

StockDataMaster主数据接口

提供统一的股票数据访问接口,集成多数据源、缓存、健康检测等功能

StockDataMaster

StockDataMaster(config_path: Optional[str] = None)

股票数据主数据接口

初始化StockDataMaster

参数:

名称 类型 描述 默认
config_path Optional[str]

配置文件路径

None
源代码位于: data_master.py
def __init__(self, config_path: Optional[str] = None):
    """
    初始化StockDataMaster

    Args:
        config_path: 配置文件路径
    """
    # 避免重复初始化
    if hasattr(self, '_initialized'):
        return

    # 加载配置
    self.config = get_config(config_path)

    # 设置日志
    self._setup_logging()
    self.logger = logging.getLogger("StockDataMaster")
    self.logger.info("初始化StockDataMaster...")

    # 创建适配器
    self.adapters = {}
    self._init_adapters()

    # 创建缓存管理器
    self.cache_manager = CacheManager(self.config, self.adapters)

    # 创建健康管理器
    self.health_manager = HealthManager(self.config, self.adapters)

    # 🔥 优化: 股票名称缓存 (减少重复查询)
    self._stock_name_cache = {}  # {code: name}

    # 🔥 优化: baostock会话状态
    self._bs_session_active = False
    self._bs_last_login_time = None

    # 🔥 baostock 冷却机制
    self._baostock_consecutive_failures = 0
    self._baostock_cooldown_until = 0  # time.time() 时间戳

    # 启动健康监控
    if self.config.is_health_check_enabled():
        self.health_manager.start_monitoring()

    self._initialized = True
    self.logger.info("StockDataMaster初始化完成")

get_kline

get_kline(code: str, freq: str = 'd', start_date: Optional[str] = None, end_date: Optional[str] = None, count: Optional[int] = None, adjust: str = 'qfq', use_cache: bool = True) -> Optional[DataFrame]

获取K线数据(统一接口)

参数:

名称 类型 描述 默认
code str

股票代码,如'600519'或'sh.600519'

必需
freq str

频率 'd'=日,'w'=周,'m'=月,'5m'=5分钟,'15m'=15分钟,'30m'=30分钟,'60m'=60分钟

'd'
start_date Optional[str]

开始日期 'YYYY-MM-DD'

None
end_date Optional[str]

结束日期 'YYYY-MM-DD'

None
count Optional[int]

获取数量(从最新往前)

None
adjust str

复权类型,'qfq'=前复权(统一使用)

'qfq'
use_cache bool

是否使用缓存(仅对日线有效)

True

返回:

类型 描述
Optional[DataFrame]

DataFrame,包含date,open,high,low,close,volume,amount列

源代码位于: data_master.py
def get_kline(
    self,
    code: str,
    freq: str = 'd',
    start_date: Optional[str] = None,
    end_date: Optional[str] = None,
    count: Optional[int] = None,
    adjust: str = 'qfq',
    use_cache: bool = True
) -> Optional[pd.DataFrame]:
    """
    获取K线数据(统一接口)

    Args:
        code: 股票代码,如'600519'或'sh.600519'
        freq: 频率 'd'=日,'w'=周,'m'=月,'5m'=5分钟,'15m'=15分钟,'30m'=30分钟,'60m'=60分钟
        start_date: 开始日期 'YYYY-MM-DD'
        end_date: 结束日期 'YYYY-MM-DD'
        count: 获取数量(从最新往前)
        adjust: 复权类型,'qfq'=前复权(统一使用)
        use_cache: 是否使用缓存(仅对日线有效)

    Returns:
        DataFrame,包含date,open,high,low,close,volume,amount列
    """
    # 边界条件检查: 空股票代码
    if not code or not code.strip():
        self.logger.warning("股票代码不能为空")
        return None

    # 边界条件检查: count为0或负数
    if count is not None and count <= 0:
        self.logger.warning(f"count参数无效: {count}, 返回空DataFrame")
        return pd.DataFrame(columns=['date', 'open', 'high', 'low', 'close', 'volume', 'amount'])

    # 标准化代码
    code = self._normalize_code(code)

    # 日线数据优先使用缓存
    if freq == 'd' and use_cache and self.cache_manager.enabled:
        cached_df = self.cache_manager.get_cached_kline(code, start_date, end_date, count)

        if cached_df is not None and not cached_df.empty:
            # 检查缓存是否足够新(传入用户请求的 end_date 参数)
            is_fresh = self._is_cache_fresh(cached_df, end_date)
            # 检查数据量是否充足 - 放宽条件,>=95%即可
            # 原因: 双源校验可能导致少数数据校验失败,如119/120
            is_sufficient = count is None or len(cached_df) >= count * 0.95
            # 检查请求的日期范围是否超出缓存范围
            is_date_range_covered = self._is_date_range_covered(cached_df, start_date, end_date)

            self.logger.debug(f"缓存检查: code={code}, is_fresh={is_fresh}, is_sufficient={is_sufficient}, is_date_range_covered={is_date_range_covered}, cached_rows={len(cached_df)}, request_count={count}")

            if is_fresh and is_sufficient and is_date_range_covered:
                # 确保source字段存在(防止attrs丢失)
                if 'source' not in cached_df.attrs:
                    cached_df.attrs['source'] = 'cache'
                self.logger.debug(f"从缓存获取{code}数据: {len(cached_df)}条, source={cached_df.attrs.get('source')}")
                return cached_df
            elif is_fresh and not is_sufficient:
                self.logger.debug(f"{code}缓存数据不足({len(cached_df)}/{count}条),从数据源补充")
            elif not is_date_range_covered:
                self.logger.debug(f"{code}缓存日期范围不覆盖请求范围,重新获取")
            else:
                self.logger.debug(f"{code}缓存数据不新鲜(latest={cached_df['date'].iloc[-1]}),重新获取")

    # 从数据源获取
    df = self._fetch_kline_from_source(code, freq, start_date, end_date, count, adjust)

    # 日线数据尝试缓存(双源校验)
    # 注意:无论 use_cache 是否为 True,已获取的数据都应尝试缓存。
    # use_cache 仅控制上面的"是否从缓存读取",不应影响写入。
    # 否则 warmup 路径传入 use_cache=False 时,拉取的新数据不会写入缓存。
    if df is not None and freq == 'd':
        self._try_cache_kline(code, df, start_date, end_date, count, adjust)

    return df

get_valuation

get_valuation(code: str, start_date: Optional[str] = None, end_date: Optional[str] = None) -> Optional[DataFrame]

获取估值数据

参数:

名称 类型 描述 默认
code str

股票代码

必需
start_date Optional[str]

开始日期

None
end_date Optional[str]

结束日期

None

返回:

类型 描述
Optional[DataFrame]

DataFrame,包含pe_ttm,pb,ps_ttm等估值指标

源代码位于: data_master.py
def get_valuation(
    self,
    code: str,
    start_date: Optional[str] = None,
    end_date: Optional[str] = None
) -> Optional[pd.DataFrame]:
    """
    获取估值数据

    Args:
        code: 股票代码
        start_date: 开始日期
        end_date: 结束日期

    Returns:
        DataFrame,包含pe_ttm,pb,ps_ttm等估值指标
    """
    code = self._normalize_code(code)

    # 获取活跃数据源
    active_source = self.health_manager.get_active_source('valuation')

    if not active_source:
        self.logger.error("没有可用的估值数据源")
        return None

    # 获取备用列表
    sources = self.config.get_sources_by_usage('valuation')

    if active_source in sources:
        sources.remove(active_source)
        sources.insert(0, active_source)

    # 依次尝试
    for source_name in sources:
        adapter = self.adapters.get(source_name)

        if not adapter:
            continue

        try:
            df = adapter.get_valuation(code, start_date, end_date)

            if df is not None and not df.empty:
                self.logger.debug(f"成功从{source_name}获取{code}估值数据")
                return df

        except Exception as e:
            self.logger.error(f"{source_name}获取估值失败: {e}")
            continue

    self.logger.error(f"所有数据源均无法获取{code}的估值数据")
    return None

get_tick

get_tick(code: str) -> Optional[Dict[str, Any]]

获取实时tick数据

参数:

名称 类型 描述 默认
code str

股票代码

必需

返回:

类型 描述
Optional[Dict[str, Any]]

实时行情字典

源代码位于: data_master.py
def get_tick(self, code: str) -> Optional[Dict[str, Any]]:
    """
    获取实时tick数据

    Args:
        code: 股票代码

    Returns:
        实时行情字典
    """
    code = self._normalize_code(code)

    # 获取活跃数据源
    active_source = self.health_manager.get_active_source('tick')

    if not active_source:
        self.logger.error("没有可用的tick数据源")
        return None

    # 获取备用列表
    sources = self.config.get_sources_by_usage('tick')

    if active_source in sources:
        sources.remove(active_source)
        sources.insert(0, active_source)

    # 依次尝试
    for source_name in sources:
        adapter = self.adapters.get(source_name)

        if not adapter:
            continue

        try:
            tick = adapter.get_tick(code)

            if tick is not None:
                # 添加数据源信息到tick数据中
                tick['source'] = source_name
                self.logger.debug(f"成功从{source_name}获取{code}实时行情")
                return tick

        except Exception as e:
            self.logger.error(f"{source_name}获取tick失败: {e}")
            continue

    self.logger.error(f"所有数据源均无法获取{code}的tick数据")
    return None

get_health_status

get_health_status() -> Dict[str, Any]

获取系统健康状态

返回:

类型 描述
Dict[str, Any]

健康状态报告

源代码位于: data_master.py
def get_health_status(self) -> Dict[str, Any]:
    """
    获取系统健康状态

    Returns:
        健康状态报告
    """
    return self.health_manager.get_health_report()

get_cache_statistics

get_cache_statistics() -> Dict[str, Any]

获取缓存统计信息

返回:

类型 描述
Dict[str, Any]

缓存统计信息

源代码位于: data_master.py
def get_cache_statistics(self) -> Dict[str, Any]:
    """
    获取缓存统计信息

    Returns:
        缓存统计信息
    """
    return self.cache_manager.get_cache_statistics()

cleanup_cache

cleanup_cache(days: Optional[int] = None)

清理缓存

参数:

名称 类型 描述 默认
days Optional[int]

保留天数

None
源代码位于: data_master.py
def cleanup_cache(self, days: Optional[int] = None):
    """
    清理缓存

    Args:
        days: 保留天数
    """
    self.cache_manager.cleanup_old_cache(days)

force_switch_source

force_switch_source(usage: str, target_source: str) -> bool

强制切换数据源

参数:

名称 类型 描述 默认
usage str

用途 'kline'/'valuation'/'tick'

必需
target_source str

目标数据源名称

必需

返回:

类型 描述
bool

是否成功

源代码位于: data_master.py
def force_switch_source(self, usage: str, target_source: str) -> bool:
    """
    强制切换数据源

    Args:
        usage: 用途 'kline'/'valuation'/'tick'
        target_source: 目标数据源名称

    Returns:
        是否成功
    """
    return self.health_manager.force_switch(usage, target_source)

get_stock_name

get_stock_name(code: str) -> str

获取股票名称 - 四级查找链

L1: 内存缓存 dict L2: baostock query_stock_basic() (免费兜底,数据最全含退市股) L3: xtquant.get_stock_name() (QMT用户快速查询) L4: tushare pro.stock_basic() (付费用户补充)

参数:

名称 类型 描述 默认
code str

股票代码,支持 '600519', 'sh.600519', '600519.SH' 格式

必需

返回:

类型 描述
str

股票名称字符串,所有级别都失败时返回代码本身

源代码位于: data_master.py
def get_stock_name(self, code: str) -> str:
    """
    获取股票名称 - 四级查找链

    L1: 内存缓存 dict
    L2: baostock query_stock_basic() (免费兜底,数据最全含退市股)
    L3: xtquant.get_stock_name() (QMT用户快速查询)
    L4: tushare pro.stock_basic() (付费用户补充)

    Args:
        code: 股票代码,支持 '600519', 'sh.600519', '600519.SH' 格式

    Returns:
        股票名称字符串,所有级别都失败时返回代码本身
    """
    if not code:
        return code

    # 标准化为6位纯数字
    pure_code = self._normalize_code(code)

    # L1: 内存缓存
    if pure_code in self._stock_name_cache:
        self.logger.debug(f"L1缓存命中: {pure_code} -> {self._stock_name_cache[pure_code]}")
        return self._stock_name_cache[pure_code]

    # L2: baostock (免费无门槛,数据最全含退市股)
    name = self._get_stock_name_from_baostock(pure_code)
    if name:
        self._stock_name_cache[pure_code] = name
        self.cache_manager.cache_stock_name(pure_code, name, 'baostock')
        self.logger.info(f"L2 baostock获取股票名称: {pure_code} -> {name}")
        return name

    # L3: xtquant (QMT用户快速查询)
    name = self._get_stock_name_from_xtquant(pure_code)
    if name:
        self._stock_name_cache[pure_code] = name
        self.cache_manager.cache_stock_name(pure_code, name, 'xtquant')
        self.logger.info(f"L3 xtquant获取股票名称: {pure_code} -> {name}")
        return name

    # L4: tushare (付费用户补充)
    name = self._get_stock_name_from_tushare(pure_code)
    if name:
        self._stock_name_cache[pure_code] = name
        self.cache_manager.cache_stock_name(pure_code, name, 'tushare')
        self.logger.info(f"L4 tushare获取股票名称: {pure_code} -> {name}")
        return name

    # 所有级别失败,返回代码本身
    self.logger.warning(f"所有级别查找失败,返回代码本身: {pure_code}")
    return pure_code

warmup_stock_names

warmup_stock_names() -> int

预热股票名称缓存 - 从 Tushare 批量获取全市场股票名称, 写入内存缓存和 SQLite 持久化缓存。

调用时机:初始化后、盘前准备阶段、或缓存为空时。

返回:

类型 描述
int

缓存的股票名称数量

源代码位于: data_master.py
def warmup_stock_names(self) -> int:
    """
    预热股票名称缓存 - 从 Tushare 批量获取全市场股票名称,
    写入内存缓存和 SQLite 持久化缓存。

    调用时机:初始化后、盘前准备阶段、或缓存为空时。

    Returns:
        缓存的股票名称数量
    """
    adapter = self.adapters.get('tushare')
    if not adapter or not adapter.is_connected:
        self.logger.warning("Tushare未连接,跳过股票名称预热")
        return 0

    try:
        names = adapter.get_all_stock_names()
        if not names:
            self.logger.warning("Tushare未返回股票名称数据")
            return 0

        # 写入内存缓存
        self._stock_name_cache.update(names)

        # 批量写入 SQLite
        count = self.cache_manager.bulk_cache_stock_names(names, source='tushare')

        self.logger.info(f"股票名称预热完成: {count} 只股票已缓存 (内存+SQLite)")
        return count

    except Exception as e:
        self.logger.error(f"股票名称预热失败: {e}")
        return 0

close

close()

关闭DataMaster,清理资源

源代码位于: data_master.py
def close(self):
    """关闭DataMaster,清理资源"""
    self.logger.info("关闭StockDataMaster...")

    # 🔥 优化: 关闭baostock会话
    if self._bs_session_active:
        try:
            import baostock as bs
            bs.logout()
            self._bs_session_active = False
            self.logger.debug("baostock会话已关闭")
        except Exception as e:
            self.logger.warning(f"关闭baostock会话失败: {e}")

    # 停止健康监控
    self.health_manager.stop_monitoring()

    # 断开所有适配器
    for adapter in self.adapters.values():
        try:
            adapter.disconnect()
        except:
            pass

    self.logger.info("StockDataMaster已关闭")

get_data_master

get_data_master(config_path: Optional[str] = None) -> StockDataMaster

获取全局StockDataMaster实例(单例)

参数:

名称 类型 描述 默认
config_path Optional[str]

配置文件路径

None

返回:

类型 描述
StockDataMaster

StockDataMaster实例

源代码位于: data_master.py
def get_data_master(config_path: Optional[str] = None) -> 'StockDataMaster':
    """
    获取全局StockDataMaster实例(单例)

    Args:
        config_path: 配置文件路径

    Returns:
        StockDataMaster实例
    """
    global _global_master
    if _global_master is None:
        _global_master = StockDataMaster(config_path)
    return _global_master