21、实时计算框架:Flink/Spark Streaming应用、窗口计算、状态管理、容错机制

做市商系统里,实时计算就是心脏。你想想看,行情每秒跳几百次,你的定价模型、风险敞口、对冲指令,全得在毫秒级完成。我这些年折腾下来,发现选对实时框架,比选对交易策略还重要——策略错了亏钱,框架崩了直接爆仓。

今天咱们就聊聊两个主流框架:Flink 和 Spark Streaming。我不会跟你念文档,咱们直接说实战中怎么用、坑在哪。

21.1 Flink vs Spark Streaming:怎么选?

先说结论:做市商场景,我个人更倾向 Flink。为什么?因为做市商对延迟敏感,对状态一致性要求极高。

Spark Streaming 本质是微批处理,把流切成小批次。延迟在秒级,适合做风控报表、历史回测。但做实时定价?嗯,有点吃力。

Flink 是真正的逐条处理。事件一来,立刻触发计算。延迟能做到毫秒级。我在项目中遇到过,用 Spark Streaming 做期权定价,行情剧烈波动时,定价结果滞后了 3 秒,直接导致报价被套利机器人吃掉。后来换成 Flink,延迟降到 20 毫秒,再没出过这问题。

特性 Flink Spark Streaming
处理模式 真正的流处理 微批处理
延迟 毫秒级 秒级
状态管理 原生支持,强一致性 需要外部存储
容错机制 精确一次(Exactly-Once) 至少一次(At-Least-Once)
适用场景 实时定价、风控、交易 报表、监控、回测

核心建议:交易链路用 Flink,分析链路用 Spark Streaming。别混着用,也别指望一个框架打天下。

21.2 窗口计算:做市商的核心武器

做市商系统里,窗口计算无处不在。比如计算过去 5 秒的成交量加权均价(VWAP),或者过去 1 分钟的波动率。窗口就是时间上的「切片」。

Flink 支持三种窗口:

  • 滚动窗口(Tumbling Window):固定大小,不重叠。比如每 5 秒算一次 VWAP。
  • 滑动窗口(Sliding Window):固定大小,有重叠。比如每 1 秒算一次过去 5 秒的 VWAP。
  • 会话窗口(Session Window):按不活动时间切分。适合用户行为分析,做市商用得少。

我习惯用滑动窗口做实时波动率计算。举个例子:

// Flink 滑动窗口计算波动率
DataStream<Trade> trades = ...;

trades
  .keyBy(trade -> trade.getSymbol())
  .window(SlidingProcessingTimeWindows.of(
      Time.seconds(10),  // 窗口大小
      Time.seconds(1)    // 滑动步长
  ))
  .aggregate(new VolatilityAggregator())
  .map(vol -> {
      // 根据波动率调整报价价差
      double spread = baseSpread * (1 + vol * 2);
      return new Quote(symbol, bid, ask, spread);
  });

避坑指南:我曾经在滑动窗口上吃过亏。窗口大小和滑动步长设置不合理,导致计算量爆炸。比如窗口 10 秒、步长 100 毫秒,意味着每秒要维护 100 个窗口状态。记住:步长越小,状态越大,GC 压力越大。做市商场景,步长别小于 1 秒。

21.3 状态管理:别让数据丢了

做市商系统里,状态就是一切。当前持仓、未成交订单、实时 Greeks、风险限额……这些都是状态。如果状态丢了,系统就瞎了。

Flink 的状态管理是我最喜欢它的地方。它提供了:

  • ValueState:存单个值,比如当前持仓量。
  • ListState:存列表,比如未成交订单列表。
  • MapState:存键值对,比如每个合约的 Greeks。
  • BroadcastState:广播状态,比如全局参数配置。

我举个例子,用 ValueState 维护每个交易对的净头寸:

public class PositionTracker extends RichFlatMapFunction<Trade, Position> {
    private transient ValueState<Double> netPosition;

    @Override
    public void open(Configuration config) {
        ValueStateDescriptor<Double> descriptor =
            new ValueStateDescriptor<>("netPosition", Types.DOUBLE);
        netPosition = getRuntimeContext().getState(descriptor);
    }

    @Override
    public void flatMap(Trade trade, Collector<Position> out) throws Exception {
        Double current = netPosition.value();
        if (current == null) current = 0.0;

        // 更新净头寸
        double newPosition = current + trade.getQuantity();
        netPosition.update(newPosition);

        // 检查是否超过限额
        if (Math.abs(newPosition) > positionLimit) {
            // 触发风控报警
            out.collect(new Position(trade.getSymbol(), newPosition, true));
        }
    }
}

注意:状态不是无限存的。每个 key 的状态都保存在内存里,key 太多会 OOM。我建议做市商系统里,状态大小控制在 10GB 以内。超过这个量,考虑用 RocksDB 状态后端,把状态存到磁盘。

21.4 容错机制:系统挂了怎么办?

做市商系统不能停。停了就是钱。Flink 的容错机制基于 CheckpointSavepoint

Checkpoint 是自动的,每隔几秒拍一次快照。保存当前所有算子的状态和消费位置。如果挂了,从最近的 Checkpoint 恢复,保证精确一次语义。

Savepoint 是手动的,用于版本升级、代码变更。我每次上线新策略,都会先打个 Savepoint。万一新代码有问题,秒级回滚。

配置 Checkpoint 的代码很简单:

StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();

// 开启 Checkpoint,每 5 秒一次
env.enableCheckpointing(5000);

// 设置模式:精确一次
env.getCheckpointConfig().setCheckpointingMode(CheckpointingMode.EXACTLY_ONCE);

// 超时时间:1 分钟
env.getCheckpointConfig().setCheckpointTimeout(60000);

// 同时进行的 Checkpoint 数:1
env.getCheckpointConfig().setMaxConcurrentCheckpoints(1);

个人经验:Checkpoint 间隔别设太短。我见过有人设 1 秒一次,结果 Checkpoint 本身成了性能瓶颈。做市商场景,5-10 秒一次比较合理。另外,一定要监控 Checkpoint 失败次数。如果频繁失败,说明状态太大或者存储有问题,得赶紧排查。

21.5 实战架构图

下面这张图是我在项目中实际使用的实时计算架构。你可以看到数据从行情接入,到 Flink 处理,再到输出报价的完整链路。

做市商实时计算架构 行情数据源 交易所/数据商 订单数据源 交易系统/OMS Kafka 消息队列 数据缓冲与解耦 Flink 实时计算集群 窗口计算(VWAP/波动率) 状态管理(持仓/Greeks) 容错机制(Checkpoint) 报价输出 定价引擎/交易接口 状态后端(RocksDB/HDFS) 图例: 数据源 消息队列 实时计算 输出 状态后端

这张图里,Kafka 负责解耦和缓冲,Flink 做核心计算,状态后端保证数据不丢。我建议你把这个架构作为模板,根据自己业务调整。

21.6 总结一下

实时计算框架这块,我的核心建议就三条:

  • 选 Flink 做交易链路,别犹豫。延迟和一致性是生命线。
  • 窗口计算要谨慎,步长别太细,状态别太大。
  • 容错机制必须配齐,Checkpoint 和 Savepoint 是救命稻草。

做市商系统里,实时计算不是锦上添花,是生存刚需。框架选对了,后面的事就顺了。

交易系统化学习资料 微信Strategy888888