ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

从数据流到因子库:量化因子工程化实践指南

从数据流到因子库:量化因子工程化实践指南 写在前面的结论这次我们来看一个偏工程、偏底层的量化话题从数据流到因子库。这个标题来自我正在推进的“365 天量化金融”系列第 90 篇走到这一步前面的碎数据、单因子、单标的逻辑已经不能再靠散装脚本支撑了。无论你是在做 stocks、加密货币还是期货策略都会在某个阶段被同一类问题卡住数据怎么稳定落地、因子怎么统一调度、回测怎么批量跑、信号怎么供到上层。本文会把“因子阶段收官”这件事拆开讲为什么先有数据流、再有因子库中间每一步解决什么问题以及后续策略研发应该怎么组织。先给出这篇内容的几个核心看点。核心能力说明主题类型量化金融工程方案数据流、因子库、批量回测核心价值把散装脚本收拢为可复用的因子计算与数据管理体系适用阶段单因子验证通过、准备上规模做多因子研究的阶段技术重点数据流管道、因子注册与缓存、批量调度、结果落库工程门槛对 Python 有一定基础能理解类、装饰器、任务队列即可推荐运行环境Linux 或 macOS本地开发建议 Python 3.10输入材料以本项目原始工程思路为准合规提示文中涉及的行情数据、因子计算仅作为技术研究不构成投资建议下面按“从数据流到因子库”的主线把这套工程方案完整讲一遍。1. 核心能力速览在动手写代码之前先明确这套方案的定位。它解决的问题不是“某个因子怎么算”而是“大量因子怎么组织、怎么批量计算、怎么避免算完就丢”。从本项目的原始材料看核心能力可以整理为下表。能力项说明项目类型量化金融工程框架侧重数据流与因子库搭建数据接入支持统一行情数据入口按标的分批拉取、落盘、校验因子计算采用因子类 装饰器注册机制新因子无需改动主流程因子存储因子表按标的 日期主键落库便于回测与信号拼接批量任务支持全市场标的循环计算并进行断点续跑缓存机制数据层与因子层均可缓存避免重复计算调度方式可对接 APScheduler、Airflow 或自行实现简单队列对外输出因子表、信号表、特征宽表供后续模型与策略调用合规边界只负责数据处理与因子研究不涉及实盘交易指令这里有一个关键逻辑数据流是因子库的底座。如果数据层还在用零散的 CSV 文件和手写 IO因子库做得再花哨也跑不稳。所以本文会先讲数据流再讲因子库最后落到批量验证和性能观察。这样安排是因为项目从第 90 篇开始已经过了“研究单个因子是否有效”的阶段正式进入“批量生产因子”的阶段。2. 适用场景与使用边界这套方案适合怎么用、不适合怎么用其实很清楚。2.1 适合的场景手里已经有几个通过初步检验的因子想统一放到一个库里管理。每天要处理多标的、多周期行情数据不想再手动读多个文件。需要把因子计算结果供多个策略复用避免每个策略各算一遍。开始做因子绩效分析、相关性分析、IC 分析需要结构化因子数据。后续要接入机器学习模型需要把因子表转成特征宽表。在这些场景下数据流 因子库的收益非常大。你可以把“计算因子”当成一次注册、批量执行、落库的过程而不是写一次、跑一次、丢掉一次。2.2 不适合的场景只做一两个标的、一两个因子手工 Excel 反而更快。完全不了解任何量化概念还没搞清楚 K 线、收益率、回测是什么就先别搭工程。需要毫秒级实时交易信号这套偏离线批处理不适合直接作为交易执行引擎。没有授权数据源靠爬取不合法行情做研究风险很高不建议。2.3 使用边界与合规提醒量化研究涉及数据版权、策略保密和资金风险。需要注意以下边界行情数据必须来自合法授权渠道不要使用未经许可抓取的数据。因子库中的计算结果用于研究与回测不构成投资建议。如果因子涉及新闻、舆情、另类数据还需要另行处理文本授权问题。回测表现好不代表实盘表现好任何策略上线前都应做样本外验证和风险评估。这也是为什么本文会把“批量任务”“缓存机制”“断点续跑”这些工程能力看得和因子公式本身一样重要。没有工程保障再好的因子也无法稳定产出。3. 环境准备与前置条件在搭建数据流和因子库之前先确认环境能支撑后续开发。多数量化工程用 Python 完成因此优先准备 Python 环境和常用库。3.1 基础环境清单环境项建议配置操作系统Linux / macOS 最佳Windows 可运行但部分调度组件支持较弱Python3.10 或 3.11 均可数据存储SQLite 起步数据量上来后转 PostgreSQL缓存本地 parquet 或 Redis按需选择调度先不用框架直接用脚本循环再考虑 APScheduler包管理uv 或 pip venv版本管理git3.2 需要安装的核心依赖这里给出一个常见的依赖清单。具体版本可以根据实际环境调整。# 数据与计算 pandas numpy pyarrow # 数据库 sqlalchemy psycopg2-binary # 如果需要连接 PostgreSQL # 任务调度可选 apscheduler # 配置管理 pydantic python-dotenv # 日志 loguru安装命令示例pip install pandas numpy pyarrow sqlalchemy apscheduler pydantic python-dotenv loguru3.3 目录结构设计从材料看项目走到因子库阶段目录结构需要重新组织。建议采用以下结构。quant_research/ ├── config/ │ └── settings.py ├── data/ │ ├── raw/ # 原始行情数据落地 │ └── processed/ # 清洗后的标准数据 ├── factors/ │ ├── base.py # 因子基类与注册表 │ ├── registry.py # 因子注册与发现 │ ├── technical/ # 技术类因子 │ ├── fundamental/ # 基本面因子 │ └── alternative/ # 另类因子 ├── pipelines/ │ ├── data_pipeline.py │ └── factor_pipeline.py ├── storage/ │ ├── database.py │ └── cache.py ├── utils/ │ ├── logger.py │ └── time_utils.py └── main.py这个目录结构把“数据层、因子层、调度层、存储层”分开目的就是让新因子、新数据源可以无侵入地加入。4. 数据流管道设计“从数据流到因子库”的第一步是把行情数据整理成一个稳定、可复现的输入。数据流管道的职责是拉取原始数据、清洗、校验、按统一格式落盘或入库提供给因子计算层。4.1 数据流核心接口我们可以把数据流抽象成一个标准接口。无论数据来自 CSV、数据库还是第三方 API对外都返回同一种结构。# pipelines/data_pipeline.py from abc import ABC, abstractmethod import pandas as pd class DataSource(ABC): 数据源抽象接口 abstractmethod def fetch(self, symbol: str, start: str, end: str) - pd.DataFrame: 获取某个标的时间区间内的行情数据 class DataPipeline: 数据管道负责清洗、校验、标准化 def __init__(self, source: DataSource): self.source source def load(self, symbol: str, start: str, end: str) - pd.DataFrame: df self.source.fetch(symbol, start, end) if df.empty: return df df self._clean(df) self._validate(df) return df def _clean(self, df: pd.DataFrame) - pd.DataFrame: # 统一列名与时间索引 df df.rename(columns{ datetime: time, open: open, high: high, low: low, close: close, volume: volume, }) df[time] pd.to_datetime(df[time]) df df.sort_values(time).drop_duplicates(subset[time]) return df def _validate(self, df: pd.DataFrame) - None: if df[close].isnull().any(): raise ValueError(close 列存在缺失值) if (df[high] df[low]).any(): raise ValueError(high 小于 low数据异常)这样设计的好处在于后续新增数据源时只需要实现DataSource接口不需要改动清洗逻辑和因子计算逻辑。4.2 数据落盘与缓存数据流中要特别注意“重复拉取”的问题。如果每天拉全量数据既慢又容易触发数据源限制。建议在数据管道中加入缓存判断。# storage/cache.py import os import pandas as pd class DataCache: 本地 parquet 缓存 def __init__(self, cache_dir: str data/processed): self.cache_dir cache_dir os.makedirs(cache_dir, exist_okTrue) def _path(self, symbol: str, freq: str) - str: return os.path.join(self.cache_dir, f{symbol}_{freq}.parquet) def exists(self, symbol: str, freq: str) - bool: return os.path.exists(self._path(symbol, freq)) def read(self, symbol: str, freq: str) - pd.DataFrame: return pd.read_parquet(self._path(symbol, freq)) def write(self, symbol: str, freq: str, df: pd.DataFrame) - None: df.to_parquet(self._path(symbol, freq), indexFalse) class FactorCache(DataCache): 因子层缓存单独目录 def __init__(self, cache_dir: str data/factors): super().__init__(cache_dir)缓存的意义不只是省时间。它能保证同一份数据在因子计算时是可复现的避免不同时间拉到的数据不一致导致因子值出现偏差。4.3 数据流验证方法完成数据管道后先做一次小样本验证。输入某标的 2020 年到 2024 年的日线数据。预期输出标准化 DataFrame列为 time/open/high/low/close/volume按时间升序且无重复。判断标准数据量合理、无空值、高低价无倒挂。常见失败原因数据源字段命名不一致导致rename失败。时区处理不一致导致日期重复。复权因子未处理导致价格跳变。建议在正式跑全量数据前先用单个标的验证。不要一上来就对几千个标的做数据流清洗否则问题会被放大。5. 因子库搭建因子库是整个阶段收官的核心。它的核心能力是“注册因子、统一调度、批量计算、自动落库”。5.1 因子基类设计先定义一个因子基类所有因子都继承它并实现compute方法。# factors/base.py from abc import ABC, abstractmethod import pandas as pd class Factor(ABC): 因子基类 name: str base_factor description: str params: dict {} def __init__(self, **kwargs): self.params.update(kwargs) abstractmethod def compute(self, data: pd.DataFrame) - pd.Series: 输入行情数据输出因子序列 def __repr__(self): return fFactor {self.name} params{self.params}每个具体因子只需要实现compute返回一个与行情数据索引对齐的Series。例如一个简单动量因子# factors/technical/momentum.py import pandas as pd from factors.base import Factor class MomentumFactor(Factor): N 日动量因子 name momentum description N 日收益率动量 def compute(self, data: pd.DataFrame) - pd.Series: n self.params.get(n, 20) return data[close].pct_change(n)这里把name当作因子唯一标识后续注册、缓存、落库都以它为 key。5.2 因子注册机制因子上规模后最忌讳在业务代码里写大量 if-else。更好的方案是维护一个全局因子注册表通过装饰器自动注册。# factors/registry.py from typing import Dict, Type from factors.base import Factor class FactorRegistry: _registry: Dict[str, Type[Factor]] {} classmethod def register(cls, factor_cls: Type[Factor]): key factor_cls.name if key in cls._registry: raise ValueError(f因子 {key} 已存在) cls._registry[key] factor_cls return factor_cls classmethod def get(cls, name: str) - Type[Factor]: if name not in cls._registry: raise KeyError(f因子 {name} 未注册) return cls._registry[name] classmethod def all_names(cls): return list(cls._registry.keys()) # 用装饰器注册 def register_factor(cls): return FactorRegistry.register(cls)这样每个新因子文件只要被 import就会被自动注册。# factors/technical/__init__.py from factors.registry import register_factor from .momentum import MomentumFactor register_factor(MomentumFactor)后续新增因子时只需要添加新文件并注册主流程完全不需要改。5.3 因子计算调度器有了注册表和因子类后需要一层调度逻辑来批量计算因子。# pipelines/factor_pipeline.py import pandas as pd from factors.registry import FactorRegistry from storage.cache import DataCache, FactorCache class FactorPipeline: 因子计算管道 def __init__(self, data_cache: DataCache, factor_cache: FactorCache): self.data_cache data_cache self.factor_cache factor_cache def compute_factor(self, factor_name: str, symbol: str, freq: str): cache_key f{factor_name}_{symbol}_{freq} if self.factor_cache.exists(symbol, cache_key): print(f读取缓存因子: {cache_key}) return self.factor_cache.read(symbol, cache_key) # 加载行情数据 data self.data_cache.read(symbol, freq) # 获取因子类 factor_cls FactorRegistry.get(factor_name) factor factor_cls() # 计算因子 factor_series factor.compute(data) factor_df pd.DataFrame({ time: data[time], factor_value: factor_series.values, }) # 缓存 self.factor_cache.write(symbol, cache_key, factor_df) print(f因子计算完成: {cache_key}) return factor_df5.4 因子批量计算单因子跑通后批量计算就是把“单标的单因子”的流程扩展到“多标的 × 多因子”。# main.py from storage.cache import DataCache, FactorCache from pipelines.factor_pipeline import FactorPipeline from factors.registry import FactorRegistry def run_batch(symbols, factors, freq1d): data_cache DataCache(data/processed) factor_cache FactorCache(data/factors) pipeline FactorPipeline(data_cache, factor_cache) for symbol in symbols: for factor_name in factors: try: df pipeline.compute_factor(factor_name, symbol, freq) if df.empty: print(f跳过空结果: {symbol} {factor_name}) except Exception as e: print(f计算失败: {symbol} {factor_name}, error{e}) if __name__ __main__: symbols [000001, 000002, 000003] factors FactorRegistry.all_names() run_batch(symbols, factors)批量任务最怕中途崩掉。建议通过缓存机制做断点续跑已经计算并缓存的因子会直接跳过失败的重跑一次即可。如果愿意引入更正式的调度可以对接 APScheduler 实现定时计算from apscheduler.schedulers.blocking import BlockingScheduler scheduler BlockingScheduler() scheduler.scheduled_job(cron, hour18, minute30) def job(): symbols load_symbols() factors FactorRegistry.all_names() run_batch(symbols, factors) scheduler.start()不过第一次部署时建议先手动运行不要直接挂调度。确认因子值和落库逻辑无误再上定时任务。6. 因子落库与数据查询因子计算完成后不能只存在 parquet 缓存里。为了后续回测和策略复用还需要把因子计算结果统一落到数据库。6.1 因子表结构因子数据的特点是“长表”比“宽表”更适合扩展。推荐结构如下。字段名类型说明symbolVARCHAR标的代码factor_nameVARCHAR因子名称timeTIMESTAMP时间戳valueDOUBLE因子值paramsVARCHAR因子参数, JSON 格式这样设计的好处是新增因子不需要改表结构查询时按factor_name过滤即可。6.2 SQLAlchemy 写入示例# storage/database.py from sqlalchemy import create_engine, Column, String, DateTime, Float, Index from sqlalchemy.orm import declarative_base, sessionmaker import pandas as pd Base declarative_base() class FactorRecord(Base): __tablename__ factor_values id Column(String, primary_keyTrue) symbol Column(String, indexTrue) factor_name Column(String, indexTrue) time Column(DateTime, indexTrue) value Column(Float) def save_factor(symbol: str, factor_name: str, df: pd.DataFrame, engine): session sessionmaker(bindengine)() try: records [] for _, row in df.iterrows(): records.append(FactorRecord( idf{symbol}_{factor_name}_{row[time]}, symbolsymbol, factor_namefactor_name, timerow[time], valuerow[factor_value], )) session.add_all(records) session.commit() except Exception as e: session.rollback() raise e finally: session.close()这里使用主键去重。重复运行同样的因子计算不会产生重复记录。6.3 查询因子数据查询时直接按symbol和factor_name拼成宽表即可。import pandas as pd from sqlalchemy import create_engine engine create_engine(sqlite:///factor.db) query SELECT time, MAX(CASE WHEN factor_namemomentum THEN value END) AS momentum, MAX(CASE WHEN factor_namevolatility THEN value END) AS volatility FROM factor_values WHERE symbol000001 GROUP BY time ORDER BY time df pd.read_sql(query, engine)这是一个典型的因子转宽表查询。后续做 IC 分析、分层回测或机器学习特征工程都可以基于这个宽表推进。7. 批量任务与调度注意事项批量任务是因子库规模化的关键。下面重点说明怎么设计批量任务以及如何避免运行过多、过慢、过乱的问题。7.1 批量任务三层设计从项目实操来看批量任务建议分成三层。第一层是数据准备任务。先确保所有标的数据都已更新到最新并且通过校验。这一步失败后面因子计算没有意义。第二层是因子计算任务。遍历因子注册表计算每个因子的值并缓存。量大的时候可以使用多进程但要注意内存占用。第三层是因子落库任务。把缓存中的因子结果写入数据库生成特征宽表供回测和策略使用。7.2 使用简单任务队列如果暂时不想引入 Celery可以先用一个简单的 FIFO 队列。import queue import threading import time def worker(): while True: task task_queue.get() if task is None: break symbol, factor_name task try: pipeline.compute_factor(factor_name, symbol) print(f完成: {symbol} {factor_name}) except Exception as e: print(f失败: {symbol} {factor_name}, error{e}) finally: task_queue.task_done() task_queue queue.Queue() threads [threading.Thread(targetworker) for _ in range(4)] for t in threads: t.start() for symbol in symbols: for factor_name in FactorRegistry.all_names(): task_queue.put((symbol, factor_name)) task_queue.join() for _ in threads: task_queue.put(None)多线程适合 IO 密集型任务因子计算如果是纯 CPU 密集多进程更合适。但第一次跑尽量先单线程确认不存在计算 bug 再考虑并发。7.3 断点续跑与日志批量跑因子时建议把日志写到文件而不是只看终端。from loguru import logger logger.add(logs/factor_batch.log, rotation500 MB, retention7 days) logger.info(f开始批量计算, 标的数{len(symbols)}, 因子数{len(FactorRegistry.all_names())})同时把每个标的、每个因子的“成功/失败”状态记录下来便于失败后重试。状态含义处理方式success已计算并缓存跳过failed计算异常查看日志修复后重试empty数据为空检查数据源是否覆盖该标的pending未计算正常执行8. 资源占用与性能观察量化工程的性能问题不像大模型那么极端但也不能忽略。数据量大的时候Pandas 全量计算一样会吃满内存。8.1 内存优化读取行情数据时只保留需要的列不要全字段加载。优先使用float32而不是float64能减少一半内存。因子计算完成后及时释放不需要的中间 DataFrame。数据量极大时用polars替代pandas性能提升明显。import polars as pl df pl.read_parquet(data/processed/000001_1d.parquet)不过输入项目本身使用的是 pandas 风格因此本文以 pandas 为主。若遇到性能瓶颈可以平滑切换。8.2 显存占用观察量化因子计算一般不涉及 GPU但如果你在因子库中加入了嵌入模型、NLP 舆情因子就会涉及显存。先用 CPU 跑小批量文本数据确认因子逻辑正确。GPU 推理关注显存占用可以用nvidia-smi实时观察。批次大小调小可以降低显存峰值。如果同时跑多个模型建议按模型单独启动进程。当前阶段如果没有 NLP 和深度学习因子可以不考虑 GPU先专注把 CPU 因子库跑稳。8.3 计算耗时观察建立性能基准非常重要。第一次跑完因子库后记录如下指标单标的单因子平均耗时。全市场全因子单轮总耗时。缓存命中后的重复查询耗时。数据库写入耗时。以单一标的数据为例当数据量从几千行增长到几百万行时因子计算耗时可能显著上升。此时建议使用增量计算只对新数据区间计算因子再与历史因子拼接。# 增量更新示例 def compute_incremental(symbol, factor_name, start_date): old_df factor_cache.read(symbol, factor_name) new_data data_module.load_since(symbol, start_date) new_factor factor.compute(new_data) final_df pd.concat([old_df, new_factor], ignore_indexTrue) final_df final_df.drop_duplicates(subset[time]) return final_df增量计算可以极大减少重复计算。数据每天更新时只算最近几个交易日即可。9. 从因子库到策略验证因子库建好之后下一步就是拿因子去验证策略效果。这里给出从因子库到策略的通用流程。9.1 因子有效性初筛拿到因子宽表后先做基础统计因子覆盖率。因子缺失值比例。因子值分布。因子与未来收益的 IC 均值。import pandas as pd factor_df pd.read_sql(query, engine) # 计算未来 5 日收益 factor_df[future_return] factor_df[close].pct_change(5).shift(-5) # 计算 IC ic factor_df.groupby(time).apply( lambda x: x[factor_value].corr(x[future_return]) ) print(ic.mean(), ic.std())IC 均值和 ICIR 可以快速判断一个因子是否值得继续研究。9.2 分层回测别只会看 IC还要做分组回测。把因子值按每日分位数分为 5 组或 10 组观察多空组合表现。factor_df[group] factor_df.groupby(time)[factor_value].transform( lambda x: pd.qcut(x, 5, labelsFalse, duplicatesdrop) ) group_return factor_df.groupby([time, group])[future_return].mean()如果多头组与空头组收益差距稳定因子才有进一步使用价值。9.3 因子库延展方向因子库本身只是第一步。后续可以扩展为基于因子库构建多因子打分模型。接入 LightGBM 做因子合成。构造行业中性化因子。将因子值接入回测引擎形成信号端到端流程。10. 常见问题与排查方法在实际部署这套数据流与因子库时容易遇到下面这些问题。问题现象可能原因排查方式解决方案数据拉取后时间索引混乱时区未处理检查原始数据的时区字段统一转 UTC 或本地时区因子值大量为 NaN数据区间不足或复权问题查看数据起始日期和缺失比例扩展数据区间或调整参数重复运行后因子记录翻倍主键设计缺失检查数据库主键使用symbol factor time作为主键批量任务中途崩溃某个标的因子计算异常查看日志定位增加异常捕获和断点续跑内存占用过高一次性加载全市场数据观察内存曲线改用分标的循环或增量计算API 拉数被限制请求频率过高检查返回状态码增加 sleep 或使用缓存因子值整体偏移收益率计算方向错误抽查单个标的因子曲线核对 pct_change 的参数和 shift 方向数据库写入过慢逐行插入查看 SQL 日志改用批量插入或 to_sql遇到问题先看日志。建议在数据管道、因子计算、落库三个阶段分别打日志这样能快速定位是哪一层出了问题。11. 最佳实践与工程建议最后把“从数据流到因子库”过程中的几条工程经验总结一下。先跑通最小闭环再铺量。不要一开始就对全市场几千个标的跑批量任务。先用 3 到 5 个标的、2 到 3 个因子跑通数据清洗、因子计算、缓存、落库、查询整个流程确认无误后再扩大到全市场。因子注册表不要轻易删除历史因子。因子库的一大价值是积累。即使某个因子当前表现一般后续换数据区间或换标的池后可能有效保留历史计算记录有助于复盘。缓存目录和数据库要定期备份。因子计算可复现的前提是数据和缓存都稳定。建议原始数据只增不改因子缓存可以定期清理但数据库中的因子历史不要随意删除。引入配置中心管理标的池和因子参数。不要把所有标的列表和参数硬编码到代码里。用 YAML 或环境变量管理方便不同环境切换。# config/factor_config.yaml symbols: - 000001 - 000002 - 000003 factors: momentum: n: 20 volatility: window: 20批量任务记录状态并支持重试。最简单的方式是把每个因子的计算状态写到一个 state 表任务失败后可以根据状态表重跑失败项。要时刻关注合规边界。行情数据的来源要合法使用的数据要遵守授权协议。如果因子涉及新闻或另类数据注意文本版权。不要把因子库的离线研究结果直接当作实盘信号使用需经过完整回测和风险评估。12. 总结与下一步这次项目做到第 90 篇最大的变化不是因子数量变多了而是整个研究过程开始工程化。数据流负责把数据稳定地送到因子层因子库负责把算力集中在因子的批量生产和复用上批量调度和缓存机制保证整个流程可以在全市场重复执行。走到这一步单因子的零散研究正式收口成了一套可迭代的因子生产系统。接下来最先应该验证的是把你手头最熟悉的几个因子按这套结构注册进去用 3 到 5 个标的跑一遍完整流程检查因子值、缓存、数据库和 IC 结果是否符合预期。最容易踩的坑是数据清洗不一致和缓存主键冲突这两块建议优先测试。熟练之后可以把数据源扩展、因子合成、机器学习模型和更完善的回测引擎接进来形成自己的量化研究基础设施。建议收藏备用。下次再提到“因子阶段收官”不只是看几个公式而是看整个数据流、因子库、批量任务和验证闭环怎么跑通。
返回列表