跳转至

adapters — 数据源适配器

所有数据源遵循 DataSourceAdapter 抽象基类,由 AdapterFactory 按配置实例化。

基类 base_adapter

base_adapter

数据源适配器基类

定义统一的数据源接口规范

DataSourceAdapter

DataSourceAdapter(name: str, config: Dict[str, Any])

Bases: ABC

数据源适配器抽象基类

初始化适配器

参数:

名称 类型 描述 默认
name str

数据源名称

必需
config Dict[str, Any]

数据源配置

必需
源代码位于: adapters/base_adapter.py
def __init__(self, name: str, config: Dict[str, Any]):
    """
    初始化适配器

    Args:
        name: 数据源名称
        config: 数据源配置
    """
    self.name = name
    self.config = config
    self.logger = logging.getLogger(f"DataMaster.{name}")
    self.is_connected = False
    self.last_error = None
    self.error_count = 0
connect abstractmethod
connect() -> bool

连接数据源

返回:

类型 描述
bool

是否连接成功

源代码位于: adapters/base_adapter.py
@abstractmethod
def connect(self) -> bool:
    """
    连接数据源

    Returns:
        是否连接成功
    """
    pass
disconnect abstractmethod
disconnect()

断开数据源连接

源代码位于: adapters/base_adapter.py
@abstractmethod
def disconnect(self):
    """断开数据源连接"""
    pass
get_kline abstractmethod
get_kline(code: str, freq: str = 'd', start_date: Optional[str] = None, end_date: Optional[str] = None, count: Optional[int] = None, adjust: str = 'qfq') -> 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'

返回:

类型 描述
Optional[DataFrame]

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

Optional[DataFrame]

返回None表示失败

源代码位于: adapters/base_adapter.py
@abstractmethod
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'
) -> 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'前复权

    Returns:
        DataFrame,列包含: date,open,high,low,close,volume,amount
        返回None表示失败
    """
    pass
get_valuation abstractmethod
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,列包含: date,pe_ttm,pb,ps_ttm,pcf_ncf等

Optional[DataFrame]

返回None表示失败

源代码位于: adapters/base_adapter.py
@abstractmethod
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,列包含: date,pe_ttm,pb,ps_ttm,pcf_ncf等
        返回None表示失败
    """
    pass
get_tick abstractmethod
get_tick(code: str) -> Optional[Dict[str, Any]]

获取实时tick数据

参数:

名称 类型 描述 默认
code str

股票代码

必需

返回:

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

字典,包含: open,high,low,close,volume,amount,last,bid,ask等

Optional[Dict[str, Any]]

返回None表示失败

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

    Args:
        code: 股票代码

    Returns:
        字典,包含: open,high,low,close,volume,amount,last,bid,ask等
        返回None表示失败
    """
    pass
health_check
health_check() -> Dict[str, Any]

健康检查

修复说明: - 在后台health_check线程中,适配器可能已断开连接 - 修复方案:每次health_check前强制重新连接 - 这样确保健康检查准确反映数据源的实际可用性

返回:

类型 描述
Dict[str, Any]

健康检查结果字典

Dict[str, Any]

{ 'status': 'ok'|'warning'|'error', 'response_time': float, # 响应时间(秒) 'data_freshness': bool, # 数据是否新鲜 'error_message': str|None

Dict[str, Any]

}

源代码位于: adapters/base_adapter.py
def health_check(self) -> Dict[str, Any]:
    """
    健康检查

    修复说明:
    - 在后台health_check线程中,适配器可能已断开连接
    - 修复方案:每次health_check前强制重新连接
    - 这样确保健康检查准确反映数据源的实际可用性

    Returns:
        健康检查结果字典
        {
            'status': 'ok'|'warning'|'error',
            'response_time': float,  # 响应时间(秒)
            'data_freshness': bool,  # 数据是否新鲜
            'error_message': str|None
        }
    """
    start_time = datetime.now()
    result = {
        'status': 'ok',
        'response_time': 0.0,
        'data_freshness': True,
        'error_message': None
    }

    try:
        # 重要修复:每次health_check前强制重新连接
        # 原因:后台线程调用时,适配器可能已断开连接
        if not self.is_connected:
            self.logger.debug(f"{self.name} health_check: 适配器未连接,尝试重新连接...")
            if not self.connect():
                result['status'] = 'error'
                result['error_message'] = f'连接失败: {self.last_error or "未知错误"}'
                result['response_time'] = (datetime.now() - start_time).total_seconds()
                return result

        # 尝试获取测试数据(浦发银行600000日K线最近5条 - 活跃股票,数据稳定)
        # 修改原因1:600519(茅台)在某些数据源可能不稳定,改用600000提高健康检查准确性
        # 修改原因2:count=1在某些数据源会被过滤,改用count=30确保能获取到数据
        test_df = self.get_kline('600000', freq='d', count=30)

        # 计算响应时间
        response_time = (datetime.now() - start_time).total_seconds()
        result['response_time'] = response_time

        # 检查数据有效性
        if test_df is None or test_df.empty:
            result['status'] = 'error'
            result['error_message'] = '无法获取测试数据'
            self.error_count += 1
            return result

        # 检查数据新鲜度(最新数据应该是最近3个交易日内的)
        latest_date = pd.to_datetime(test_df['date'].iloc[-1])
        days_diff = (datetime.now() - latest_date).days
        if days_diff > self.config.get('data_freshness_days', 3):
            result['status'] = 'warning'
            result['data_freshness'] = False
            result['error_message'] = f'数据不新鲜,最新数据日期: {latest_date.date()}'

        # 检查响应时间
        threshold = self.config.get('timeout', 5)
        if response_time > threshold:
            result['status'] = 'warning'
            result['error_message'] = f'响应时间过长: {response_time:.2f}秒'

        # 成功则重置错误计数
        self.error_count = 0

    except Exception as e:
        result['status'] = 'error'
        result['error_message'] = str(e)
        result['response_time'] = (datetime.now() - start_time).total_seconds()
        self.error_count += 1
        self.last_error = str(e)
        self.logger.error(f"健康检查失败: {e}")

    return result
normalize_code
normalize_code(code: str) -> str

标准化股票代码为无前缀格式

参数:

名称 类型 描述 默认
code str

股票代码,可能带有'sh.'或'sz.'前缀

必需

返回:

类型 描述
str

标准化后的6位代码

源代码位于: adapters/base_adapter.py
def normalize_code(self, code: str) -> str:
    """
    标准化股票代码为无前缀格式

    Args:
        code: 股票代码,可能带有'sh.'或'sz.'前缀

    Returns:
        标准化后的6位代码
    """
    if code.startswith(('sh.', 'sz.')):
        return code.split('.')[1]
    return code
add_prefix
add_prefix(code: str) -> str

为股票代码添加交易所前缀

参数:

名称 类型 描述 默认
code str

6位股票代码

必需

返回:

类型 描述
str

带前缀的代码,如 'sh.600519' 或 'sz.000001'

源代码位于: adapters/base_adapter.py
def add_prefix(self, code: str) -> str:
    """
    为股票代码添加交易所前缀

    Args:
        code: 6位股票代码

    Returns:
        带前缀的代码,如 'sh.600519' 或 'sz.000001'
    """
    code = self.normalize_code(code)
    if code.startswith(('6', '5')):
        return f'sh.{code}'
    elif code.startswith(('0', '3')):
        return f'sz.{code}'
    return code
standardize_dataframe
standardize_dataframe(df: DataFrame) -> DataFrame

标准化DataFrame格式

参数:

名称 类型 描述 默认
df DataFrame

原始DataFrame

必需

返回:

类型 描述
DataFrame

标准化后的DataFrame,包含统一列名

源代码位于: adapters/base_adapter.py
def standardize_dataframe(self, df: pd.DataFrame) -> pd.DataFrame:
    """
    标准化DataFrame格式

    Args:
        df: 原始DataFrame

    Returns:
        标准化后的DataFrame,包含统一列名
    """
    if df is None or df.empty:
        return df

    # 统一列名映射
    column_mapping = {
        'datetime': 'date',
        'vol': 'volume',
        'trade': 'amount'
    }

    # 重命名列
    df = df.rename(columns=column_mapping)

    # 确保date列为字符串格式 YYYY-MM-DD
    if 'date' in df.columns:
        df['date'] = pd.to_datetime(df['date']).dt.strftime('%Y-%m-%d')

    return df

工厂 AdapterFactory

adapters

适配器工厂

用于创建和管理所有数据源适配器实例

AdapterFactory

适配器工厂类

create_adapter classmethod
create_adapter(name: str, config: Dict) -> DataSourceAdapter

创建数据源适配器实例

参数:

名称 类型 描述 默认
name str

数据源名称

必需
config Dict

数据源配置

必需

返回:

类型 描述
DataSourceAdapter

DataSourceAdapter实例

引发:

类型 描述
ValueError

不支持的数据源类型

源代码位于: adapters/__init__.py
@classmethod
def create_adapter(cls, name: str, config: Dict) -> DataSourceAdapter:
    """
    创建数据源适配器实例

    Args:
        name: 数据源名称
        config: 数据源配置

    Returns:
        DataSourceAdapter实例

    Raises:
        ValueError: 不支持的数据源类型
    """
    adapter_class = cls.ADAPTER_MAP.get(name)

    if adapter_class is None:
        raise ValueError(f"不支持的数据源类型: {name}")

    return adapter_class(name, config)
get_supported_sources classmethod
get_supported_sources() -> list

获取支持的数据源列表

返回:

类型 描述
list

数据源名称列表

源代码位于: adapters/__init__.py
@classmethod
def get_supported_sources(cls) -> list:
    """
    获取支持的数据源列表

    Returns:
        数据源名称列表
    """
    return list(cls.ADAPTER_MAP.keys())

Tushare 适配器

tushare_adapter

Tushare数据源适配器

使用tushare库获取数据,主要用于估值数据和K线数据

TushareAdapter

TushareAdapter(name: str, config: Dict[str, Any])

Bases: DataSourceAdapter

Tushare数据源适配器

源代码位于: adapters/tushare_adapter.py
def __init__(self, name: str, config: Dict[str, Any]):
    super().__init__(name, config)
    self.pro = None
    self.token = config.get('token', '')
    self.use_adj_factor = config.get('use_adj_factor', False)
    # 滑动窗口限速:每分钟最多 max_calls_per_minute 次 API 调用
    self._max_cpm: int = config.get('max_calls_per_minute', 500)
    self._rate_times: deque = deque()
    self._rate_lock = threading.Lock()
connect
connect() -> bool

连接Tushare数据源

源代码位于: adapters/tushare_adapter.py
def connect(self) -> bool:
    """连接Tushare数据源"""
    try:
        if not self.token:
            self.logger.error(f"{self.name} token未配置")
            self.is_connected = False
            return False

        ts.set_token(self.token)
        self.pro = ts.pro_api()
        self.is_connected = True
        self.logger.info(f"{self.name} 连接成功")
        return True

    except Exception as e:
        self.logger.error(f"{self.name} 连接失败: {e}")
        self.last_error = str(e)
        self.is_connected = False
        return False
disconnect
disconnect()

断开连接

源代码位于: adapters/tushare_adapter.py
def disconnect(self):
    """断开连接"""
    self.pro = None
    self.is_connected = False
    self.logger.info(f"{self.name} 已断开连接")
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', fetch_turn: bool = True) -> Optional[DataFrame]

获取K线数据

参数:

名称 类型 描述 默认
code str

股票代码

必需
freq str

频率

'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'
fetch_turn bool

是否额外调用 daily_basic 获取换手率(默认True;预热时设为False节省API调用)

True

返回:

类型 描述
Optional[DataFrame]

DataFrame

源代码位于: adapters/tushare_adapter.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',
    fetch_turn: bool = True
) -> Optional[pd.DataFrame]:
    """
    获取K线数据

    Args:
        code: 股票代码
        freq: 频率
        start_date: 开始日期 YYYY-MM-DD
        end_date: 结束日期 YYYY-MM-DD
        count: 获取数量
        adjust: 复权类型 'qfq'=前复权
        fetch_turn: 是否额外调用 daily_basic 获取换手率(默认True;预热时设为False节省API调用)

    Returns:
        DataFrame
    """
    if not self.is_connected:
        if not self.connect():
            return None

    try:
        ts_code = self._convert_code(code)

        # 转换日期格式: YYYY-MM-DD -> YYYYMMDD
        if start_date:
            start_date = start_date.replace('-', '')
        if end_date:
            end_date = end_date.replace('-', '')

        # 如果提供count但没有start_date,计算start_date
        if count and not start_date:
            from datetime import datetime, timedelta
            # 估算天数(count * 2倍,充分考虑节假日和周末)
            days = int(count * 2)
            start_date = (datetime.now() - timedelta(days=days)).strftime('%Y%m%d')

        # 默认结束日期
        if not end_date:
            from datetime import datetime
            end_date = datetime.now().strftime('%Y%m%d')

        # ---- 数据获取 + 复权 ----
        #
        # use_adj_factor=false (默认, 推荐):
        #   ts.pro_bar(adj='qfq') 原生接口返回前复权数据
        #   优点: 数据准确, 一步到位
        #   要求: tushare 积分 >= 2000
        #
        # use_adj_factor=true (低积分备选):
        #   pro.daily() 取未复权数据 + pro.adj_factor() 取复权因子
        #   手动计算: qfq_price = price * adj_factor / latest_adj_factor
        #   优点: 仅需 120 积分
        #   缺点: 多一次 API 调用; 部分边缘数据可能有微小差异
        #
        if self.use_adj_factor:
            # --- 模式 B: daily + adj_factor 手动复权 (120 积分) ---
            df = self._call(self.pro.daily,
                           ts_code=ts_code,
                           start_date=start_date,
                           end_date=end_date)

            if df is None or df.empty:
                self.logger.debug(f"Tushare未获取到数据: {code}")
                return None

            # 先重命名 + 格式化, 再做手动复权
            df = df.rename(columns={
                'trade_date': 'date',
                'ts_code': 'code',
                'vol': 'volume',
                'pct_chg': 'pct_change'
            })
            if 'date' in df.columns:
                df['date'] = pd.to_datetime(df['date']).dt.strftime('%Y-%m-%d')
            df = df.sort_values('date').reset_index(drop=True)

            if adjust in ['qfq', 'hfq']:
                df = self._apply_adj_factor(df, ts_code, start_date, end_date, adjust)

        else:
            # --- 模式 A: pro_bar 原生复权 (2000 积分, 默认) ---
            # ts.pro_bar 直接返回已复权数据, adj 参数: qfq/hfq/None
            adj_param = adjust if adjust in ('qfq', 'hfq') else None
            df = self._call(ts.pro_bar,
                           ts_code=ts_code,
                           start_date=start_date,
                           end_date=end_date,
                           adj=adj_param,
                           freq=freq if freq in ['d', 'w', 'm'] else 'd')

            if df is None or df.empty:
                self.logger.debug(f"Tushare未获取到数据: {code}")
                return None

            # pro_bar 返回的列名: trade_date, open, high, low, close, vol, amount ...
            df = df.rename(columns={
                'trade_date': 'date',
                'ts_code': 'code',
                'vol': 'volume',
                'pct_chg': 'pct_change'
            })
            if 'date' in df.columns:
                df['date'] = pd.to_datetime(df['date']).dt.strftime('%Y-%m-%d')
            df = df.sort_values('date').reset_index(drop=True)

        # 如果指定了count,只取最近count条
        if count:
            df = df.tail(count)

        # 统一成交量单位为"股" (daily接口的vol单位是"手",需要*100)
        if 'volume' in df.columns:
            df['volume'] = df['volume'] * 100
            self.logger.debug(f"成交量单位转换: 手 -> 股 (*100)")

        # 统一成交额单位为"元" (daily接口的amount单位是"千元",需要*1000)
        if 'amount' in df.columns:
            df['amount'] = df['amount'] * 1000
            self.logger.debug(f"成交额单位转换: 千元 -> 元 (*1000)")

        # 日K线:可选获取换手率(fetch_turn=False 时跳过,节省 1 次 API 调用)
        # 预热场景传 fetch_turn=False;换手率可通过 get_daily_basic_batch 批量补充
        if freq == 'd' and fetch_turn:
            try:
                turn_df = self._call(self.pro.daily_basic,
                                    ts_code=ts_code,
                                    start_date=start_date,
                                    end_date=end_date,
                                    fields='trade_date,turnover_rate')
                if turn_df is not None and not turn_df.empty:
                    turn_df = turn_df.rename(columns={
                        'trade_date': 'date',
                        'turnover_rate': 'turn'
                    })
                    turn_df['date'] = pd.to_datetime(turn_df['date']).dt.strftime('%Y-%m-%d')
                    df = df.merge(turn_df[['date', 'turn']], on='date', how='left')
            except Exception as te:
                self.logger.debug(f"获取换手率失败 {code}: {te}")

        # 确保必要列
        required_cols = ['date', 'open', 'high', 'low', 'close', 'volume']
        if not all(col in df.columns for col in required_cols):
            self.logger.error(f"数据列不完整: {df.columns.tolist()}")
            return None

        # 标准化
        df = self.standardize_dataframe(df)

        return df

    except Exception as e:
        self.logger.error(f"Tushare获取K线失败 {code}: {e}")
        self.last_error = str(e)
        self.error_count += 1
        return None
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

源代码位于: adapters/tushare_adapter.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
    """
    if not self.is_connected:
        if not self.connect():
            return None

    try:
        ts_code = self._convert_code(code)

        # 转换日期格式
        if start_date:
            start_date = start_date.replace('-', '')
        if end_date:
            end_date = end_date.replace('-', '')

        # 默认日期范围
        if not end_date:
            from datetime import datetime
            end_date = datetime.now().strftime('%Y%m%d')
        if not start_date:
            from datetime import datetime, timedelta
            start_date = (datetime.now() - timedelta(days=365)).strftime('%Y%m%d')

        # 查询每日指标
        df = self._call(self.pro.daily_basic,
                        ts_code=ts_code,
                        start_date=start_date,
                        end_date=end_date,
                        fields='ts_code,trade_date,pe_ttm,pb,ps_ttm,total_mv,circ_mv')

        if df is None or df.empty:
            self.logger.warning(f"Tushare未获取到估值数据: {code}")
            return None

        # 列名映射
        df = df.rename(columns={
            'trade_date': 'date',
            'ts_code': 'code'
        })

        # 日期格式转换
        df['date'] = pd.to_datetime(df['date']).dt.strftime('%Y-%m-%d')

        # 排序
        df = df.sort_values('date').reset_index(drop=True)

        return df

    except Exception as e:
        self.logger.error(f"Tushare获取估值失败 {code}: {e}")
        self.last_error = str(e)
        self.error_count += 1
        return None
get_tick
get_tick(code: str) -> Optional[Dict[str, Any]]

获取实时tick数据

注意: Tushare实时行情需要高级权限,这里使用最新日线数据模拟

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

    注意: Tushare实时行情需要高级权限,这里使用最新日线数据模拟
    """
    try:
        # 获取最近1天数据
        df = self.get_kline(code, freq='d', count=1, adjust='qfq')

        if df is None or df.empty:
            return None

        latest = df.iloc[-1]

        tick_data = {
            'code': self.normalize_code(code),
            'open': float(latest['open']),
            'high': float(latest['high']),
            'low': float(latest['low']),
            'close': float(latest['close']),
            'last': float(latest['close']),
            'volume': float(latest['volume']),
            'amount': float(latest.get('amount', 0))
        }

        return tick_data

    except Exception as e:
        self.logger.error(f"Tushare模拟tick失败 {code}: {e}")
        self.last_error = str(e)
        return None
get_daily_batch
get_daily_batch(trade_date: str) -> Optional[DataFrame]

按交易日拉取全市场未复权日线行情(daily 接口,约 0.25s/次,覆盖全部 A 股)。

返回列: date, ts_code, open, high, low, close, volume, amount 调用方应自行做前复权(与 get_adj_factor_batch 配合)。

源代码位于: adapters/tushare_adapter.py
def get_daily_batch(self, trade_date: str) -> Optional[pd.DataFrame]:
    """
    按交易日拉取全市场未复权日线行情(daily 接口,约 0.25s/次,覆盖全部 A 股)。

    返回列: date, ts_code, open, high, low, close, volume, amount
    调用方应自行做前复权(与 get_adj_factor_batch 配合)。
    """
    if not self.is_connected:
        if not self.connect():
            return None
    try:
        df = self._call(self.pro.daily, trade_date=trade_date)
        if df is None or df.empty:
            return None
        df = df.rename(columns={'trade_date': 'date', 'vol': 'volume'})
        df['date'] = pd.to_datetime(df['date']).dt.strftime('%Y-%m-%d')
        df['volume'] = df['volume'] * 100    # 手 -> 股
        df['amount'] = df['amount'] * 1000   # 千元 -> 元
        return df.sort_values(['ts_code', 'date']).reset_index(drop=True)
    except Exception as e:
        self.logger.error(f"get_daily_batch {trade_date}: {e}")
        self.last_error = str(e)
        return None
get_adj_factor_batch
get_adj_factor_batch(trade_date: str) -> Optional[DataFrame]

按交易日拉取全市场复权因子(adj_factor 接口,约 0.14s/次)。

返回列: ts_code, date, adj_factor

源代码位于: adapters/tushare_adapter.py
def get_adj_factor_batch(self, trade_date: str) -> Optional[pd.DataFrame]:
    """
    按交易日拉取全市场复权因子(adj_factor 接口,约 0.14s/次)。

    返回列: ts_code, date, adj_factor
    """
    if not self.is_connected:
        if not self.connect():
            return None
    try:
        df = self._call(self.pro.adj_factor, trade_date=trade_date)
        if df is None or df.empty:
            return None
        df = df.rename(columns={'trade_date': 'date'})
        df['date'] = pd.to_datetime(df['date']).dt.strftime('%Y-%m-%d')
        return df.reset_index(drop=True)
    except Exception as e:
        self.logger.error(f"get_adj_factor_batch {trade_date}: {e}")
        self.last_error = str(e)
        return None
get_daily_basic_batch
get_daily_basic_batch(trade_date: str, fields: str = 'ts_code,trade_date,turnover_rate') -> Optional[DataFrame]

按交易日拉取全市场每日指标(daily_basic 接口,约 0.3s/次)。

默认只拉换手率,调用方可通过 fields 参数扩展。 返回列取决于 fields,date 列已格式化为 YYYY-MM-DD。

源代码位于: adapters/tushare_adapter.py
def get_daily_basic_batch(self, trade_date: str,
                          fields: str = 'ts_code,trade_date,turnover_rate') -> Optional[pd.DataFrame]:
    """
    按交易日拉取全市场每日指标(daily_basic 接口,约 0.3s/次)。

    默认只拉换手率,调用方可通过 fields 参数扩展。
    返回列取决于 fields,date 列已格式化为 YYYY-MM-DD。
    """
    if not self.is_connected:
        if not self.connect():
            return None
    try:
        df = self._call(self.pro.daily_basic, trade_date=trade_date, fields=fields)
        if df is None or df.empty:
            return None
        df = df.rename(columns={'trade_date': 'date'})
        df['date'] = pd.to_datetime(df['date']).dt.strftime('%Y-%m-%d')
        return df.reset_index(drop=True)
    except Exception as e:
        self.logger.error(f"get_daily_basic_batch {trade_date}: {e}")
        self.last_error = str(e)
        return None
get_stock_name
get_stock_name(code: str) -> Optional[str]

通过 pro.stock_basic() 获取股票名称

源代码位于: adapters/tushare_adapter.py
def get_stock_name(self, code: str) -> Optional[str]:
    """通过 pro.stock_basic() 获取股票名称"""
    if not self.is_connected:
        if not self.connect():
            return None

    try:
        ts_code = self._convert_code(code)
        df = self._call(self.pro.stock_basic,
                        ts_code=ts_code,
                        fields='ts_code,name')
        if df is not None and not df.empty:
            return df.iloc[0]['name']
        return None
    except Exception as e:
        self.logger.debug(f"Tushare获取股票名称失败 {code}: {e}")
        return None
get_all_stock_names
get_all_stock_names() -> Optional[Dict[str, str]]

获取全市场股票名称映射 {纯6位代码: 名称},单次 API 调用

源代码位于: adapters/tushare_adapter.py
def get_all_stock_names(self) -> Optional[Dict[str, str]]:
    """获取全市场股票名称映射 {纯6位代码: 名称},单次 API 调用"""
    if not self.is_connected:
        if not self.connect():
            return None

    try:
        df = self._call(self.pro.stock_basic,
                        exchange='',
                        list_status='L',
                        fields='ts_code,name')
        if df is None or df.empty:
            return None
        return {row['ts_code'].split('.')[0]: row['name']
                for _, row in df.iterrows()}
    except Exception as e:
        self.logger.error(f"Tushare获取全市场股票名称失败: {e}")
        return None

Mootdx 适配器

mootdx_adapter

Mootdx数据源适配器

使用mootdx库获取通达信数据,主要用于K线和实时行情

MootdxAdapter

MootdxAdapter(name: str, config: Dict[str, Any])

Bases: DataSourceAdapter

Mootdx数据源适配器

源代码位于: adapters/mootdx_adapter.py
def __init__(self, name: str, config: Dict[str, Any]):
    super().__init__(name, config)
    self.client = None
connect
connect() -> bool

连接Mootdx数据源

源代码位于: adapters/mootdx_adapter.py
def connect(self) -> bool:
    """连接Mootdx数据源"""
    try:
        self.client = Quotes.factory('std')  # 使用标准版通达信
        self.is_connected = True
        self.logger.info(f"{self.name} 连接成功")
        return True
    except Exception as e:
        self.logger.error(f"{self.name} 连接失败: {e}")
        self.last_error = str(e)
        self.is_connected = False
        return False
disconnect
disconnect()

断开连接

源代码位于: adapters/mootdx_adapter.py
def disconnect(self):
    """断开连接"""
    self.client = None
    self.is_connected = False
    self.logger.info(f"{self.name} 已断开连接")
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') -> Optional[DataFrame]

获取K线数据

参数:

名称 类型 描述 默认
code str

股票代码

必需
freq str

频率

'd'
start_date Optional[str]

开始日期(mootdx不支持日期范围,此参数用于后过滤)

None
end_date Optional[str]

结束日期(mootdx不支持日期范围,此参数用于后过滤)

None
count Optional[int]

获取数量(offset参数)

None
adjust str

复权类型 'qfq'=前复权, 'hfq'=后复权, None=不复权

'qfq'

返回:

类型 描述
Optional[DataFrame]

DataFrame

源代码位于: adapters/mootdx_adapter.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'
) -> Optional[pd.DataFrame]:
    """
    获取K线数据

    Args:
        code: 股票代码
        freq: 频率
        start_date: 开始日期(mootdx不支持日期范围,此参数用于后过滤)
        end_date: 结束日期(mootdx不支持日期范围,此参数用于后过滤)
        count: 获取数量(offset参数)
        adjust: 复权类型 'qfq'=前复权, 'hfq'=后复权, None=不复权

    Returns:
        DataFrame
    """
    if not self.is_connected:
        if not self.connect():
            return None

    try:
        # 标准化代码
        code = self.normalize_code(code)

        # 转换频率
        if freq not in self.FREQ_MAP:
            self.logger.error(f"不支持的频率: {freq}")
            return None

        mootdx_freq = self.FREQ_MAP[freq]

        # 默认获取数量
        if count is None:
            count = 800  # mootdx默认值

        # 获取数据
        df = self.client.bars(
            symbol=code,
            frequency=mootdx_freq,
            offset=count,
            adjust=adjust
        )

        if df is None or df.empty:
            self.logger.warning(f"Mootdx未获取到数据: {code}")
            return None

        # 修复mootdx重复列问题: 原始数据包含'vol'和'volume'两列
        # 删除'vol'列,保留'volume'列(因为它是复权后的成交量)
        if 'vol' in df.columns and 'volume' in df.columns:
            df = df.drop(columns=['vol'])
            self.logger.debug("删除重复的vol列,保留复权后的volume列")

        # 标准化列名
        df = self.standardize_dataframe(df)

        # 统一成交量单位为"股" (mootdx返回的是手,需要*100)
        if 'volume' in df.columns:
            df['volume'] = df['volume'] * 100
            self.logger.debug(f"成交量单位转换: 手 -> 股 (*100)")

        # 日期过滤(如果提供了日期范围)
        if start_date:
            df = df[df['date'] >= start_date]
        if end_date:
            df = df[df['date'] <= end_date]

        # 确保必要列存在
        required_cols = ['date', 'open', 'high', 'low', 'close', 'volume']
        if not all(col in df.columns for col in required_cols):
            self.logger.error(f"数据列不完整: {df.columns.tolist()}")
            return None

        # 复权因子统一化处理(修复mootdx部分复权问题)
        if adjust in ['qfq', 'hfq'] and 'factor' in df.columns:
            # 前复权(qfq):以最新日期为基准(最新的factor),历史数据按factor调整
            # 找到最新交易日的factor(DataFrame末尾,时间最新)
            latest_factor = df['factor'].iloc[-1] if len(df) > 0 else 1.0

            # 如果存在factor变化,需要统一化
            if df['factor'].nunique() > 1:
                # 前复权公式: 调整后价格 = 原价格 × (最新factor / 历史factor)
                # 这样可以消除除权除息导致的价格跳跃
                price_cols = ['open', 'high', 'low', 'close']
                for col in price_cols:
                    if col in df.columns:
                        # 关键修复:历史数据(factor>latest_factor的)需要向下调整
                        # 最新数据(factor=latest_factor的)保持不变
                        df[col] = df[col] * (latest_factor / df['factor'])

                self.logger.info(
                    f"复权因子统一化: {code}, "
                    f"factor范围 {df['factor'].min():.6f}-{df['factor'].max():.6f}, "
                    f"以最新factor {latest_factor:.6f}为基准"
                )

        # 清理NaN值
        df = df.dropna(subset=['close'])  # 至少close不能为空

        if df.empty:
            self.logger.warning(f"Mootdx数据清理后为空: {code}")
            return None

        return df

    except Exception as e:
        self.logger.error(f"Mootdx获取K线失败 {code}: {e}")
        self.last_error = str(e)
        self.error_count += 1
        return None
get_valuation
get_valuation(code: str, start_date: Optional[str] = None, end_date: Optional[str] = None) -> Optional[DataFrame]

获取估值数据

注意: Mootdx不提供估值数据,返回None

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

    注意: Mootdx不提供估值数据,返回None
    """
    self.logger.warning(f"Mootdx不支持估值数据获取")
    return None
get_tick
get_tick(code: str) -> Optional[Dict[str, Any]]

获取实时tick数据

参数:

名称 类型 描述 默认
code str

股票代码

必需

返回:

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

实时行情字典

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

    Args:
        code: 股票代码

    Returns:
        实时行情字典
    """
    if not self.is_connected:
        if not self.connect():
            return None

    try:
        code = self.normalize_code(code)

        # 获取实时行情
        quotes = self.client.quotes([code])

        # 处理DataFrame或列表返回
        if isinstance(quotes, pd.DataFrame):
            if quotes.empty:
                self.logger.warning(f"Mootdx未获取到实时行情: {code}")
                return None
            quote = quotes.iloc[0].to_dict()
        elif isinstance(quotes, list):
            if not quotes or len(quotes) == 0:
                self.logger.warning(f"Mootdx未获取到实时行情: {code}")
                return None
            quote = quotes[0] if isinstance(quotes[0], dict) else quotes[0].to_dict()
        else:
            self.logger.warning(f"Mootdx未获取到实时行情: {code}")
            return None

        # 转换为统一格式
        tick_data = {
            'code': code,
            'name': quote.get('name', ''),
            'open': quote.get('open', 0),
            'high': quote.get('high', 0),
            'low': quote.get('low', 0),
            'close': quote.get('close', 0),
            'last': quote.get('price', 0),
            'volume': quote.get('vol', 0),
            'amount': quote.get('amount', 0),
            'bid': quote.get('bid1', 0),
            'ask': quote.get('ask1', 0),
            'yesterday_close': quote.get('last_close', 0)
        }

        return tick_data

    except Exception as e:
        self.logger.error(f"Mootdx获取tick失败 {code}: {e}")
        self.last_error = str(e)
        self.error_count += 1
        return None

Baostock 适配器

baostock_adapter

Baostock数据源适配器

使用baostock库获取数据,主要用于K线和估值数据

BaostockAdapter

BaostockAdapter(name: str, config: Dict[str, Any])

Bases: DataSourceAdapter

Baostock数据源适配器

源代码位于: adapters/baostock_adapter.py
def __init__(self, name: str, config: Dict[str, Any]):
    super().__init__(name, config)
    self.login_result = None
    self.api_key = config.get('api_key', '')
    self._consecutive_failures = 0
    self._last_failure_time = 0.0
connect
connect() -> bool

连接Baostock数据源(带超时保护,最多等待 _CONNECT_TIMEOUT 秒)

源代码位于: adapters/baostock_adapter.py
def connect(self) -> bool:
    """连接Baostock数据源(带超时保护,最多等待 _CONNECT_TIMEOUT 秒)"""
    import socket
    original_timeout = socket.getdefaulttimeout()
    socket.setdefaulttimeout(self._CONNECT_TIMEOUT)
    try:
        # 新版 baostock 0.9.x: 登录前应用 API Key(旧版无 set_API_key 时自动跳过,匿名访问)
        apply_api_key(self.api_key, self.logger)
        self.login_result = bs.login()
    except Exception as e:
        self.logger.error(f"{self.name} 连接失败: {e}")
        self.last_error = str(e)
        self.is_connected = False
        return False
    finally:
        socket.setdefaulttimeout(original_timeout)

    if self.login_result.error_code == '0':
        self.is_connected = True
        self.logger.info(f"{self.name} 连接成功")
        return True
    else:
        error_detail = describe_error(self.login_result.error_code, self.login_result.error_msg)
        self.logger.error(f"{self.name} 连接失败: {error_detail}")
        self.last_error = error_detail
        self.is_connected = False
        return False
disconnect
disconnect()

断开连接

源代码位于: adapters/baostock_adapter.py
def disconnect(self):
    """断开连接"""
    try:
        bs.logout()
        self.is_connected = False
        self.logger.info(f"{self.name} 已断开连接")
    except:
        pass
health_check
health_check() -> Dict[str, Any]

Baostock 专属健康检查。

不能复用基类实现,因为基类会调用 get_kline(),而 get_kline() 带有快速失败冷却机制。业务查询连续失败后,基类 health_check 会被 冷却逻辑直接短路为 None,导致健康检查无法主动探测恢复。

源代码位于: adapters/baostock_adapter.py
def health_check(self) -> Dict[str, Any]:
    """
    Baostock 专属健康检查。

    不能复用基类实现,因为基类会调用 get_kline(),而 get_kline()
    带有快速失败冷却机制。业务查询连续失败后,基类 health_check 会被
    冷却逻辑直接短路为 None,导致健康检查无法主动探测恢复。
    """
    from datetime import datetime, timedelta

    start_time = time.time()
    result = {
        'status': 'ok',
        'response_time': 0.0,
        'data_freshness': True,
        'error_message': None
    }

    try:
        if not self.is_connected:
            self.logger.debug(f"{self.name} 健康检查: 未连接,尝试重连")
            if not self.connect():
                result['status'] = 'error'
                result['error_message'] = f'连接失败: {self.last_error or "未知错误"}'
                result['response_time'] = time.time() - start_time
                return result

        end_date = datetime.now().strftime('%Y-%m-%d')
        start_date = (datetime.now() - timedelta(days=90)).strftime('%Y-%m-%d')

        rs = bs.query_history_k_data_plus(
            'sh.600000',
            'date,code,close,volume',
            start_date=start_date,
            end_date=end_date,
            frequency='d',
            adjustflag='2'
        )

        result['response_time'] = time.time() - start_time

        if rs.error_code != '0':
            error_detail = describe_error(rs.error_code, rs.error_msg)
            self.last_error = error_detail
            self._consecutive_failures += 1
            self._last_failure_time = time.time()
            self.error_count += 1
            result['status'] = 'error'
            result['error_message'] = f'Baostock查询失败: {error_detail}'
            return result

        rows = []
        while (rs.error_code == '0') and rs.next():
            rows.append(rs.get_row_data())

        if not rows:
            self._consecutive_failures += 1
            self._last_failure_time = time.time()
            self.error_count += 1
            result['status'] = 'error'
            result['error_message'] = '无法获取测试数据'
            return result

        latest_date = pd.to_datetime(rows[-1][0])
        days_diff = (datetime.now() - latest_date).days
        if days_diff > self.config.get('data_freshness_days', 3):
            result['status'] = 'warning'
            result['data_freshness'] = False
            result['error_message'] = f'数据不新鲜,最新数据日期: {latest_date.date()}'

        threshold = self.config.get('timeout', 5)
        if result['response_time'] > threshold:
            result['status'] = 'warning'
            result['error_message'] = f'响应时间过长: {result["response_time"]:.2f}秒'

        self._consecutive_failures = 0
        self.error_count = 0
        self.last_error = None

    except Exception as e:
        result['status'] = 'error'
        result['error_message'] = str(e)
        result['response_time'] = time.time() - start_time
        self.last_error = str(e)
        self.error_count += 1
        self._consecutive_failures += 1
        self._last_failure_time = time.time()
        self.logger.error(f"Baostock健康检查异常: {e}")

    return result
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') -> Optional[DataFrame]

获取K线数据

参数:

名称 类型 描述 默认
code str

股票代码

必需
freq str

频率

'd'
start_date Optional[str]

开始日期

None
end_date Optional[str]

结束日期

None
count Optional[int]

获取数量(baostock通过日期范围获取,此参数用于计算start_date)

None
adjust str

复权类型 'qfq'=前复权(adjustflag='2')

'qfq'

返回:

类型 描述
Optional[DataFrame]

DataFrame

源代码位于: adapters/baostock_adapter.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'
) -> Optional[pd.DataFrame]:
    """
    获取K线数据

    Args:
        code: 股票代码
        freq: 频率
        start_date: 开始日期
        end_date: 结束日期
        count: 获取数量(baostock通过日期范围获取,此参数用于计算start_date)
        adjust: 复权类型 'qfq'=前复权(adjustflag='2')

    Returns:
        DataFrame
    """
    # 快速失败:连续失败超过阈值且在冷却窗口内,立即返回 None
    if (self._consecutive_failures >= self._FAST_FAIL_THRESHOLD and
            time.time() - self._last_failure_time < self._FAST_FAIL_COOLDOWN):
        self.logger.debug(
            f"baostock快速失败(连续{self._consecutive_failures}次网络故障,"
            f"冷却剩余{self._FAST_FAIL_COOLDOWN - (time.time() - self._last_failure_time):.0f}s)"
        )
        return None

    if not self.is_connected:
        if not self.connect():
            self._consecutive_failures += 1
            self._last_failure_time = time.time()
            return None

    try:
        # 添加前缀
        bs_code = self._add_bs_prefix(code)

        # 转换频率
        if freq not in self.FREQ_MAP:
            self.logger.error(f"不支持的频率: {freq}")
            return None

        bs_freq = self.FREQ_MAP[freq]

        # 如果提供count但没有start_date,计算start_date
        if count and not start_date:
            from datetime import datetime, timedelta
            # 估算天数(count * 2倍,充分考虑节假日和周末)
            days = int(count * 2)
            start_date = (datetime.now() - timedelta(days=days)).strftime('%Y-%m-%d')

        # 默认结束日期为今天
        if not end_date:
            from datetime import datetime
            end_date = datetime.now().strftime('%Y-%m-%d')

        # 设置复权类型: 归一化为 baostock 接受的 '1'(后复权)/'2'(前复权)/'3'(不复权)
        adjustflag = normalize_adjustflag(adjust)

        # 字段定义:
        # - 日线/周线/月线:包含 turn(换手率),无 time 字段
        # - 分钟线:包含 time 字段,无 turn/preclose(baostock不支持)
        is_daily_or_higher = bs_freq in ('d', 'w', 'm')
        if is_daily_or_higher:
            fields = "date,code,open,high,low,close,preclose,volume,amount,turn,adjustflag"
        else:
            fields = "date,time,code,open,high,low,close,volume,amount,adjustflag"

        # 使用官方推荐迭代模式获取数据(比 get_data() 更稳定)
        rs = bs.query_history_k_data_plus(
            bs_code,
            fields,
            start_date=start_date,
            end_date=end_date,
            frequency=bs_freq,
            adjustflag=adjustflag
        )

        if rs.error_code != '0':
            error_detail = describe_error(rs.error_code, rs.error_msg)
            self.logger.error(f"Baostock查询失败 {bs_code}: {error_detail}")
            self.last_error = error_detail
            self._consecutive_failures += 1
            self._last_failure_time = time.time()
            return None

        # 迭代收集数据(官方推荐模式)
        data_list = []
        while (rs.error_code == '0') and rs.next():
            data_list.append(rs.get_row_data())

        if not data_list:
            self.logger.debug(f"Baostock未获取到数据: {bs_code}")
            return None

        df = pd.DataFrame(data_list, columns=rs.fields)

        # 分钟线:合并 date + time 为统一 date 列(格式 YYYY-MM-DD HH:MM:SS)
        if not is_daily_or_higher and 'time' in df.columns:
            df['date'] = df['date'] + ' ' + df['time'].str.strip()
            df = df.drop(columns=['time'])

        # 空字符串替换为 NaN,再做数值转换(baostock 停牌时 turn 为空字符串)
        numeric_cols = ['open', 'high', 'low', 'close', 'preclose', 'volume', 'amount', 'turn']
        for col in numeric_cols:
            if col in df.columns:
                df[col] = df[col].replace('', None)
                df[col] = pd.to_numeric(df[col], errors='coerce')

        # 如果指定了count,只取最近count条
        if count:
            df = df.tail(count)

        # 标准化
        df = self.standardize_dataframe(df)

        # 成功:重置连续失败计数
        self._consecutive_failures = 0
        return df

    except Exception as e:
        self.logger.error(f"Baostock获取K线失败 {code}: {e}")
        self.last_error = str(e)
        self.error_count += 1
        self._consecutive_failures += 1
        self._last_failure_time = time.time()
        return None
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,pcf_ncf等估值指标

源代码位于: adapters/baostock_adapter.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,pcf_ncf等估值指标
    """
    if not self.is_connected:
        if not self.connect():
            return None

    try:
        bs_code = self._add_bs_prefix(code)

        # 默认结束日期为今天
        if not end_date:
            from datetime import datetime
            end_date = datetime.now().strftime('%Y-%m-%d')

        # 默认开始日期为1年前
        if not start_date:
            from datetime import datetime, timedelta
            start_date = (datetime.now() - timedelta(days=365)).strftime('%Y-%m-%d')

        # 查询估值数据(迭代模式)
        rs = bs.query_history_k_data_plus(
            bs_code,
            "date,code,peTTM,pbMRQ,psTTM,pcfNcfTTM",
            start_date=start_date,
            end_date=end_date,
            frequency="d",
            adjustflag="3"
        )

        if rs.error_code != '0':
            error_detail = describe_error(rs.error_code, rs.error_msg)
            self.logger.error(f"Baostock估值查询失败 {bs_code}: {error_detail}")
            self.last_error = error_detail
            return None

        data_list = []
        while (rs.error_code == '0') and rs.next():
            data_list.append(rs.get_row_data())

        if not data_list:
            self.logger.warning(f"Baostock未获取到估值数据: {bs_code}")
            return None

        df = pd.DataFrame(data_list, columns=rs.fields)

        # 重命名列
        df = df.rename(columns={
            'peTTM': 'pe_ttm',
            'pbMRQ': 'pb',
            'psTTM': 'ps_ttm',
            'pcfNcfTTM': 'pcf_ncf'
        })

        # 数值转换
        numeric_cols = ['pe_ttm', 'pb', 'ps_ttm', 'pcf_ncf']
        for col in numeric_cols:
            if col in df.columns:
                df[col] = pd.to_numeric(df[col], errors='coerce')

        return df

    except Exception as e:
        self.logger.error(f"Baostock获取估值失败 {code}: {e}")
        self.last_error = str(e)
        self.error_count += 1
        return None
get_tick
get_tick(code: str) -> Optional[Dict[str, Any]]

获取实时tick数据

注意: Baostock不提供实时tick,可以使用5分钟K线最新一条模拟

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

    注意: Baostock不提供实时tick,可以使用5分钟K线最新一条模拟
    """
    try:
        # 获取最近2条5分钟K线
        df = self.get_kline(code, freq='5m', count=2, adjust='qfq')

        if df is None or df.empty:
            return None

        # 使用最新一条数据
        latest = df.iloc[-1]

        tick_data = {
            'code': self.normalize_code(code),
            'open': float(latest['open']),
            'high': float(latest['high']),
            'low': float(latest['low']),
            'close': float(latest['close']),
            'last': float(latest['close']),
            'volume': float(latest['volume']),
            'amount': float(latest.get('amount', 0)),
            'yesterday_close': float(df.iloc[-2]['close']) if len(df) > 1 else float(latest['close'])
        }

        return tick_data

    except Exception as e:
        self.logger.error(f"Baostock模拟tick失败 {code}: {e}")
        self.last_error = str(e)
        return None

xtquant 适配器

xtquant_adapter

xtquant数据源适配器

使用xtquant(迅投QMT)获取实时行情数据 参考: https://github.com/weihong-su/miniQMT 和 https://github.com/weihong-su/khQuant

增强功能(2025-11-19): - 智能重试机制: 带指数退避的重试策略 - 连接保活: 心跳检测和自动重连 - 连接统计: 重试、重连次数等监控指标

XtquantAdapter

XtquantAdapter(name: str, config: Dict[str, Any])

Bases: DataSourceAdapter

xtquant数据源适配器 - 增强版(带重试和心跳保活)

源代码位于: adapters/xtquant_adapter.py
def __init__(self, name: str, config: Dict[str, Any]):
    super().__init__(name, config)
    self.xt_data = None
    self.fallback_to_5m = config.get('fallback_to_5m', True)
    self._lib_loader = None  # LibLoader实例
    self.cache_manager = None  # 缓存管理器(由外部注入)

    # 缓存配置
    self.cache_enabled = config.get('cache_enabled', True)  # 是否启用缓存
    self.cache_only_daily = config.get('cache_only_daily', True)  # 只缓存日K线

    # 重试配置
    self.retry_times = config.get('retry_times', 2)
    self.retry_delay = config.get('retry_delay', 0.5)
    self.retry_backoff = config.get('retry_backoff_factor', 2.0)
    self.max_retry_delay = config.get('max_retry_delay', 5.0)

    # 连接配置
    self.connect_timeout = config.get('connect_timeout', 10)
    self.connect_retry_times = config.get('connect_retry_times', 3)
    self.connect_retry_delay = config.get('connect_retry_delay', 1.0)

    # 心跳保活配置
    self.heartbeat_enabled = config.get('heartbeat_enabled', True)
    self.heartbeat_interval = config.get('heartbeat_interval', 30)  # 秒
    self.heartbeat_timeout = config.get('heartbeat_timeout', 5)  # 秒
    self.auto_reconnect = config.get('auto_reconnect', True)
    self.max_reconnect_attempts = config.get('max_reconnect_attempts', 5)

    # 心跳线程相关
    self._heartbeat_thread = None
    self._heartbeat_running = False
    self._heartbeat_lock = threading.Lock()
    self.last_heartbeat = None

    # 统计数据
    self.retry_stats = RetryStats()
    self.connection_stats = {
        'connect_count': 0,
        'disconnect_count': 0,
        'reconnect_count': 0,
        'heartbeat_failures': 0,
        'last_connect_time': None,
        'last_disconnect_time': None
    }
connect
connect() -> bool

连接xtquant数据源 - 三层验证机制

三层验证: 1. QMT服务器连接 (xtdata.connect) 2. 数据接口可用性 (get_trading_dates) 3. 市场数据有效性 (get_full_tick + 价格>0.1)

注意: xtquant依赖本地QMT客户端运行 券商可能不定期关闭接口,交易时段较稳定

源代码位于: adapters/xtquant_adapter.py
def connect(self) -> bool:
    """
    连接xtquant数据源 - 三层验证机制

    三层验证:
    1. QMT服务器连接 (xtdata.connect)
    2. 数据接口可用性 (get_trading_dates)
    3. 市场数据有效性 (get_full_tick + 价格>0.1)

    注意: xtquant依赖本地QMT客户端运行
    券商可能不定期关闭接口,交易时段较稳定
    """
    import time

    try:
        # 尝试导入xtquant库(使用LibLoader支持内置库)
        try:
            # 优先尝试使用LibLoader
            if self._lib_loader is None:
                try:
                    from ..utils.lib_loader import get_lib_loader
                    from ..config import Config
                except (ImportError, ValueError):
                    # 绝对导入回退
                    from utils.lib_loader import get_lib_loader
                    from config import Config
                cfg = Config()
                self._lib_loader = get_lib_loader({'use_builtin_libs': cfg.get('use_builtin_libs', True)})

            # 加载xtquant包(确保lib目录在sys.path中)
            xtquant_module = self._lib_loader.load_library('xtquant', fallback=True)
            if xtquant_module:
                # xtdata是xtquant的子模块,需要导入
                from xtquant import xtdata
                self.xt_data = xtdata
            else:
                # LibLoader加载失败,尝试直接导入
                from xtquant import xtdata
                self.xt_data = xtdata
        except ImportError as e:
            self.logger.warning(f"{self.name} xtquant库未安装: {e}")
            self.is_connected = False
            return False

        # 第一层验证: QMT服务器连接
        try:
            # 注意: connect()返回客户端对象(不是状态码)
            # 连接成功返回对象,失败抛出异常
            connect_result = self.xt_data.connect()
            if connect_result is None:
                self.logger.warning(f"{self.name} QMT服务器连接失败: 返回None")
                self.is_connected = False
                return False
            self.logger.debug(f"{self.name} 第1层验证通过: QMT服务器已连接")
        except Exception as e:
            self.logger.warning(f"{self.name} QMT服务器连接异常: {e}")
            self.is_connected = False
            return False

        # 第二层验证: 数据接口可用性
        try:
            trading_dates = self.xt_data.get_trading_dates('SH', start_time='20241101', end_time='20241201')
            if trading_dates is None or len(trading_dates) == 0:
                self.logger.warning(f"{self.name} 数据接口不可用: get_trading_dates返回空")
                self.is_connected = False
                return False
            self.logger.debug(f"{self.name} 第2层验证通过: 数据接口可用 ({len(trading_dates)}个交易日)")
        except Exception as e:
            self.logger.warning(f"{self.name} 数据接口验证失败: {e}")
            self.is_connected = False
            return False

        # 第三层验证: 市场数据有效性
        # 使用浦发银行(600000)测试,比茅台更稳定
        try:
            test_tick = self.xt_data.get_full_tick(['600000.SH'])
            if not test_tick or '600000.SH' not in test_tick:
                self.logger.warning(f"{self.name} 市场数据获取失败: get_full_tick返回空")
                self.is_connected = False
                return False

            # 验证价格有效性(关键: 价格>0.1)
            tick_data = test_tick['600000.SH']
            last_price = tick_data.get('lastPrice', 0)
            if last_price < 0.1:
                self.logger.warning(
                    f"{self.name} 市场数据无效: lastPrice={last_price:.4f} (应>0.1)"
                )
                self.is_connected = False
                return False

            self.logger.debug(
                f"{self.name} 第3层验证通过: 市场数据有效 (价格={last_price:.2f})"
            )
        except Exception as e:
            self.logger.warning(f"{self.name} 市场数据验证失败: {e}")
            self.is_connected = False
            return False

        # 三层验证全部通过
        self.is_connected = True
        self.connection_stats['connect_count'] += 1
        self.connection_stats['last_connect_time'] = datetime.now()
        self.logger.info(f"{self.name} 连接成功 (三层验证通过)")

        # 自动启动心跳检测
        if self.heartbeat_enabled and not self._heartbeat_running:
            self._start_heartbeat()

        return True

    except Exception as e:
        self.logger.error(f"{self.name} 连接失败: {e}")
        self.last_error = str(e)
        self.is_connected = False
        return False
get_connection_stats
get_connection_stats() -> Dict[str, Any]

获取连接统计信息

返回:

类型 描述
Dict[str, Any]

连接统计字典

源代码位于: adapters/xtquant_adapter.py
def get_connection_stats(self) -> Dict[str, Any]:
    """
    获取连接统计信息

    Returns:
        连接统计字典
    """
    with self._heartbeat_lock:
        stats = self.connection_stats.copy()
        stats['is_connected'] = self.is_connected
        stats['last_heartbeat'] = self.last_heartbeat
        stats['heartbeat_running'] = self._heartbeat_running

        # 计算连接健康度
        if self.last_heartbeat:
            elapsed = (datetime.now() - self.last_heartbeat).total_seconds()
            stats['heartbeat_age_seconds'] = elapsed
            if elapsed < self.heartbeat_interval * 2:
                stats['connection_health'] = 'good'
            elif elapsed < self.heartbeat_interval * 5:
                stats['connection_health'] = 'degraded'
            else:
                stats['connection_health'] = 'poor'
        else:
            stats['heartbeat_age_seconds'] = None
            stats['connection_health'] = 'unknown'

        # 添加重试统计
        stats['retry_stats'] = self.retry_stats.get_stats()

        return stats
disconnect
disconnect()

断开连接

源代码位于: adapters/xtquant_adapter.py
def disconnect(self):
    """断开连接"""
    # 停止心跳检测
    self._stop_heartbeat()

    self.xt_data = None
    self.is_connected = False
    self.connection_stats['disconnect_count'] += 1
    self.connection_stats['last_disconnect_time'] = datetime.now()
    self.logger.info(f"{self.name} 已断开连接")
set_cache_manager
set_cache_manager(cache_manager)

设置缓存管理器(由外部注入)

源代码位于: adapters/xtquant_adapter.py
def set_cache_manager(self, cache_manager):
    """设置缓存管理器(由外部注入)"""
    self.cache_manager = cache_manager
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') -> Optional[DataFrame]

获取K线数据(download_history_data + get_market_data_ex 工作流)

工作流: download_history_data → 等待 → get_market_data_ex - download_history_data: 异步下载到本地(返回None是正常的) - get_market_data_ex: 从本地读取已下载数据(比get_local_data更可靠)

阶段2优化: 集成智能缓存机制 - 日K线优先从缓存获取 - 缓存未命中时从xtquant获取并自动缓存 - 盘中时段不缓存当日数据,盘后自动缓存

参数:

名称 类型 描述 默认
code str

股票代码(6位数字或带前缀sh./sz.)

必需
freq str

K线周期('1m'/'5m'/'15m'/'30m'/'60m'/'d')

'd'
start_date Optional[str]

开始日期 'YYYYMMDD' 或 'YYYY-MM-DD'(可选)

None
end_date Optional[str]

结束日期 'YYYYMMDD' 或 'YYYY-MM-DD'(可选)

None
count Optional[int]

获取条数(可选)

None
adjust str

复权类型('qfq'前复权)

'qfq'

返回:

名称 类型 描述
DataFrame Optional[DataFrame]

包含date,open,high,low,close,volume,amount列,失败返回None

源代码位于: adapters/xtquant_adapter.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'
) -> Optional[pd.DataFrame]:
    """
    获取K线数据(download_history_data + get_market_data_ex 工作流)

    工作流: download_history_data → 等待 → get_market_data_ex
    - download_history_data: 异步下载到本地(返回None是正常的)
    - get_market_data_ex: 从本地读取已下载数据(比get_local_data更可靠)

    阶段2优化: 集成智能缓存机制
    - 日K线优先从缓存获取
    - 缓存未命中时从xtquant获取并自动缓存
    - 盘中时段不缓存当日数据,盘后自动缓存

    Args:
        code: 股票代码(6位数字或带前缀sh./sz.)
        freq: K线周期('1m'/'5m'/'15m'/'30m'/'60m'/'d')
        start_date: 开始日期 'YYYYMMDD' 或 'YYYY-MM-DD'(可选)
        end_date: 结束日期 'YYYYMMDD' 或 'YYYY-MM-DD'(可选)
        count: 获取条数(可选)
        adjust: 复权类型('qfq'前复权)

    Returns:
        DataFrame: 包含date,open,high,low,close,volume,amount列,失败返回None
    """
    import time
    from datetime import datetime, timedelta

    # ========== 阶段2: 智能缓存检查 ==========
    if self.cache_enabled and self.cache_manager and freq == 'd':
        # 只缓存日K线
        cache_start_date = start_date.replace('-', '') if start_date else None
        cache_end_date = end_date.replace('-', '') if end_date else None

        cached_df = self.cache_manager.get_cached_kline(
            code=code,
            start_date=cache_start_date,
            end_date=cache_end_date,
            count=count
        )

        if cached_df is not None and not cached_df.empty:
            self.logger.debug(f"xtquant缓存命中: {code} ({freq}) {len(cached_df)}条")
            cached_df.attrs['source'] = 'xtquant_cache'
            return cached_df
        else:
            self.logger.debug(f"xtquant缓存未命中: {code} ({freq}), 从API获取")
    # ==========================================

    if not self.is_connected:
        if not self.connect():
            return None

    # 周期映射(基于测试验证的可用周期)
    period_map = {
        '1m': '1m',
        '5m': '5m',
        '15m': '15m',
        '30m': '30m',
        '60m': '60m',
        'd': '1d'
    }

    if freq not in period_map:
        self.logger.error(f"不支持的周期: {freq}, 仅支持: {list(period_map.keys())}")
        return None

    xt_period = period_map[freq]
    xt_code = self._convert_code(code)

    # 标准化日期格式(移除横线,转为YYYYMMDD)
    if start_date:
        start_date = start_date.replace('-', '')
    if end_date:
        end_date = end_date.replace('-', '')

    # 计算时间范围(如果未提供)
    if not start_date or not end_date:
        end_dt = datetime.now()

        if count and count > 0:
            days_per_bar = {
                '1m': 0.01, '5m': 0.05, '15m': 0.15,
                '30m': 0.3, '60m': 0.5, '1d': 1
            }
            days = int(count * days_per_bar.get(xt_period, 1) * 1.5)
            start_dt = end_dt - timedelta(days=max(days, 30))
        else:
            start_dt = end_dt - timedelta(days=30)

        start_date = start_dt.strftime('%Y%m%d')
        end_date = end_dt.strftime('%Y%m%d')

    # 重试机制
    max_retries = 2
    for attempt in range(max_retries):
        try:
            # 步骤1: 下载历史数据到本地
            # download_history_data 是异步操作,返回None是正常的
            self.logger.debug(f"xtquant下载历史数据: {xt_code} ({xt_period}) {start_date}-{end_date}")

            self.xt_data.download_history_data(
                stock_code=xt_code,
                period=xt_period,
                start_time=start_date,
                end_time=end_date
            )

            # 步骤2: 等待下载完成
            # 分钟线数据量大,等待更久
            wait_time = 5 if xt_period in ['1m', '5m'] else 3
            time.sleep(wait_time)

            # 步骤3: 从本地读取数据
            # 使用 get_market_data_ex(比 get_local_data 更可靠)
            data = self.xt_data.get_market_data_ex(
                field_list=[],  # 空列表=获取所有字段
                stock_list=[xt_code],
                period=xt_period,
                start_time=start_date,
                end_time=end_date,
                count=-1,
                dividend_type='front' if adjust == 'qfq' else 'none',
                fill_data=True
            )

            # 步骤4: 处理返回数据
            df = None
            if data is not None:
                if isinstance(data, dict):
                    if xt_code in data and data[xt_code] is not None:
                        df = data[xt_code]
                    else:
                        self.logger.debug(f"xtquant数据中没有{xt_code}的数据")
                elif hasattr(data, 'empty'):
                    df = data
            else:
                self.logger.warning(f"xtquant返回None: {code} ({freq})")

            if df is None or df.empty:
                self.logger.warning(
                    f"xtquant未获取到K线数据: {code} ({freq}) "
                    f"[尝试{attempt+1}/{max_retries}] "
                    f"(可能数据未订阅或盘后不可用)"
                )
                if attempt < max_retries - 1:
                    time.sleep(1)
                    continue
                return None

            # 步骤5: 转换为StockDataMaster标准格式
            result_df = self._convert_to_standard_format(df, freq)

            if result_df is None or result_df.empty:
                self.logger.warning(f"xtquant数据转换失败: {code} ({freq})")
                if attempt < max_retries - 1:
                    time.sleep(1)
                    continue
                return None

            # 步骤6: 根据count截取数据
            if count and count > 0:
                result_df = result_df.tail(count).reset_index(drop=True)

            # 步骤7: 数据校验
            if not self._validate_kline_data(result_df, code, freq):
                self.logger.warning(f"xtquant K线数据校验失败: {code} ({freq})")
                if attempt < max_retries - 1:
                    time.sleep(1)
                    continue
                return None

            # 步骤8: 成功后保存到缓存
            # 注意:xtquant 单源不再自动标 validated=1。读取端 (CacheManager.get_cached_kline)
            # 已收紧为"必须 validated=1 且 source2 不是 none/null/'' ",否则不命中。
            # 单源缓存仅作为"占位",等待 warmup 阶段做真正的双源升级。
            if self.cache_enabled and self.cache_manager and freq == 'd':
                try:
                    self.cache_manager.save_to_cache(
                        code=code,
                        df=result_df,
                        source1='xtquant',
                        source2=None,
                        validated=False
                    )
                    self.logger.debug(f"xtquant数据已保存到缓存(单源待校验): {code} ({freq})")
                except Exception as cache_error:
                    # 缓存失败不影响数据返回
                    self.logger.warning(f"缓存保存失败: {cache_error}")

            # 成功
            self.logger.info(f"xtquant获取K线成功: {code} ({freq}) {len(result_df)}条")
            result_df.attrs['source'] = 'xtquant'
            return result_df

        except Exception as e:
            self.logger.error(f"xtquant获取K线失败: {code} ({freq}) [尝试{attempt+1}/{max_retries}] {e}")
            self.last_error = str(e)
            self.error_count += 1
            if attempt < max_retries - 1:
                time.sleep(1)
                continue

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

获取估值数据

注意: xtquant不提供估值数据

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

    注意: xtquant不提供估值数据
    """
    self.logger.warning(f"xtquant不支持估值数据获取")
    return None
get_tick
get_tick(code: str) -> Optional[Dict[str, Any]]

获取实时tick数据 - 增强版(带智能重试)

参数:

名称 类型 描述 默认
code str

股票代码

必需

返回:

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

实时行情字典

重试策略: 1. 连接失败: 使用_connect_with_retry重试连接 2. 数据获取失败: 重试retry_times次,带指数退避 3. 数据校验失败: 不重试(数据问题,非临时错误)

源代码位于: adapters/xtquant_adapter.py
def get_tick(self, code: str) -> Optional[Dict[str, Any]]:
    """
    获取实时tick数据 - 增强版(带智能重试)

    Args:
        code: 股票代码

    Returns:
        实时行情字典

    重试策略:
    1. 连接失败: 使用_connect_with_retry重试连接
    2. 数据获取失败: 重试retry_times次,带指数退避
    3. 数据校验失败: 不重试(数据问题,非临时错误)
    """
    xt_code = self._convert_code(code)
    delay = self.retry_delay
    attempts = 0

    for attempt in range(self.retry_times + 1):
        attempts += 1
        try:
            # 检查连接状态
            if not self.is_connected:
                if not self._connect_with_retry():
                    if attempt < self.retry_times:
                        self.logger.debug(
                            f"get_tick连接失败,重试 {attempt+1}/{self.retry_times}"
                        )
                        time.sleep(delay)
                        delay = min(delay * self.retry_backoff, self.max_retry_delay)
                        continue
                    self.retry_stats.record_failure(attempts)
                    return None

            # 获取实时行情
            tick = self.xt_data.get_full_tick([xt_code])

            if not tick or xt_code not in tick:
                self.logger.warning(
                    f"xtquant未获取到实时行情: {code} [尝试{attempt+1}/{self.retry_times+1}]"
                )
                if attempt < self.retry_times:
                    time.sleep(delay)
                    delay = min(delay * self.retry_backoff, self.max_retry_delay)
                    continue
                self.retry_stats.record_failure(attempts)
                return None

            tick_data_raw = tick[xt_code]

            # 转换为统一格式
            tick_data = {
                'code': self.normalize_code(code),
                'name': tick_data_raw.get('stockName', ''),
                'open': tick_data_raw.get('open', 0),
                'high': tick_data_raw.get('high', 0),
                'low': tick_data_raw.get('low', 0),
                'close': tick_data_raw.get('lastClose', 0),
                'last': tick_data_raw.get('lastPrice', 0),
                'volume': tick_data_raw.get('volume', 0),
                'amount': tick_data_raw.get('amount', 0),
                'bid': tick_data_raw.get('bidPrice', [0])[0] if tick_data_raw.get('bidPrice') else 0,
                'ask': tick_data_raw.get('askPrice', [0])[0] if tick_data_raw.get('askPrice') else 0,
                'yesterday_close': tick_data_raw.get('lastClose', 0),
                'source': 'xtquant',
                # 修复: 添加时间戳字段映射
                'timestamp': self._convert_tick_timestamp(tick_data_raw.get('time'))
            }

            # 数据校验(校验失败不重试)
            if not self._validate_tick_data(tick_data, code):
                self.logger.warning(f"tick数据校验失败: {code}")
                self.retry_stats.record_failure(attempts)
                return None

            # 成功
            self.retry_stats.record_success(attempts)
            return tick_data

        except Exception as e:
            self.logger.error(
                f"xtquant获取tick失败 {code}: {e} [尝试{attempt+1}/{self.retry_times+1}]"
            )
            self.last_error = str(e)
            self.error_count += 1

            if attempt < self.retry_times:
                time.sleep(delay)
                delay = min(delay * self.retry_backoff, self.max_retry_delay)
                continue

            self.retry_stats.record_failure(attempts)
            return None

    self.retry_stats.record_failure(attempts)
    return None
get_realtime_quotes
get_realtime_quotes(codes: List[str]) -> Optional[Dict[str, Dict[str, Any]]]

批量获取实时行情 - 新增方法

参数:

名称 类型 描述 默认
codes List[str]

股票代码列表

必需

返回:

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

{code: tick_data} 字典,失败返回None

性能优势: - 批量获取比循环调用get_tick效率高 - 单次网络请求获取多只股票

源代码位于: adapters/xtquant_adapter.py
def get_realtime_quotes(self, codes: List[str]) -> Optional[Dict[str, Dict[str, Any]]]:
    """
    批量获取实时行情 - 新增方法

    Args:
        codes: 股票代码列表

    Returns:
        {code: tick_data} 字典,失败返回None

    性能优势:
    - 批量获取比循环调用get_tick效率高
    - 单次网络请求获取多只股票
    """
    if not codes:
        return {}

    # 转换代码格式
    xt_codes = [self._convert_code(code) for code in codes]
    delay = self.retry_delay

    for attempt in range(self.retry_times + 1):
        try:
            if not self.is_connected:
                if not self._connect_with_retry():
                    if attempt < self.retry_times:
                        time.sleep(delay)
                        delay = min(delay * self.retry_backoff, self.max_retry_delay)
                        continue
                    return None

            # 批量获取
            ticks = self.xt_data.get_full_tick(xt_codes)

            if not ticks:
                if attempt < self.retry_times:
                    self.logger.warning(
                        f"批量获取实时行情为空 [尝试{attempt+1}/{self.retry_times+1}]"
                    )
                    time.sleep(delay)
                    delay = min(delay * self.retry_backoff, self.max_retry_delay)
                    continue
                return None

            # 转换格式
            result = {}
            for code, xt_code in zip(codes, xt_codes):
                if xt_code in ticks:
                    tick_raw = ticks[xt_code]
                    tick_data = {
                        'code': self.normalize_code(code),
                        'name': tick_raw.get('stockName', ''),
                        'last': tick_raw.get('lastPrice', 0),
                        'open': tick_raw.get('open', 0),
                        'high': tick_raw.get('high', 0),
                        'low': tick_raw.get('low', 0),
                        'volume': tick_raw.get('volume', 0),
                        'amount': tick_raw.get('amount', 0),
                        'source': 'xtquant'
                    }
                    result[code] = tick_data

            return result if result else None

        except Exception as e:
            self.logger.error(f"批量获取实时行情失败: {e}")
            if attempt < self.retry_times:
                time.sleep(delay)
                delay = min(delay * self.retry_backoff, self.max_retry_delay)
                continue
            return None

    return None
health_check
health_check() -> Dict[str, Any]

健康检查 - 时段感知机制

检查策略: 1. 交易时段(9:15-15:00): 严格检查,数据必须有效 2. 非交易时段: 宽松检查,允许数据为空但连接正常 3. 强制重连: 每次检查前重新连接(后台线程可能断开) 4. 数据校验: 使用_validate_tick_data验证数据有效性

源代码位于: adapters/xtquant_adapter.py
def health_check(self) -> Dict[str, Any]:
    """
    健康检查 - 时段感知机制

    检查策略:
    1. 交易时段(9:15-15:00): 严格检查,数据必须有效
    2. 非交易时段: 宽松检查,允许数据为空但连接正常
    3. 强制重连: 每次检查前重新连接(后台线程可能断开)
    4. 数据校验: 使用_validate_tick_data验证数据有效性
    """
    import time
    from datetime import datetime

    result = {
        'status': 'ok',
        'response_time': 0.0,
        'data_freshness': True,
        'error_message': None
    }

    try:
        # 判断当前是否为交易时段 (9:15-15:00)
        now = datetime.now()
        current_time = now.time()
        from datetime import time as time_type
        market_start = time_type(9, 15)
        market_end = time_type(15, 0)
        is_trading_hours = market_start <= current_time <= market_end

        # 强制重连(关键: 后台线程可能已断开)
        if not self.is_connected:
            self.logger.debug(f"{self.name} 健康检查: 未连接,尝试重连")
            if not self.connect():
                result['status'] = 'error'
                result['error_message'] = 'xtquant连接失败'
                return result

        # 测试获取实时数据
        start_time = time.time()
        test_tick = self.get_tick('600000')  # 使用浦发银行测试
        elapsed = (time.time() - start_time) * 1000  # 毫秒
        result['response_time'] = round(elapsed, 2)

        # 数据有效性校验
        if test_tick is None:
            if is_trading_hours:
                # 交易时段: 数据为空视为严重错误
                result['status'] = 'error'
                result['error_message'] = 'xtquant交易时段数据获取失败'
            else:
                # 非交易时段: 数据为空仅为警告
                result['status'] = 'warning'
                result['error_message'] = 'xtquant非交易时段数据为空(正常)'
        else:
            # 使用_validate_tick_data严格校验
            is_valid = self._validate_tick_data(test_tick, '600000')
            if not is_valid:
                if is_trading_hours:
                    # 交易时段: 数据无效视为错误
                    result['status'] = 'error'
                    result['error_message'] = 'xtquant数据校验失败(价格<0.1或逻辑错误)'
                else:
                    # 非交易时段: 数据无效视为警告
                    result['status'] = 'warning'
                    result['error_message'] = 'xtquant数据校验失败(非交易时段)'
            else:
                # 数据有效
                result['status'] = 'ok'
                result['data_freshness'] = True

        # 响应时间检查(交易时段要求<1秒)
        if is_trading_hours and result['response_time'] > 1000:
            result['status'] = 'warning'
            result['error_message'] = f'响应时间过长: {result["response_time"]}ms'

    except Exception as e:
        result['status'] = 'error'
        result['error_message'] = str(e)
        self.error_count += 1
        self.logger.error(f"{self.name} 健康检查异常: {e}")

    return result
get_financial_data
get_financial_data(code: str, table: str = 'Balance', start_date: Optional[str] = None, end_date: Optional[str] = None) -> Optional[DataFrame]

获取财务数据(新增功能 - 阶段1)

参数:

名称 类型 描述 默认
code str

股票代码

必需
table str

财务表类型 - 'Balance': 资产负债表 - 'Income': 利润表 - 'CashFlow': 现金流量表

'Balance'
start_date Optional[str]

开始日期 'YYYYMMDD'

None
end_date Optional[str]

结束日期 'YYYYMMDD'

None

返回:

名称 类型 描述
DataFrame Optional[DataFrame]

财务数据,失败返回None

性能优势: - xtquant提供官方财务数据,权威性高 - 直接本地缓存,速度快

源代码位于: adapters/xtquant_adapter.py
def get_financial_data(
    self,
    code: str,
    table: str = 'Balance',
    start_date: Optional[str] = None,
    end_date: Optional[str] = None
) -> Optional[pd.DataFrame]:
    """
    获取财务数据(新增功能 - 阶段1)

    Args:
        code: 股票代码
        table: 财务表类型
            - 'Balance': 资产负债表
            - 'Income': 利润表
            - 'CashFlow': 现金流量表
        start_date: 开始日期 'YYYYMMDD'
        end_date: 结束日期 'YYYYMMDD'

    Returns:
        DataFrame: 财务数据,失败返回None

    性能优势:
    - xtquant提供官方财务数据,权威性高
    - 直接本地缓存,速度快
    """
    if not self.is_connected:
        if not self.connect():
            return None

    try:
        from datetime import datetime, timedelta

        xt_code = self._convert_code(code)

        # 默认获取最近3年数据
        if not start_date or not end_date:
            end_dt = datetime.now()
            start_dt = end_dt - timedelta(days=3*365)
            start_date = start_dt.strftime('%Y%m%d')
            end_date = end_dt.strftime('%Y%m%d')

        self.logger.debug(f"xtquant获取财务数据: {code} {table} {start_date}-{end_date}")

        # 调用xtquant财务数据接口
        financial_data = self.xt_data.get_financial_data(
            stock_list=[xt_code],
            table_list=[table],
            start_time=start_date,
            end_time=end_date,
            report_type='report_time'  # 按报告期
        )

        # 处理dict格式返回值
        if financial_data is None:
            self.logger.warning(f"xtquant未获取到财务数据: {code} {table} (返回None)")
            return None

        # xtquant财务数据返回dict格式
        if isinstance(financial_data, dict):
            if xt_code in financial_data:
                df = financial_data[xt_code]
                if df is None:
                    self.logger.warning(f"xtquant未获取到财务数据: {code} {table} (数据为None)")
                    return None
                if hasattr(df, 'empty') and df.empty:
                    self.logger.warning(f"xtquant未获取到财务数据: {code} {table} (数据为空)")
                    return None
                financial_data = df
            else:
                self.logger.warning(f"xtquant未获取到财务数据: {code} {table} (code不在返回dict中)")
                return None

        # 数据标准化
        result_df = financial_data.reset_index()

        # 转换日期格式
        if 'm_anntime' in result_df.columns:
            result_df['date'] = pd.to_datetime(result_df['m_anntime']).dt.strftime('%Y-%m-%d')

        self.logger.info(f"xtquant获取财务数据成功: {code} {table} {len(result_df)}条")
        return result_df

    except Exception as e:
        self.logger.error(f"xtquant获取财务数据失败: {code} {table} {e}")
        self.last_error = str(e)
        self.error_count += 1
        return None
get_trading_calendar
get_trading_calendar(market: str = 'SH', start_date: Optional[str] = None, end_date: Optional[str] = None) -> Optional[DataFrame]

获取交易日历(新增功能 - 阶段1)

参数:

名称 类型 描述 默认
market str

市场代码 'SH'(上交所)或 'SZ'(深交所)

'SH'
start_date Optional[str]

开始日期 'YYYYMMDD'

None
end_date Optional[str]

结束日期 'YYYYMMDD'

None

返回:

名称 类型 描述
DataFrame Optional[DataFrame]

交易日历,包含trading_date列,失败返回None

性能优势: - 本地数据,毫秒级响应 - 支持未来交易日查询

源代码位于: adapters/xtquant_adapter.py
def get_trading_calendar(
    self,
    market: str = 'SH',
    start_date: Optional[str] = None,
    end_date: Optional[str] = None
) -> Optional[pd.DataFrame]:
    """
    获取交易日历(新增功能 - 阶段1)

    Args:
        market: 市场代码 'SH'(上交所)或 'SZ'(深交所)
        start_date: 开始日期 'YYYYMMDD'
        end_date: 结束日期 'YYYYMMDD'

    Returns:
        DataFrame: 交易日历,包含trading_date列,失败返回None

    性能优势:
    - 本地数据,毫秒级响应
    - 支持未来交易日查询
    """
    if not self.is_connected:
        if not self.connect():
            return None

    try:
        from datetime import datetime, timedelta

        # 默认获取最近1年到未来30天
        if not start_date or not end_date:
            start_dt = datetime.now() - timedelta(days=365)
            end_dt = datetime.now() + timedelta(days=30)
            start_date = start_dt.strftime('%Y%m%d')
            end_date = end_dt.strftime('%Y%m%d')

        self.logger.debug(f"xtquant获取交易日历: {market} {start_date}-{end_date}")

        # 调用xtquant交易日历接口
        trading_dates = self.xt_data.get_trading_dates(
            market=market,
            start_time=start_date,
            end_time=end_date
        )

        if trading_dates is None or len(trading_dates) == 0:
            self.logger.warning(f"xtquant未获取到交易日历: {market}")
            return None

        # 转换为DataFrame
        result_df = pd.DataFrame({
            'trading_date': trading_dates
        })

        # 日期格式标准化 (xtquant返回的是字符串'YYYYMMDD')
        # 过滤掉无效日期(如'9200000')
        def parse_trading_date(date_str):
            try:
                # 只保留8位日期
                if len(str(date_str)) == 8:
                    return pd.to_datetime(date_str, format='%Y%m%d').strftime('%Y-%m-%d')
                else:
                    return None
            except:
                return None

        result_df['trading_date'] = result_df['trading_date'].apply(parse_trading_date)
        # 移除无效日期
        result_df = result_df[result_df['trading_date'].notna()].reset_index(drop=True)

        self.logger.info(f"xtquant获取交易日历成功: {market} {len(result_df)}个交易日")
        return result_df

    except Exception as e:
        self.logger.error(f"xtquant获取交易日历失败: {market} {e}")
        self.last_error = str(e)
        self.error_count += 1
        return None
get_stock_list
get_stock_list(market: Optional[str] = None) -> Optional[DataFrame]

获取股票列表(新增功能 - 阶段1)

参数:

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

市场代码(可选) - 'SH': 上交所 - 'SZ': 深交所 - None: 所有市场

None

返回:

名称 类型 描述
DataFrame Optional[DataFrame]

股票列表,包含code, name列,失败返回None

性能优势: - 本地数据,秒级响应 - 包含股票名称

源代码位于: adapters/xtquant_adapter.py
def get_stock_list(
    self,
    market: Optional[str] = None
) -> Optional[pd.DataFrame]:
    """
    获取股票列表(新增功能 - 阶段1)

    Args:
        market: 市场代码(可选)
            - 'SH': 上交所
            - 'SZ': 深交所
            - None: 所有市场

    Returns:
        DataFrame: 股票列表,包含code, name列,失败返回None

    性能优势:
    - 本地数据,秒级响应
    - 包含股票名称
    """
    if not self.is_connected:
        if not self.connect():
            return None

    try:
        self.logger.debug(f"xtquant获取股票列表: {market or '全市场'}")

        # 获取股票详情(包含名称)
        if market:
            # 指定市场
            stock_list = self.xt_data.get_stock_list_in_sector(f'{market}A股')
        else:
            # 全市场
            stock_list = self.xt_data.get_stock_list_in_sector('沪深A股')

        if stock_list is None or len(stock_list) == 0:
            self.logger.warning(f"xtquant未获取到股票列表: {market or '全市场'}")
            return None

        # 获取股票名称
        result_list = []
        for xt_code in stock_list:
            try:
                # 获取股票详情
                detail = self.xt_data.get_instrument_detail(xt_code)
                if detail:
                    stock_name = detail.get('InstrumentName', '')
                    # 标准化代码(去除.SH/.SZ后缀)
                    code = xt_code.split('.')[0] if '.' in xt_code else xt_code
                    result_list.append({
                        'code': code,
                        'name': stock_name
                    })
            except:
                continue

        if len(result_list) == 0:
            self.logger.warning(f"xtquant股票列表为空")
            return None

        result_df = pd.DataFrame(result_list)

        self.logger.info(f"xtquant获取股票列表成功: {len(result_df)}只股票")
        return result_df

    except Exception as e:
        self.logger.error(f"xtquant获取股票列表失败: {e}")
        self.last_error = str(e)
        self.error_count += 1
        return None
get_adjust_factors
get_adjust_factors(code: str, start_date: Optional[str] = None, end_date: Optional[str] = None) -> Optional[DataFrame]

获取复权因子(新增功能 - 阶段1)

参数:

名称 类型 描述 默认
code str

股票代码

必需
start_date Optional[str]

开始日期 'YYYYMMDD'

None
end_date Optional[str]

结束日期 'YYYYMMDD'

None

返回:

名称 类型 描述
DataFrame Optional[DataFrame]

复权因子,包含date, adjust_factor列,失败返回None

用途: - 用于数据一致性校验 - 跨数据源复权因子对比 - 确保复权数据无波动

实现说明: - xtquant使用 get_divid_factors 获取除权除息信息 - 计算累积复权因子

源代码位于: adapters/xtquant_adapter.py
def get_adjust_factors(
    self,
    code: str,
    start_date: Optional[str] = None,
    end_date: Optional[str] = None
) -> Optional[pd.DataFrame]:
    """
    获取复权因子(新增功能 - 阶段1)

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

    Returns:
        DataFrame: 复权因子,包含date, adjust_factor列,失败返回None

    用途:
    - 用于数据一致性校验
    - 跨数据源复权因子对比
    - 确保复权数据无波动

    实现说明:
    - xtquant使用 get_divid_factors 获取除权除息信息
    - 计算累积复权因子
    """
    if not self.is_connected:
        if not self.connect():
            return None

    try:
        from datetime import datetime, timedelta

        xt_code = self._convert_code(code)

        # 默认获取最近3年数据
        if not start_date or not end_date:
            end_dt = datetime.now()
            start_dt = end_dt - timedelta(days=3*365)
            start_date = start_dt.strftime('%Y%m%d')
            end_date = end_dt.strftime('%Y%m%d')

        self.logger.debug(f"xtquant获取复权因子: {code} {start_date}-{end_date}")

        # 调用xtquant复权因子接口
        # 注意: get_divid_factors的参数名是stock_code(无其他参数)
        divid_factors = self.xt_data.get_divid_factors(xt_code)

        if divid_factors is None or divid_factors.empty:
            self.logger.warning(f"xtquant未获取到复权因子: {code}")
            return None

        # 数据标准化
        result_df = divid_factors.reset_index()

        # 转换日期格式
        if 'date' in result_df.columns:
            result_df['date'] = pd.to_datetime(result_df['date']).dt.strftime('%Y-%m-%d')

        # 重命名列(如果需要)
        if 'factor' in result_df.columns:
            result_df = result_df.rename(columns={'factor': 'adjust_factor'})

        self.logger.info(f"xtquant获取复权因子成功: {code} {len(result_df)}条")
        return result_df

    except Exception as e:
        self.logger.error(f"xtquant获取复权因子失败: {code} {e}")
        self.last_error = str(e)
        self.error_count += 1
        return None
get_stock_name
get_stock_name(code: str) -> Optional[str]

通过 xtdata 获取股票名称

参数:

名称 类型 描述 默认
code str

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

必需

返回:

类型 描述
Optional[str]

股票名称,失败返回 None

源代码位于: adapters/xtquant_adapter.py
def get_stock_name(self, code: str) -> Optional[str]:
    """
    通过 xtdata 获取股票名称

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

    Returns:
        股票名称,失败返回 None
    """
    if not self.is_connected:
        if not self.connect():
            return None

    try:
        # 标准化为 xtquant 格式
        xt_code = self._normalize_to_xt_code(code)
        detail = self.xt_data.get_instrument_detail(xt_code)

        if detail:
            name = detail.get('InstrumentName') or detail.get('instrumentName') or detail.get('name')
            if name:
                self.logger.debug(f"xtquant获取股票名称: {code} -> {name}")
                return name

        self.logger.debug(f"xtquant未获取到{code}的股票详情")
        return None

    except Exception as e:
        self.logger.debug(f"xtquant获取股票名称失败 {code}: {e}")
        return None