Python 金融数据处理:Wind/聚源数据接入与标准化处理

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(fWind 数据拉取失败, error_code{err}) # 构造 DataFrame dates pd.to_datetime(data.Times) df pd.DataFrame(indexdates) 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( dfdf, metaself._build_meta(fields), sourceDataSource.WIND, fetch_timedatetime.now(), ) except Exception as e: print(fWind 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[foutliers_{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[flogic_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_okTrue) filepath os.path.join(self.cache_dir, filename) table pa.Table.from_pandas(data.df) pq.write_table( table, filepath, compressionsnappy, # 快速压缩 ) print(f数据已缓存: {filepath})四、边界分析与 Trade-offs数据口径的对齐不同数据源对同一指标的计算方式可能不同必须在元数据中记录计算口径供下游分析模型参考不能假设PE 都是 PE复权处理股票行情数据需要统一复权方式前复权 / 后复权Wind 和聚源的复权方式可能不同建议在适配器层统一为前复权数据时效性Wind 和聚源的数据更新频率不同T0 vs T1需要在元数据中标记数据获取时间和数据对应日期回测时要注意未来数据问题本地缓存策略金融数据拉取受限速和配额限制建议缓存已拉取的数据Parquet 格式按月分文件增量更新而非全量重拉五、总结金融数据处理的核心不是能拿到数据而是拿到的是正确的数据适配器模式——每个数据源一个 Adapter屏蔽 API 差异统一字段模型——DataField 枚举定义所有统一字段元数据追溯——每个字段记录来源、口径、获取时间多源降级——主源失败时自动切换到备用源数据校验——缺失值、异常值、逻辑错误的自动化检查金融数据处理的 80% 工作不在代码在搞清楚每个字段的口径是什么。