开篇导语
上一讲我们搭建了Python量化开发环境,并通过akshare完成了第一次"从数据到决策"的多因子选股实验。但这里有一个隐患你可能已经察觉了:每次运行策略都需要临时从网上抓数据——慢、不稳定、不可复现。
真正能打仗的量化团队,绝不会每跑一次回测就去网上重新下载一次数据。他们把数据采集、清洗、存储做成一条自动化管道(Pipeline),需要时从本地数据库秒级读取。这条数据管道就是量化交易的"弹药库"——没有稳定的数据供给,策略研究就是纸上谈兵。
今天,小晴(情报官)正式上岗。他将负责把分散在各个数据源的市场信息——行情、财务、宏观、研报、新闻——统一收拢、清洗、标准化后存入本地数据库,并在后续课程中逐步进阶为能阅读研报、感知舆情的AI投研专家。
学完本节后,你将能够:
- 搭建覆盖行情、财务、宏观三类数据的本地采集管道
- 处理复权、停牌、缺失值等常见数据质量问题
- 利用大模型从文本中提取结构化"催化剂事件"
- 建立本地HDF5/CSV格式的高质量数据库
4.1 小晴的使命——数据管道的全景设计
概念讲解
先看一个真实场景:你要回测一个基于财务因子(PE、ROE)和技术指标(MACD、RSI)的组合策略。这需要两类数据:日线行情(用来算技术指标)和季度财务报告(用来算估值因子)。
如果每次回测前都从网上下载——akshare取行情(约10秒)、取财务数据(约30秒)、取宏观指标(约15秒)——光数据准备就要一分钟。一天迭代测试50个参数组合,等待时间就接近一小时了。更致命的是,网络波动可能导致某次下载失败,你的批量回测脚本就会卡在半路,一行一行排查错误。
小晴的数据管道要解决的就是这个问题。它的设计架构如下:
┌─────────┐ ┌─────────┐ ┌─────────┐
│ 数据源A │ │ 数据源B │ │ 数据源C │
│ akshare │ │ tushare │ │ QMT行情 │
└────┬─────┘ └────┬─────┘ └────┬─────┘
│ │ │
└──────────────┼──────────────┘
▼
┌────────────────┐
│ 小晴采集器 │ ← 定时/手动触发
│ 多源并行拉取 │
└───────┬────────┘
▼
┌────────────────┐
│ 清洗与标准化 │ ← 复权处理、停牌标记、去重
│ 缺失值填补 │
└───────┬────────┘
▼
┌────────────────┐
│ 本地数据库 │ ← HDF5 / CSV / SQLite
│ (结构化存储) │
└───────┬────────┘
▼
┌────────────────┐
│ 小雅策略回测 │ ← 秒级读取,无需联网
└────────────────┘
核心设计原则有三条:
代码演示
下面是小晴数据管道的核心采集器框架:
"""
小晴数据管道 – 核心采集框架
职责:统一管理数据采集、增量更新、本地存储
"""
import akshare as ak
import pandas as pd
import os
from datetime import datetime, timedelta
class XiaoQingDataPipeline:
"""小晴情报官的数据管道"""
def __init__(self, data_dir="./local_data"):
"""初始化:设置本地存储目录,创建子文件夹"""
self.data_dir = data_dir
self.daily_dir = os.path.join(data_dir, "daily") # 日线行情
self.finance_dir = os.path.join(data_dir, "finance") # 财务数据
self.macro_dir = os.path.join(data_dir, "macro") # 宏观数据
for d in [self.daily_dir, self.finance_dir, self.macro_dir]:
os.makedirs(d, exist_ok=True)
# — 模块一:日线行情采集 —
def fetch_daily_kline(self, symbol, start_date="20200101",
end_date=None, force_refresh=False):
"""
采集单只股票的日K线数据,支持增量更新
参数:
symbol: 股票代码,如 '600519'
start_date: 起始日期
end_date: 结束日期,默认为今天
force_refresh: 是否强制全量更新
"""
if end_date is None:
end_date = datetime.now().strftime("%Y%m%d")
file_path = os.path.join(self.daily_dir, f"{symbol}.csv")
# 增量更新逻辑:如果本地已有数据,只拉取新增部分
if not force_refresh and os.path.exists(file_path):
existing = pd.read_csv(file_path, index_col=0, parse_dates=True)
last_date = existing.index.max().strftime("%Y%m%d")
if last_date >= end_date:
print(f"[{symbol}] 数据已是最新,跳过")
return existing
# 从最后日期+1天开始拉取
start_date = (pd.to_datetime(last_date) +
timedelta(days=1)).strftime("%Y%m%d")
print(f"[{symbol}] 增量更新:{start_date} → {end_date}")
# 从akshare拉取数据
df = ak.stock_zh_a_hist(symbol=symbol, period="daily",
start_date=start_date, end_date=end_date,
adjust="qfq")
df['日期'] = pd.to_datetime(df['日期'])
df.set_index('日期', inplace=True)
# 如果有历史数据,合并
if os.path.exists(file_path) and not force_refresh:
existing = pd.read_csv(file_path, index_col=0, parse_dates=True)
df = pd.concat([existing, df])
df = df[~df.index.duplicated(keep='last')] # 去重,保留最新
# 存入本地CSV
df.to_csv(file_path)
print(f"[{symbol}] 保存完成,共 {len(df)} 条记录")
return df
# — 模块二:宏观数据采集 —
def fetch_macro_data(self):
"""采集关键宏观经济指标"""
indicators = {}
# 获取货币供应量数据(M1/M2是衡量市场流动性的关键指标)
try:
money_supply = ak.macro_china_money_supply()
indicators['money_supply'] = money_supply
print("宏观数据:货币供应量采集完成")
except Exception as e:
print(f"宏观数据采集失败(货币供应量):{e}")
# 获取CPI数据(通货膨胀率指标)
try:
cpi = ak.macro_china_cpi_monthly()
indicators['cpi'] = cpi
print("宏观数据:CPI采集完成")
except Exception as e:
print(f"宏观数据采集失败(CPI):{e}")
# 保存到本地
for name, data in indicators.items():
file_path = os.path.join(self.macro_dir, f"{name}.csv")
data.to_csv(file_path, index=False)
return indicators
# — 模块三:财务数据采集 —
def fetch_financial_data(self, symbol, indicator="按报告期"):
"""
采集个股财务指标(以同花顺接口为例)
返回值包含:每股收益、每股净资产、净资产收益率、毛利率等
"""
try:
df = ak.stock_financial_abstract_ths(symbol=symbol,
indicator=indicator)
file_path = os.path.join(self.finance_dir, f"{symbol}_finance.csv")
df.to_csv(file_path, index=False)
print(f"[{symbol}] 财务数据采集完成,共 {len(df)} 条")
return df
except Exception as e:
print(f"[{symbol}] 财务数据采集失败:{e}")
return None
# === 使用示例:让小晴跑一次全量采集 ===
xiaoQing = XiaoQingDataPipeline(data_dir="./local_data")
# 采集股票池的日线数据
stock_pool = ['600519', '000001', '000858', '600036']
for code in stock_pool:
xiaoQing.fetch_daily_kline(code, start_date="20240101")
关键要点
- 增量更新机制是数据管道的生命线——每次只拉新增数据,日线一次增量不到1KB
- HDF5格式比CSV更适合大规模行情存储(单文件可存数千只股票的全部历史)
- 采集频率要合理:宏观数据按月更新,日线行情按天更新,tick级数据慎采(一天就是几个GB)
- 小晴的代码应该模块化——每个数据源独立一个函数,方便后续扩展和维护
4.2 数据清洗——修复"脏数据"的常见套路
概念讲解
从外部数据源获取的原始数据几乎100%存在质量问题。以下是量化数据中最常见的四类问题及其处理方案:
| 复权缺失 | 历史价格在除权除息日出现"断崖式"跳空 | 统一使用前复权(qfq),避免含权价格影响回测 |
| 停牌空值 | 交易日无成交,OHLC为NaN或重复前一日 | 标记停牌期,回测中跳过或填充前收盘价 |
| 财务数据滞后 | 12月31日的年报,实际次年4月才公布 | 使用时点匹配:4月之前用三季度报,4月之后可用年报 |
| 异常极值 | 错误录入导致PE=99999或ROE=500% | 设定合理阈值,超过阈值的标记为缺失并另行处理 |
其中财务数据的"时点匹配"是最容易被忽视但影响最大的问题。例如,2024年1月做回测,你"看到"的是2023年年报的数据(通常在3-4月才公布),这属于"偷看答案"——实盘中1月份根本不知道这个数据。
代码演示
下面实现一个数据清洗器,自动处理上述四类问题:
class DataCleaner:
"""小晴的清洗工具箱"""
@staticmethod
def clean_daily_kline(df: pd.DataFrame) -> pd.DataFrame:
"""
清洗日线数据:处理停牌、缺失值、异常价格
参数:
df: 日线DataFrame,index为日期,包含开盘/收盘/最高/最低/成交量
"""
cleaned = df.copy()
# 1. 标记停牌日:成交量=0 或 最高价=最低价(一字停牌)
cleaned['is_suspended'] = (
(cleaned['成交量'] == 0) |
(cleaned['最高'] == cleaned['最低'])
)
# 2. 前向填充停牌日的收盘价(用于计算收益率时保持连续性)
cleaned['close_ffill'] = cleaned['收盘'].replace(
cleaned[cleaned['is_suspended']]['收盘'], pd.NA
).ffill()
# 3. 检测异常价格变动:单日涨跌幅超过涨跌停限制
cleaned['pct_change'] = cleaned['收盘'].pct_change()
cleaned['is_abnormal'] = abs(cleaned['pct_change']) > 0.11 # 主板±10%阈值
if cleaned['is_abnormal'].sum() > 0:
abnormal_dates = cleaned[cleaned['is_abnormal']].index
print(f"检测到 {len(abnormal_dates)} 个异常变动日:")
for d in abnormal_dates:
pct = cleaned.loc[d, 'pct_change']
print(f" {d.strftime('%Y-%m-%d')}: "
f"{pct:+.2%}(可能为复权或数据错误)")
return cleaned
@staticmethod
def filter_financial_extremes(df: pd.DataFrame,
columns: list) -> pd.DataFrame:
"""
过滤财务数据中的极端值
参数:
df: 财务数据DataFrame
columns: 需要过滤的列名列表(如['pe', 'pb'])
"""
for col in columns:
if col not in df.columns:
continue
# 计算1%和99%分位数
lower_bound = df[col].quantile(0.01)
upper_bound = df[col].quantile(0.99)
# 将超出范围的标记为NaN
mask = (df[col] < lower_bound) | (df[col] > upper_bound)
df.loc[mask, col] = pd.NA
extreme_count = mask.sum()
if extreme_count > 0:
print(f" [{col}] 剔除 {extreme_count} 个极端值 "
f"(范围:{lower_bound:.2f} ~ {upper_bound:.2f})")
return df
@staticmethod
def align_report_date(df: pd.DataFrame,
report_col: str = '报告期') -> pd.DataFrame:
"""
财务数据时点对齐:将报告期映射到实际可用的最早日期
A股年报截止日为4月30日,中报为8月31日,季报为当季末次月
例如:2023-12-31的年报,actual_date = 2024-04-30
"""
df = df.copy()
df[report_col] = pd.to_datetime(df[report_col])
def get_available_date(report_date):
month = report_date.month
year = report_date.year
if month == 12: # 年报:次年4月30日可用
return pd.Timestamp(year + 1, 4, 30)
elif month == 6: # 中报:当年8月31日可用
return pd.Timestamp(year, 8, 31)
elif month == 9: # 三季报:当年10月31日可用
return pd.Timestamp(year, 10, 31)
elif month == 3: # 一季报:当年4月30日可用
return pd.Timestamp(year, 4, 30)
return report_date
df['available_date'] = df[report_col].apply(get_available_date)
return df
关键要点
- 前复权是A股回测的"默认选项"——它保证除权除息不会在K线图上留"缺口"
- 停牌期间的数据处理影响回测结果——如果策略在停牌期间产生信号但无法执行,需要标记为"无效信号"
- 财务数据时点匹配严重且普遍被忽视——很多策略回测净值漂亮,实盘惨不忍睹,原因之一就是用了"未来财务数据"
- 清洗规则不能一刀切——科创板涨跌停±20%,需要动态调整异常波动的检测阈值
4.3 催化剂识别——用大模型读懂市场事件
概念讲解
除了结构化数据(价格、财务),市场还存在大量非结构化信息——新闻公告、政策文件、行业研报、社交媒体讨论。这些信息中蕴藏着可能的"催化剂事件":
- 超预期财报(净利润大幅超出分析师一致预期)
- 政策利好(如新能源补贴、半导体减税)
- 重大合同公告(如中标大额项目)
- 负面事件(如被立案调查、高管出事、产品召回)
传统方法用关键词匹配来识别事件(如"同比增长 + 净利润 > 100%"),但会漏掉很多"不像关键词但确实是事件"的情况。大模型在这一步大有用武之地——它能理解句子的语义,而不只是扫描关键词。
小晴的工作流程是:
代码演示
下面是一个利用OpenAI API(或其他兼容接口)做催化剂事件识别的方法:
"""
小晴催化剂事件识别模块
—— 让大模型帮你读公告,提取结构化事件信息
"""
import json
class CatalystDetector:
"""催化剂事件识别器"""
# 定义提示词模板
DETECTION_PROMPT = """你是一个A股市场事件分析专家。请阅读以下新闻/公告,判断它是否是一个值得量化策略关注的"催化剂事件"。
输出要求:严格返回JSON格式,包含以下字段:
– is_catalyst: true/false,是否构成催化剂
– event_type: 事件类型(超预期业绩/政策利好/重大合同/行业变革/风险警示/其他)
– direction: 对股价的可能影响(positive/negative/neutral)
– impact_level: 影响程度(high/medium/low)
– summary: 一句话总结事件(不超过30字)
新闻内容:
{news_text}
请仅返回JSON,不要有其他文字。"""
def __init__(self, api_key=None, base_url=None):
"""初始化:配置大模型API(使用OpenAI兼容接口)"""
self.api_key = api_key or os.getenv("OPENAI_API_KEY")
self.base_url = base_url or os.getenv("OPENAI_BASE_URL")
def detect(self, news_text: str) -> dict:
"""
识别单条新闻是否包含催化剂事件
在实际部署中,这个函数会被批量调用,
处理每天新出现的数百条公告和新闻
"""
try:
# 使用openai库(pip install openai)
from openai import OpenAI
client = OpenAI(api_key=self.api_key, base_url=self.base_url)
response = client.chat.completions.create(
model="gpt-4o-mini", # 使用轻量模型降低API成本
messages=[
{"role": "system", "content": "你是一个A股事件分析助手,请严格按JSON格式回复。"},
{"role": "user", "content": self.DETECTION_PROMPT.format(
news_text=news_text
)}
],
temperature=0.1, # 低温度保证输出格式稳定
max_tokens=200
)
result = json.loads(response.choices[0].message.content)
return result
except Exception as e:
print(f"催化剂识别失败:{e}")
return {"is_catalyst": False, "error": str(e)}
def batch_detect(self, news_list: list) -> list:
"""批量识别,返回所有被标记为催化剂的事件"""
catalysts = []
for i, news in enumerate(news_list):
result = self.detect(news)
if result.get('is_catalyst'):
catalysts.append({
**result,
'original_news': news,
'detected_at': datetime.now().isoformat()
})
if (i + 1) % 50 == 0:
print(f"已处理 {i+1}/{len(news_list)} 条新闻…")
print(f"共识别出 {len(catalysts)} 条催化剂事件")
return catalysts
# === 使用示例 ===
# detector = CatalystDetector(api_key="your-api-key")
# test_news = "贵州茅台2024年第三季度净利润同比增长25%,超出市场预期的18%"
# result = detector.detect(test_news)
# print(json.dumps(result, ensure_ascii=False, indent=2))
print("催化剂识别模块加载完成。实际使用需配置大模型API密钥。")
关键要点
- 大模型识别事件的准确率远高于关键词匹配(约85% vs 60%),但API调用有成本,需要合理控制每日调用量
- temperature设为0.1是为了保证JSON格式输出可解析,避免模型"发挥创意"破坏结构化数据
- 催化剂事件是"加分项"而非"必需项"——即使不用大模型,纯结构化数据的策略也能盈利
- 事件影响的时间衰减很快——一条利好公告的市场效应通常在3-5个交易日内消化完毕
4.4 本地数据库——让所有数据触手可及
概念讲解
前面三节分别采集了行情数据、财务数据和事件数据。现在需要把它们统一存储到一个结构化的本地数据库中。推荐使用两种格式的组合:
| HDF5 | 大规模行情数据(数千只股票 × 十年日线) | 读写极快,支持压缩,单文件存储 | 并发写入能力弱 |
| SQLite | 财务数据、事件数据(字段结构固定) | 支持SQL查询,多表关联 | 写入速度不如HDF5 |
| Parquet | 中间计算结果、因子值 | 列式存储,压缩率高 | 不支持原地更新 |
代码演示
下面构建一个统一的本地数据访问层:
import sqlite3
import pandas as pd
class LocalDataWarehouse:
"""小晴的本地数据仓库——所有数据源的统一访问入口"""
def __init__(self, base_dir="./local_data"):
self.base_dir = base_dir
# 初始化SQLite数据库
self.db_path = os.path.join(base_dir, "quant_db.sqlite")
self._init_db()
def _init_db(self):
"""创建数据库表结构"""
conn = sqlite3.connect(self.db_path)
cursor = conn.cursor()
# 事件表:存储识别出的催化剂事件
cursor.execute("""
CREATE TABLE IF NOT EXISTS catalyst_events (
id INTEGER PRIMARY KEY AUTOINCREMENT,
stock_code TEXT NOT NULL, — 股票代码
event_type TEXT, — 事件类型
direction TEXT, — 利好/利空
impact_level TEXT, — 影响程度
summary TEXT, — 事件摘要
detected_at TEXT, — 检测时间
source_url TEXT — 来源链接
)
""")
# 交易日历表:用于快速判断某天是否为交易日
cursor.execute("""
CREATE TABLE IF NOT EXISTS trading_calendar (
trade_date TEXT PRIMARY KEY,
is_trading_day INTEGER DEFAULT 1
)
""")
conn.commit()
conn.close()
print(f"数据仓库初始化完成:{self.db_path}")
def insert_event(self, stock_code: str, event_data: dict):
"""写入一条催化剂事件"""
conn = sqlite3.connect(self.db_path)
conn.execute("""
INSERT INTO catalyst_events
(stock_code, event_type, direction, impact_level, summary, detected_at)
VALUES (?, ?, ?, ?, ?, ?)
""", (
stock_code,
event_data.get('event_type'),
event_data.get('direction'),
event_data.get('impact_level'),
event_data.get('summary'),
datetime.now().isoformat()
))
conn.commit()
conn.close()
def query_events(self, stock_code=None, days=7):
"""查询最近的催化剂事件"""
cutoff = (datetime.now() – timedelta(days=days)).isoformat()
conn = sqlite3.connect(self.db_path)
query = "SELECT * FROM catalyst_events WHERE detected_at >= ?"
params = [cutoff]
if stock_code:
query += " AND stock_code = ?"
params.append(stock_code)
df = pd.read_sql_query(query, conn, params=params)
conn.close()
return df
def get_daily_data(self, stock_code):
"""从CSV快速读取日线数据"""
file_path = os.path.join(self.base_dir, "daily", f"{stock_code}.csv")
if not os.path.exists(file_path):
raise FileNotFoundError(f"日线数据不存在:{file_path},请先运行采集器")
return pd.read_csv(file_path, index_col=0, parse_dates=True)
关键要点
- HDF5适合"一次写入、多次读取"的场景,写入时注意加锁避免并发冲突
- SQLite虽然轻量,但单文件数据库模式非常适合个人量化团队的规模(数百只股票、数年数据)
- 事件表要记录"检测时间"而非"事件发生时间"——同一事件可能在小晴和多处来源同时被捕获
- 数据备份是必须的——本地数据库不是云存储,硬盘故障会丢失所有历史数据
实战总结
本节小晴(情报官)完成了他的第一次"入职任务"——搭建从数据采集到本地存储的完整管道:
常见误区提醒:
- 误区一:数据采集频率越高越好。日线级别策略用日线数据足够了,tick级数据只会徒增存储和计算负担
- 误区二:所有数据都用CSV存。百只股票以上的日线数据用HDF5是必须的,CSV的I/O会成为瓶颈
- 误区三:忽略财务数据的时点对齐。这是导致回测与实际偏差的最隐蔽原因——"偷看未来数据"
下一讲,小雅(分析师)将接过小晴准备好的数据,正式登场——用Backtrader回测引擎跑通第一个双均线策略,并学会解读收益、回撤、夏普比率等核心绩效指标。
下一篇预告
第05讲:小雅初试锋芒——Backtrader回测引擎入门实战:数据和环境已就绪,是时候让策略接受历史的检验了。小雅将带着双均线策略走进Backtrader的回测战场,用真实数据回答"这个策略到底能不能赚钱"。
关注专栏不错过后续内容。
网硕互联帮助中心





评论前必须登录!
注册