云计算百科
云计算领域专业知识百科平台

【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_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虽然轻量,但单文件数据库模式非常适合个人量化团队的规模(数百只股票、数年数据)
    • 事件表要记录"检测时间"而非"事件发生时间"——同一事件可能在小晴和多处来源同时被捕获
    • 数据备份是必须的——本地数据库不是云存储,硬盘故障会丢失所有历史数据

    实战总结

    本节小晴(情报官)完成了他的第一次"入职任务"——搭建从数据采集到本地存储的完整管道:

  • 采集管道:统一管理行情、财务、宏观三类数据的拉取逻辑,支持增量更新
  • 清洗规则:处理复权、停牌、财务时点匹配和异常值检测,保证数据质量
  • 催化剂识别:利用大模型从新闻公告中提取结构化事件信息,辅助策略决策
  • 数据仓库:HDF5+SQLite组合存储,让所有后续分析都能在本地秒级读取
  • 常见误区提醒:

    • 误区一:数据采集频率越高越好。日线级别策略用日线数据足够了,tick级数据只会徒增存储和计算负担
    • 误区二:所有数据都用CSV存。百只股票以上的日线数据用HDF5是必须的,CSV的I/O会成为瓶颈
    • 误区三:忽略财务数据的时点对齐。这是导致回测与实际偏差的最隐蔽原因——"偷看未来数据"

    下一讲,小雅(分析师)将接过小晴准备好的数据,正式登场——用Backtrader回测引擎跑通第一个双均线策略,并学会解读收益、回撤、夏普比率等核心绩效指标。


    下一篇预告

    第05讲:小雅初试锋芒——Backtrader回测引擎入门实战:数据和环境已就绪,是时候让策略接受历史的检验了。小雅将带着双均线策略走进Backtrader的回测战场,用真实数据回答"这个策略到底能不能赚钱"。


    关注专栏不错过后续内容。

    赞(0)
    未经允许不得转载:网硕互联帮助中心 » 【AI量化交易实战】第04讲:小晴(情报官)上岗——搭建本地数据管道与清洗流水线
    分享到: 更多 (0)

    评论 抢沙发

    评论前必须登录!