主题切换
17.4 批处理与流处理
名词解释:Lambda vs Kappa 架构
Lambda 架构和 Kappa 架构是大数据处理中的两种主流范式。Lambda 架构同时维护批处理层(负责精确但高延迟的计算)和流处理层(负责近似但低延迟的计算),通过服务层合并两者结果。Kappa 架构则认为批处理是流处理的特殊情况,只使用流处理引擎处理所有数据,通过重放(Replay)日志来实现历史重新计算。
一、Lambda vs Kappa 架构
1.1 两种架构的对比
| 维度 | Lambda 架构 | Kappa 架构 |
|---|---|---|
| 处理层数 | 三层(批处理 + 流处理 + 服务层) | 单层(纯流处理) |
| 数据一致性 | 最终一致性(批处理校正流处理近似结果) | 单一真相来源 |
| 运维复杂度 | 高(需维护两套代码逻辑) | 中(单一代码库) |
| 延迟 | 流层低延迟,批层高延迟 | 始终低延迟 |
| 历史重算 | 批处理自然支持 | 需要重放(Replay)日志 |
| 适用场景 | 海量历史数据 + 实时数据兼顾 | 以实时为主、历史为辅 |
1.2 量化交易中的选择
在量化交易场景下,大部分工作负载仍以批处理为主(日终因子计算、风险归因、回测),但实时特征计算是流处理的典型场景。因此,实际落地通常是混合架构:
[Lambda 架构适配量化]
实时层 (Kafka Streams/Flink) --> 实时特征: 当前价格、日内动量、实时VPIN
批处理层 (Spark/Dask) --> 日终因子: 估值因子、质量因子、行业暴露
服务层 (Redis/Feast) --> 统一查询接口: 在线/离线特征合并此处有展示代码展开 ▼
python
import numpy as np
import pandas as pd
from abc import ABC, abstractmethod
from typing import Dict, Any, Optional
class ProcessingArchitecture(ABC):
"""
量化数据处理的抽象基类,定义批处理和流处理的统一接口。
"""
@abstractmethod
def process(self, data, **kwargs):
pass
@abstractmethod
def get_state(self) -> Dict[str, Any]:
pass
class BatchProcessor(ProcessingArchitecture):
"""
日终批处理器:在每日收盘后运行,处理全量历史数据。
适用场景:因子计算、风险模型估计、日频回测数据准备。
"""
def __init__(self, data_source: str, output_path: str):
self.data_source = data_source
self.output_path = output_path
self.state = {'last_run': None, 'records_processed': 0}
def process(self,
date: str,
symbols: Optional[list] = None) -> pd.DataFrame:
"""
对指定日期的全量数据进行批处理。
典型流程:
1. 读取原始数据
2. 数据清洗和标准化
3. 计算日频因子
4. 写入结果存储
"""
# 读取原始数据
raw_data = self._read_raw_data(date, symbols)
# 数据清洗
clean_data = self._cleanse(raw_data)
# 因子计算
factors = self._compute_factors(clean_data)
# 写入结果
self._write_output(factors, date)
# 更新状态
self.state['last_run'] = date
self.state['records_processed'] += len(clean_data)
return factors
def _read_raw_data(self, date: str, symbols: list) -> pd.DataFrame:
"""从分区存储中读取指定日期的原始数据"""
date_path = f"{self.data_source}/{date.replace('-', '')}"
return pd.read_parquet(date_path)
def _cleanse(self, data: pd.DataFrame) -> pd.DataFrame:
"""数据清洗管道"""
# 去重
data = data.drop_duplicates()
# 过滤无效价格
data = data[data['price'] > 0]
# 过滤非交易时段
data = data[
(data['timestamp'].dt.time >= pd.Timestamp('09:30').time()) &
(data['timestamp'].dt.time <= pd.Timestamp('16:00').time())
]
return data
def _compute_factors(self, data: pd.DataFrame) -> pd.DataFrame:
"""因子计算引擎"""
grouped = data.groupby('symbol')
factors = pd.DataFrame(index=data['symbol'].unique())
# 日频因子
factors['daily_return'] = grouped['price'].apply(
lambda x: x.iloc[-1] / x.iloc[0] - 1
)
factors['realized_vol'] = grouped['price'].apply(
lambda x: np.std(np.diff(np.log(x))) * np.sqrt(252)
)
factors['volume'] = grouped['size'].sum()
factors['vwap'] = grouped.apply(
lambda g: np.average(g['price'], weights=g['size'])
)
return factors
def _write_output(self, factors: pd.DataFrame, date: str):
"""写入处理结果"""
output_file = f"{self.output_path}/factors/{date.replace('-', '')}.parquet"
factors.to_parquet(output_file)
def get_state(self) -> Dict[str, Any]:
return self.state
class StreamProcessor(ProcessingArchitecture):
"""
实时流处理器:持续处理流入的数据,维护增量状态。
适用场景:实时特征、实时风险监控、实时信号生成。
"""
def __init__(self, window_size: int = 100):
self.window_size = window_size
self.window_buffer: Dict[str, list] = {} # {symbol: [prices]}
self.state = {'messages_processed': 0, 'last_update': None}
def process(self,
tick: Dict[str, Any],
symbol: str) -> Optional[Dict[str, float]]:
"""
处理单个 tick,更新滚动窗口状态并计算实时特征。
返回:
当前窗口的实时特征字典,如果窗口数据不足则返回 None
"""
# 初始化窗口
if symbol not in self.window_buffer:
self.window_buffer[symbol] = []
# 更新窗口
self.window_buffer[symbol].append(tick['price'])
# 维护窗口大小
if len(self.window_buffer[symbol]) > self.window_size:
self.window_buffer[symbol] = \
self.window_buffer[symbol][-self.window_size:]
# 窗口数据不足
if len(self.window_buffer[symbol]) < self.window_size:
return None
prices = np.array(self.window_buffer[symbol])
returns = np.diff(np.log(prices))
# 实时特征
features = {
'real_time_vwap': np.mean(prices),
'real_time_vol': np.std(returns) * np.sqrt(252),
'real_time_skew': pd.Series(returns).skew(),
'real_time_momentum': prices[-1] / prices[-self.window_size] - 1,
'bid_ask_spread': tick.get('ask', 0) - tick.get('bid', 0)
}
self.state['messages_processed'] += 1
self.state['last_update'] = tick.get('timestamp')
return features
def get_state(self) -> Dict[str, Any]:
return self.state点击展开可浏览运行结果
📘 本段代码定义了 3 个函数/类:类 `ProcessingArchitecture`(量化数据处理的抽象基类,定义批处理和流处理的统一接口)、类 `BatchProcessor`(日终批处理器:在每日收盘后运行,处理全量历史数据)、类 `StreamProcessor`(实时流处理器:持续处理流入的数据,维护增量状态)。该片段为教学展示(未包含独立运行的输入数据),可在实战练习中结合真实数据调用。
二、日终批处理 vs 实时流处理
2.1 量化策略中的处理模式划分
此处有展示代码展开 ▼
python
class QuantProcessingScheduler:
"""
量化策略的数据处理调度器。
协调批处理和流处理任务的执行。
"""
def __init__(self):
self.daily_batch_tasks = []
self.intraday_stream_tasks = []
def register_daily_task(self, task_name: str, func, priority: int = 0):
"""
注册日终批处理任务。
执行顺序按 priority 升序(数字越小越先执行)。
"""
self.daily_batch_tasks.append({
'name': task_name,
'func': func,
'priority': priority
})
self.daily_batch_tasks.sort(key=lambda x: x['priority'])
def register_stream_task(self, task_name: str, func):
"""注册实时流处理任务"""
self.intraday_stream_tasks.append({
'name': task_name,
'func': func
})
def run_daily_batch(self, date: str):
"""
运行日终批处理管道。
典型任务链:
1. 数据下载和校验
2. 日频因子计算
3. 组合估值和风险计算
4. 策略信号生成
5. 报告生成
6. 数据备份
"""
results = {}
for task in self.daily_batch_tasks:
try:
result = task['func'](date)
results[task['name']] = 'SUCCESS'
except Exception as e:
results[task['name']] = f'FAILED: {str(e)}'
# 关键任务失败时停止
if task['priority'] <= 2:
raise
return results点击展开可浏览运行结果
📘 本段代码定义了 1 个函数/类:类 `QuantProcessingScheduler`(量化策略的数据处理调度器)。该片段为教学展示(未包含独立运行的输入数据),可在实战练习中结合真实数据调用。
三、数据血缘(Data Lineage)追踪
3.1 数据血缘的重要性
在复杂的量化数据管道中,理解每个数据字段的来源、转换历史和依赖关系至关重要:
- 调试:当发现数据异常时,快速追溯到根源
- 影响分析:上游数据变更时,评估哪些下游受影响
- 合规:满足监管对数据溯源的要求
此处有展示代码展开 ▼
python
import hashlib
import json
from datetime import datetime
from typing import List, Dict, Any
class DataLineageTracker:
"""
数据血缘追踪器。
记录每个数据集的来源、转换步骤和依赖关系。
使用有向无环图(DAG)表示数据流转。
"""
def __init__(self):
self.lineage_graph: Dict[str, Dict[str, Any]] = {}
def record_transformation(self,
input_datasets: List[str],
output_dataset: str,
transformation_name: str,
parameters: Dict[str, Any],
code_version: str = 'latest') -> str:
"""
记录一次数据转换的血缘信息。
参数:
input_datasets: 上游输入数据集名称列表
output_dataset: 输出数据集名称
transformation_name: 转换函数/任务名称
parameters: 转换参数
code_version: 代码版本
返回:
本次转换的唯一 ID
"""
# 生成 lineage ID
content = f"{transformation_name}:{output_dataset}:{datetime.now().isoformat()}"
lineage_id = hashlib.md5(content.encode()).hexdigest()[:12]
# 记录输入数据集的哈希(确保可复现)
input_hashes = {}
for ds in input_datasets:
if ds in self.lineage_graph:
input_hashes[ds] = self.lineage_graph[ds].get('output_hash')
record = {
'lineage_id': lineage_id,
'timestamp': datetime.now().isoformat(),
'input_datasets': input_datasets,
'input_hashes': input_hashes,
'output_dataset': output_dataset,
'transformation': transformation_name,
'parameters': parameters,
'code_version': code_version,
'dependencies': input_datasets.copy()
}
self.lineage_graph[output_dataset] = record
return lineage_id
def trace_lineage(self, dataset: str) -> Dict[str, Any]:
"""
追溯指定数据集的完整上游血缘链。
递归追溯所有上游依赖,构建完整的 DAG 路径。
"""
if dataset not in self.lineage_graph:
return {'dataset': dataset, 'error': 'No lineage recorded'}
lineage = []
to_visit = [dataset]
visited = set()
while to_visit:
current = to_visit.pop(0)
if current in visited:
continue
visited.add(current)
if current in self.lineage_graph:
record = self.lineage_graph[current]
lineage.append({
'dataset': current,
'transformation': record['transformation'],
'lineage_id': record['lineage_id']
})
for upstream in record['dependencies']:
if upstream not in visited:
to_visit.append(upstream)
return {
'dataset': dataset,
'upstream_lineage': lineage,
'depth': len(lineage)
}点击展开可浏览运行结果
📘 本段代码定义了 1 个函数/类:类 `DataLineageTracker`(数据血缘追踪器)。该片段为教学展示(未包含独立运行的输入数据),可在实战练习中结合真实数据调用。
四、计算引擎对比
4.1 Spark / Flink / Dask 的比较
| 引擎 | 处理范式 | 最佳场景 | 编程模型 | 延迟 |
|---|---|---|---|---|
| Apache Spark | 批处理为主,微批处理(Structured Streaming) | 大规模历史数据 ETL、因子计算、回测 | RDD/DataFrame | 秒级 |
| Apache Flink | 原生流处理(Event-by-Event) | 实时特征、实时风控、实时信号 | DataStream/Table API | 毫秒级 |
| Dask | Python 原生并行计算 | 中等规模的分析、Pandas 生态兼容 | Dask DataFrame/Array | 秒级 |
此处有展示代码展开 ▼
python
# 量化场景下的引擎选择指南
def recommend_processing_engine(data_volume_gb: float,
latency_requirement_ms: float,
python_ecosystem_required: bool,
existing_infrastructure: str) -> str:
"""
根据需求推荐最适合的计算引擎。
参数:
data_volume_gb: 数据量(GB)
latency_requirement_ms: 延迟要求(毫秒)
python_ecosystem_required: 是否必须 Python 生态
existing_infrastructure: 现有基础设施
"""
if latency_requirement_ms < 1000:
return 'Apache Flink' # 唯一支持亚秒级延迟的
elif data_volume_gb > 100:
if existing_infrastructure == 'Hadoop/YARN':
return 'Apache Spark'
else:
return 'Apache Spark (standalone)'
elif python_ecosystem_required:
return 'Dask' # 与 Pandas API 完全兼容
else:
return 'Apache Spark' # 默认推荐点击展开可浏览运行结果
📘 本段代码定义了 1 个函数/类:函数 `recommend_processing_engine`(根据需求推荐最适合的计算引擎)。该片段为教学展示(未包含独立运行的输入数据),可在实战练习中结合真实数据调用。
批处理与流处理并不是非此即彼的选择。量化交易的最佳实践是根据计算任务的延迟要求和数据规模,选择最合适的处理范式。无论使用哪种计算引擎,最终产出的特征都需要被标准化地管理和服务——这正是 Feature Store(特征存储)所要解决的问题。