Skip to content

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)          --> 统一查询接口: 在线/离线特征合并
PYTHON164 行 · 5.1 KB
📄此处有展示代码164 行 · 5.1 KB展开 ▼
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 量化策略中的处理模式划分

PYTHON55 行 · 1.6 KB
📄此处有展示代码55 行 · 1.6 KB展开 ▼
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 数据血缘的重要性

在复杂的量化数据管道中,理解每个数据字段的来源、转换历史和依赖关系至关重要:

  • 调试:当发现数据异常时,快速追溯到根源
  • 影响分析:上游数据变更时,评估哪些下游受影响
  • 合规:满足监管对数据溯源的要求
PYTHON96 行 · 3.1 KB
📄此处有展示代码96 行 · 3.1 KB展开 ▼
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`(数据血缘追踪器)。该片段为教学展示(未包含独立运行的输入数据),可在实战练习中结合真实数据调用。

四、计算引擎对比

引擎处理范式最佳场景编程模型延迟
Apache Spark批处理为主,微批处理(Structured Streaming)大规模历史数据 ETL、因子计算、回测RDD/DataFrame秒级
Apache Flink原生流处理(Event-by-Event)实时特征、实时风控、实时信号DataStream/Table API毫秒级
DaskPython 原生并行计算中等规模的分析、Pandas 生态兼容Dask DataFrame/Array秒级
PYTHON25 行 · 1013 B
📄此处有展示代码25 行 · 1013 B展开 ▼
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(特征存储)所要解决的问题。