Python 金融数据处理:Wind/聚源数据接入与标准化处理
Python 金融数据处理:Wind/聚源数据接入与标准化处理
一、同一只股票,Wind 和聚源返回的 PE 不一样——数据处理的最大坑
金融数据处理的难点不是"能不能拿到数据",而是"不同数据源的数据格式、口径、时效性各不相同"。
以市盈率 PE 为例:
- Wind 返回的 PE 是静态市盈率(TTM),用 trailing 12 months 的净利润计算
- 聚源返回的 PE 可能是动态市盈率(Forward),用预测净利润计算
- Bloomberg 有自己的一套计算口径
更麻烦的是,同一数据源的不同版本 API 返回格式也可能不同。如果没有统一的数据标准化层,下游的分析模型会一直吃进"看起来很对但实际口径不同的数据"。
二、金融数据标准化架构
三、Python 实现数据适配与标准化
统一数据模型
from dataclasses import dataclass, field from datetime import datetime, date from typing import Optional, Dict, List, Any from enum import Enum import pandas as pd class DataSource(Enum): WIND = "wind" JOINQUANT = "joinquant" TUSHARE = "tushare" BLOOMBERG = "bloomberg" MANUAL = "manual" class DataField(Enum): """统一字段定义——所有数据源都映射到这个枚举""" # 基础信息 STOCK_CODE = "stock_code" STOCK_NAME = "stock_name" LIST_DATE = "list_date" # 行情 OPEN = "open" HIGH = "high" LOW = "low" CLOSE = "close" VOLUME = "volume" AMOUNT = "amount" # 估值 PE_TTM = "pe_ttm" # 市盈率(TTM) PE_FORWARD = "pe_forward" # 市盈率(预测) PB = "pb" # 市净率 PS_TTM = "ps_ttm" # 市销率(TTM) # 财务 REVENUE = "revenue" # 营业收入 NET_PROFIT = "net_profit" # 净利润 TOTAL_ASSETS = "total_assets" TOTAL_LIABILITIES = "total_liabilities" # 现金流 OPERATING_CF = "operating_cf" FREE_CF = "free_cf" @dataclass class DataFieldMeta: """字段元数据——记录数据口径和来源""" field: DataField source: DataSource source_field: str # 原始字段名 calculation_method: str # 计算口径说明 fetch_time: datetime # 数据获取时间 data_date: date # 数据对应日期 @dataclass class StandardizedDataFrame: """标准化后的数据——包含数据 + 元数据""" df: pd.DataFrame meta: Dict[str, DataFieldMeta] # 字段名 → 元数据 source: DataSource fetch_time: datetime数据源适配器
from abc import ABC, abstractmethod class DataAdapter(ABC): """数据适配器抽象基类""" @abstractmethod def get_source(self) -> DataSource: """返回数据源标识""" pass @abstractmethod def fetch_daily(self, stock_codes: List[str], start_date: str, end_date: str, fields: List[DataField]) -> StandardizedDataFrame: """拉取日频数据""" pass @abstractmethod def fetch_financial(self, stock_codes: List[str], report_periods: List[str], fields: List[DataField]) -> StandardizedDataFrame: """拉取财务数据""" pass # 字段映射表:DataSource → Unified Field @property @abstractmethod def field_mapping(self) -> Dict[DataField, str]: """原始字段名 → 统一字段名的映射""" pass class WindAdapter(DataAdapter): """Wind 数据适配器""" def __init__(self): try: from WindPy import w self.w = w self.w.start() except ImportError: raise ImportError("请安装 WindPy: pip install WindPy") def get_source(self) -> DataSource: return DataSource.WIND @property def field_mapping(self) -> Dict[DataField, str]: return { DataField.OPEN: "open", DataField.HIGH: "high", DataField.LOW: "low", DataField.CLOSE: "close", DataField.VOLUME: "volume", DataField.AMOUNT: "amt", DataField.PE_TTM: "pe_ttm", DataField.PB: "pb", DataField.REVENUE: "or_yoy", # 营业收入同比增长率 DataField.NET_PROFIT: "profit_yoy", # 净利润同比增长率 } def fetch_daily(self, stock_codes, start_date, end_date, fields): """Wind 日频数据拉取""" # 转换字段 wind_fields = [self.field_mapping[f] for f in fields if f in self.field_mapping] # 调用 Wind API field_str = ",".join(wind_fields) code_str = ",".join(stock_codes) try: # WindPy 调用 err, data = self.w.wsd(code_str, field_str, start_date, end_date, "") if err != 0: raise RuntimeError(f"Wind 数据拉取失败, error_code={err}") # 构造 DataFrame dates = pd.to_datetime(data.Times) df = pd.DataFrame(index=dates) for i, code in enumerate(stock_codes): for j, field in enumerate(fields): if field in self.field_mapping: col_name = f"{code}_{field.value}" df[col_name] = data.Data[j * len(stock_codes) + i] return StandardizedDataFrame( df=df, meta=self._build_meta(fields), source=DataSource.WIND, fetch_time=datetime.now(), ) except Exception as e: print(f"Wind API 调用失败: {e}") raise class JoinQuantAdapter(DataAdapter): """聚源/JoinQuant 数据适配器""" def __init__(self): try: import jqdatasdk as jq self.jq = jq # 登录(需要提前配置用户名密码) except ImportError: raise ImportError("请安装 jqdatasdk") def get_source(self) -> DataSource: return DataSource.JOINQUANT @property def field_mapping(self) -> Dict[DataField, str]: # 注意:聚源的字段名和 Wind 不同! # 这就是为什么需要适配器——统一字段对外 return { DataField.OPEN: "open", DataField.CLOSE: "close", DataField.VOLUME: "volume", DataField.PE_TTM: "pe_ratio", # Wind: pe_ttm, 聚源: pe_ratio DataField.PB: "pb_ratio", # Wind: pb, 聚源: pb_ratio DataField.NET_PROFIT: "net_profit_margin", # 口径也可能不同! }数据标准化处理
class DataStandardizer: """数据标准化器——将多源数据统一为单一格式""" def __init__(self): self.adapters: Dict[DataSource, DataAdapter] = {} self.cache_dir = "./data_cache" def register_adapter(self, adapter: DataAdapter): """注册数据源适配器""" self.adapters[adapter.get_source()] = adapter def get_unified_data( self, stock_codes: List[str], start_date: str, end_date: str, fields: List[DataField], primary_source: DataSource = DataSource.WIND, fallback_sources: List[DataSource] = None, ) -> StandardizedDataFrame: """ 获取统一数据——优先从主数据源获取,失败则降级到备用源 """ # 尝试主数据源 if primary_source in self.adapters: try: print(f"从 {primary_source.value} 拉取数据...") return self.adapters[primary_source].fetch_daily( stock_codes, start_date, end_date, fields, ) except Exception as e: print(f"主数据源 {primary_source.value} 失败: {e}") # 降级到备用数据源 for source in (fallback_sources or []): if source in self.adapters: try: print(f"降级到 {source.value} 拉取数据...") return self.adapters[source].fetch_daily( stock_codes, start_date, end_date, fields, ) except Exception as e: print(f"备用源 {source.value} 也失败了: {e}") raise RuntimeError("所有数据源均不可用") def validate_data(self, data: StandardizedDataFrame) -> Dict[str, List[str]]: """ 数据校验——检查完整性和合理性 """ issues = {} df = data.df # 1. 缺失值检查 missing_cols = df.columns[df.isnull().any()].tolist() if missing_cols: issues["missing_values"] = missing_cols # 2. 异常值检查(3σ 规则) for col in df.select_dtypes(include=['float64', 'int64']).columns: mean = df[col].mean() std = df[col].std() if std > 0: outliers = df[abs(df[col] - mean) > 3 * std] if len(outliers) > 0: issues[f"outliers_{col}"] = [ f"{date.strftime('%Y-%m-%d')}: {val:.2f}" for date, val in zip(outliers.index, outliers[col]) ] # 3. 逻辑校验(例如:最高价 >= 最低价) for code in set(col.split('_')[0] for col in df.columns if '_high' in col): high_col = f"{code}_high" low_col = f"{code}_low" if high_col in df.columns and low_col in df.columns: invalid = df[df[high_col] < df[low_col]] if len(invalid) > 0: issues[f"logic_error_{code}"] = [ f"{d.strftime('%Y-%m-%d')}: high < low" for d in invalid.index ] return issues def cache_to_parquet(self, data: StandardizedDataFrame, filename: str): """缓存到本地 Parquet 格式(高效压缩存储)""" import os import pyarrow as pa import pyarrow.parquet as pq os.makedirs(self.cache_dir, exist_ok=True) filepath = os.path.join(self.cache_dir, filename) table = pa.Table.from_pandas(data.df) pq.write_table( table, filepath, compression='snappy', # 快速压缩 ) print(f"数据已缓存: {filepath}")四、边界分析与 Trade-offs
数据口径的对齐:
- 不同数据源对同一指标的计算方式可能不同
- 必须在元数据中记录计算口径,供下游分析模型参考
- 不能假设"PE 都是 PE"
复权处理:
- 股票行情数据需要统一复权方式(前复权 / 后复权)
- Wind 和聚源的复权方式可能不同
- 建议在适配器层统一为前复权
数据时效性:
- Wind 和聚源的数据更新频率不同(T+0 vs T+1)
- 需要在元数据中标记数据获取时间和数据对应日期
- 回测时要注意"未来数据"问题
本地缓存策略:
- 金融数据拉取受限速和配额限制
- 建议缓存已拉取的数据(Parquet 格式,按月分文件)
- 增量更新而非全量重拉
五、总结
金融数据处理的核心不是"能拿到数据",而是"拿到的是正确的数据":
- 适配器模式——每个数据源一个 Adapter,屏蔽 API 差异
- 统一字段模型——DataField 枚举定义所有统一字段
- 元数据追溯——每个字段记录来源、口径、获取时间
- 多源降级——主源失败时自动切换到备用源
- 数据校验——缺失值、异常值、逻辑错误的自动化检查
金融数据处理的 80% 工作不在代码,在"搞清楚每个字段的口径是什么"。
