【AI量化交易实战】第04讲:小晴(情报官)上岗——搭建本地数据管道与清洗流水线

📅 2026/7/27 19:10:50
【AI量化交易实战】第04讲:小晴(情报官)上岗——搭建本地数据管道与清洗流水线
开篇导语上一讲我们搭建了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 │ (结构化存储) │ └───────┬────────┘ ▼ ┌────────────────┐ │ 小雅策略回测 │ ← 秒级读取无需联网 └────────────────┘核心设计原则有三条一次采集多次使用数据更新只做增量不重复拉历史数据存储格式高效HDF5比CSV读取快5-10倍适合大数据量的行情存储原始数据保留清洗过程中始终保留原始下载文件方便回溯审计代码演示下面是小晴数据管道的核心采集器框架 小晴数据管道 - 核心采集框架 职责统一管理数据采集、增量更新、本地存储 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_okTrue) # --- 模块一日线行情采集 --- def fetch_daily_kline(self, symbol, start_date20200101, end_dateNone, force_refreshFalse): 采集单只股票的日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_col0, parse_datesTrue) 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(days1)).strftime(%Y%m%d) print(f[{symbol}] 增量更新{start_date} → {end_date}) # 从akshare拉取数据 df ak.stock_zh_a_hist(symbolsymbol, perioddaily, start_datestart_date, end_dateend_date, adjustqfq) df[日期] pd.to_datetime(df[日期]) df.set_index(日期, inplaceTrue) # 如果有历史数据合并 if os.path.exists(file_path) and not force_refresh: existing pd.read_csv(file_path, index_col0, parse_datesTrue) df pd.concat([existing, df]) df df[~df.index.duplicated(keeplast)] # 去重保留最新 # 存入本地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, indexFalse) return indicators # --- 模块三财务数据采集 --- def fetch_financial_data(self, symbol, indicator按报告期): 采集个股财务指标以同花顺接口为例 返回值包含每股收益、每股净资产、净资产收益率、毛利率等 try: df ak.stock_financial_abstract_ths(symbolsymbol, indicatorindicator) file_path os.path.join(self.finance_dir, f{symbol}_finance.csv) df.to_csv(file_path, indexFalse) 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_date20240101)关键要点增量更新机制是数据管道的生命线——每次只拉新增数据日线一次增量不到1KBHDF5格式比CSV更适合大规模行情存储单文件可存数千只股票的全部历史采集频率要合理宏观数据按月更新日线行情按天更新tick级数据慎采一天就是几个GB小晴的代码应该模块化——每个数据源独立一个函数方便后续扩展和维护4.2 数据清洗——修复脏数据的常见套路概念讲解从外部数据源获取的原始数据几乎100%存在质量问题。以下是量化数据中最常见的四类问题及其处理方案问题类型表现处理方案复权缺失历史价格在除权除息日出现断崖式跳空统一使用前复权qfq避免含权价格影响回测停牌空值交易日无成交OHLC为NaN或重复前一日标记停牌期回测中跳过或填充前收盘价财务数据滞后12月31日的年报实际次年4月才公布使用时点匹配4月之前用三季度报4月之后可用年报异常极值错误录入导致PE99999或ROE500%设定合理阈值超过阈值的标记为缺失并另行处理其中财务数据的时点匹配是最容易被忽视但影响最大的问题。例如2024年1月做回测你看到的是2023年年报的数据通常在3-4月才公布这属于偷看答案——实盘中1月份根本不知道这个数据。代码演示下面实现一个数据清洗器自动处理上述四类问题class DataCleaner: 小晴的清洗工具箱 staticmethod def clean_daily_kline(df: pd.DataFrame) - pd.DataFrame: 清洗日线数据处理停牌、缺失值、异常价格 参数 df: 日线DataFrameindex为日期包含开盘/收盘/最高/最低/成交量 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_keyNone, base_urlNone): 初始化配置大模型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_keyself.api_key, base_urlself.base_url) response client.chat.completions.create( modelgpt-4o-mini, # 使用轻量模型降低API成本 messages[ {role: system, content: 你是一个A股事件分析助手请严格按JSON格式回复。}, {role: user, content: self.DETECTION_PROMPT.format( news_textnews_text )} ], temperature0.1, # 低温度保证输出格式稳定 max_tokens200 ) 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已处理 {i1}/{len(news_list)} 条新闻...) print(f共识别出 {len(catalysts)} 条催化剂事件) return catalysts # 使用示例 # detector CatalystDetector(api_keyyour-api-key) # test_news 贵州茅台2024年第三季度净利润同比增长25%超出市场预期的18% # result detector.detect(test_news) # print(json.dumps(result, ensure_asciiFalse, indent2)) print(催化剂识别模块加载完成。实际使用需配置大模型API密钥。)关键要点大模型识别事件的准确率远高于关键词匹配约85% vs 60%但API调用有成本需要合理控制每日调用量temperature设为0.1是为了保证JSON格式输出可解析避免模型发挥创意破坏结构化数据催化剂事件是加分项而非必需项——即使不用大模型纯结构化数据的策略也能盈利事件影响的时间衰减很快——一条利好公告的市场效应通常在3-5个交易日内消化完毕4.4 本地数据库——让所有数据触手可及概念讲解前面三节分别采集了行情数据、财务数据和事件数据。现在需要把它们统一存储到一个结构化的本地数据库中。推荐使用两种格式的组合存储格式适用数据优点缺点HDF5大规模行情数据数千只股票 × 十年日线读写极快支持压缩单文件存储并发写入能力弱SQLite财务数据、事件数据字段结构固定支持SQL查询多表关联写入速度不如HDF5Parquet中间计算结果、因子值列式存储压缩率高不支持原地更新代码演示下面构建一个统一的本地数据访问层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_codeNone, days7): 查询最近的催化剂事件 cutoff (datetime.now() - timedelta(daysdays)).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, paramsparams) 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_col0, parse_datesTrue)关键要点HDF5适合一次写入、多次读取的场景写入时注意加锁避免并发冲突SQLite虽然轻量但单文件数据库模式非常适合个人量化团队的规模数百只股票、数年数据事件表要记录检测时间而非事件发生时间——同一事件可能在小晴和多处来源同时被捕获数据备份是必须的——本地数据库不是云存储硬盘故障会丢失所有历史数据实战总结本节小晴情报官完成了他的第一次入职任务——搭建从数据采集到本地存储的完整管道采集管道统一管理行情、财务、宏观三类数据的拉取逻辑支持增量更新清洗规则处理复权、停牌、财务时点匹配和异常值检测保证数据质量催化剂识别利用大模型从新闻公告中提取结构化事件信息辅助策略决策数据仓库HDF5SQLite组合存储让所有后续分析都能在本地秒级读取常见误区提醒误区一数据采集频率越高越好。日线级别策略用日线数据足够了tick级数据只会徒增存储和计算负担误区二所有数据都用CSV存。百只股票以上的日线数据用HDF5是必须的CSV的I/O会成为瓶颈误区三忽略财务数据的时点对齐。这是导致回测与实际偏差的最隐蔽原因——偷看未来数据下一讲小雅分析师将接过小晴准备好的数据正式登场——用Backtrader回测引擎跑通第一个双均线策略并学会解读收益、回撤、夏普比率等核心绩效指标。下一篇预告第05讲小雅初试锋芒——Backtrader回测引擎入门实战数据和环境已就绪是时候让策略接受历史的检验了。小雅将带着双均线策略走进Backtrader的回测战场用真实数据回答这个策略到底能不能赚钱。关注专栏不错过后续内容。