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 的容错机制基于 Checkpoint 和 Savepoint。
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 处理,再到输出报价的完整链路。
这张图里,Kafka 负责解耦和缓冲,Flink 做核心计算,状态后端保证数据不丢。我建议你把这个架构作为模板,根据自己业务调整。
21.6 总结一下
实时计算框架这块,我的核心建议就三条:
- 选 Flink 做交易链路,别犹豫。延迟和一致性是生命线。
- 窗口计算要谨慎,步长别太细,状态别太大。
- 容错机制必须配齐,Checkpoint 和 Savepoint 是救命稻草。
做市商系统里,实时计算不是锦上添花,是生存刚需。框架选对了,后面的事就顺了。
交易系统化学习资料 微信Strategy888888