一、为什么流处理需要"全局快照"
在批处理的世界里,容错是简单的:作业失败了,重跑即可。但实时流处理面对的是一个源源不断、永不终止的数据流。系统需要在任意时刻回答一个问题:"如果此刻整个集群崩溃,我应该从哪里恢复,以及恢复后的状态应该是什么?"
这不是一个简单问题。流处理系统的状态分布在多个算子(operator)的本地存储中,数据在算子之间的网络通道(channel)中传输。任何时间点对集群"拍照",都面对一个根本性矛盾——你无法在全局同一时刻冻结所有并发执行的计算节点。
这个问题的理论基础来自 1985 年的一篇论文:K. Mani Chandy 和 Leslie Lamport 合作发表的《Distributed Snapshots: Determining Global States of Distributed Systems》。这篇论文提出的算法优雅地解决了"在异步分布式系统中如何获取一致性全局快照"的问题,成为现代流处理检查点机制的理论支柱。
二、Chandy-Lamport 算法:理论基础
2.1 模型定义
考虑一个分布式系统,由若干通过消息通道连通的进程组成。系统的全局状态(global state) = 所有进程的局部状态之和 + 所有通道的状态之和。
要进行快照,需要满足两个条件:
- 一致性(Consistency):快照必须反映一个可能的系统执行历史(causal consistency)
- 非侵入性(Non-interference):快照过程不能影响系统正常运行
2.2 核心算法
Chandy-Lamport 算法的核心思想是采用一种标记(marker)机制:
Algorithm: Chandy-Lamport Snapshot
Initiator Process P0:
1. 记录自己的局部状态
2. 向所有出站通道发送 MARKER 消息
3. 开始记录所有入站通道上的消息
Other Process Pi 收到第一条 MARKER(来自通道 C)时:
1. 如果尚未记录过局部状态:
a. 记录局部状态
b. 将通道 C 标记为空
c. 向所有出站通道发送 MARKER
d. 开始记录其他入站通道的消息
2. 如果已经记录过局部状态:
a. 停止记录通道 C 上的消息,并将之前记录的消息作为通道 C 的状态
让我们用一段伪代码理解这个算法的精妙之处:
class Process:
def __init__(self, pid, channels):
self.pid = pid
self.state = {}
self.in_channels = {c.id: [] for c in channels if c.dest == self}
self.out_channels = {c.id: c for c in channels if c.source == self}
self.recorded_state = None
self.channel_recorded = set() # 已经收到marker并完成记录的通道
def initiate_snapshot(self):
"""快照发起者调用"""
self.recorded_state = self.state.copy()
for ch in self.out_channels.values():
ch.send(MarkerMessage(sender=self.pid, snapshot_id=uuid4()))
# 开始记录所有入站通道
for ch_id, msgs in self.in_channels.items():
msgs.clear()
def on_receive(self, msg):
if isinstance(msg, MarkerMessage):
self._handle_marker(msg)
else:
# 正常业务消息
self.process_message(msg)
# 如果某个通道正在记录,把该消息加入
for ch_id in self.in_channels:
if ch_id not in self.channel_recorded:
self.in_channels[ch_id].append(msg)
def _handle_marker(self, marker):
if self.recorded_state is None:
# 首次收到marker,记录自身状态
self.recorded_state = self.state.copy()
ch_id = marker.channel_id
self.channel_recorded.add(ch_id)
# 向所有出站通道广播marker
for ch in self.out_channels.values():
ch.send(MarkerMessage(sender=self.pid,
snapshot_id=marker.snapshot_id,
channel_id=ch.id))
else:
# 已经记录过,这个marker对应通道C的状态 = 临时记录的消息
ch_id = marker.channel_id
self.channel_recorded.add(ch_id)
# 通道C的状态 = 上次快照到本次marker之间收到的消息
def get_snapshot(self):
"""返回此进程贡献的全局快照片段"""
return {
'process_state': self.recorded_state,
'channel_states': {
ch_id: list(msgs)
for ch_id, msgs in self.in_channels.items()
if ch_id in self.channel_recorded
}
}
2.3 正确性证明的关键洞察
Chandy-Lamport 算法保证快照的一致性切割(consistent cut)。所谓一致性切割,是指如果进程 Pj 在快照中显示它收到了进程 Pi 发送的消息 m,那么 Pi 在快照中也必须显示它发送了 m。换句话说,快照中不能出现"果先于因"的情况。
这一性质的保证来自两个机制:
- Marker 总是走在业务消息之前(FIFO 通道保证)
- 进程在收到 Marker 之前不会遗漏任何属于"快照前"的消息
三、从理论到工程:Flink 的异步屏障快照
Apache Flink 的容错机制基于对 Chandy-Lamport 算法的重要工程化改进,称为异步屏障快照(Asynchronous Barrier Snapshotting,ABS)。
3.1 设计动机
原始的 Chandy-Lamport 算法要求进程在收到 Marker 时"暂停处理、保存状态"。这在流处理场景中代价太高——一个算子可能持有 GB 甚至 TB 的状态,暂停会导致严重的背压(backpressure)。
Flink 的核心创新是:将状态快照的保存与正常数据处理解耦。算子收到 Barrier 后,立即异步地将当前状态的副本写入远程存储(如 HDFS 或 S3),然后继续处理后续数据。这一机制被称为"异步屏障"。
3.2 Barrier 机制详解
Flink 的 Checkpoint Coordinator 周期性地在 source 算子注入一种特殊消息——Checkpoint Barrier。Barrier 随着数据流向下游传播,每个算子收到来自所有输入流的 Barrier 后,触发一次本地快照。
时间线示例(两个并行算子):
Source-1 ---- Barrier-N -----> Map-1 ---- Barrier-N -----> Sink-1
Source-2 ---- Barrier-N -----> Map-2 ---- Barrier-N -----> Sink-2
关键规则(Barrier Alignment):
当某个输入通道的 Barrier 先到时:
- 经典模式(Aligned):暂停处理其他输入通道的数据,将其缓存,直到所有 Barrier 都到达
- 现代模式(Unaligned):继续处理,将正在处理的数据一同存入快照,减少停顿时间
3.3 算法完整流程
Flink ABS Algorithm:
Coordinator 端(JobManager):
1. 周期 T 触发(可配置,如每 60s 一次)
2. 在所有 Source 算子注入 Barrier-N(N 为递增 ID)
3. 等待所有算子的 ACK(acknowledgment)
4. 收到所有 ACK 后,向元数据系统写入 checkpoint 完成标记
5. 如有算子在超时内未 ACK,abort 本次 checkpoint
算子端:
1. 收到来自 stream-1 的 Barrier-N
2. 如果是首次收到此 Barrier:
a. 阻塞 stream-1(不再从此通道取数据)
b. 将 stream-2,3... 的数据缓存到输入缓冲
c. 等所有输入通道的 Barrier-N 都到达
d. 异步将本地状态快照写入 StateBackend
3. 快照写入完成后,向 Coordinator 发送 ACK
4. 恢复处理,将缓存的数据写出
5. 将 Barrier 向下游转发
Flink 的源码中,CheckpointBarrierHandler 就是这一算法的核心实现类。生产环境中常用的三种模式:
- EXACTLY_ONCE:Barrier 对齐,保证精确一次语义,但会有对齐延迟
- AT_LEAST_ONCE:跳过对齐,减少延迟,但存在重复
- Unaligned Checkpoints:Flink 1.2+ 引入,将 in-flight 数据纳入快照,大幅减少对齐开销
3.4 有状态算子的快照实现
让我们看一个自定义有状态 Flink 算子的快照实现:
public class DeduplicationFunction
extends KeyedProcessFunction<String, Event, Event> {
// StateTtlConfig: 状态 TTL,避免状态无限膨胀
private transient ValueState<Boolean> seenState;
private transient ListState<Event> bufferedEvents;
@Override
public void open(Configuration parameters) {
ValueStateDescriptor<Boolean> descriptor =
new ValueStateDescriptor<>("is-seen", Boolean.class);
// 配置 TTL: 状态 24 小时未访问则自动清除
StateTtlConfig ttlConfig = StateTtlConfig
.newBuilder(Time.hours(24))
.setUpdateType(StateTtlConfig.UpdateType.OnReadAndWrite)
.setStateVisibility(StateTtlConfig.StateVisibility.NeverReturnExpired)
.build();
descriptor.enableTimeToLive(ttlConfig);
seenState = getRuntimeContext().getState(descriptor);
}
@Override
public void processElement(
Event event, Context ctx, Collector<Event> out) throws Exception {
// 去重:同一 eventId 仅输出一次
if (seenState.value() == null) {
seenState.update(true);
out.collect(event);
}
}
// 这是 Flink 快照机制的核心回调
@Override
public void snapshotState(FunctionSnapshotContext context) throws Exception {
// Flink 会自动持久化所有被声明的状态 (ValueState, ListState, MapState)
// 无需手动序列化——框架会调用 StateBackend 的序列化器
// 如果需要手动清理临时缓冲,在此处执行
}
@Override
public void initializeState(FunctionInitializationContext context) throws Exception {
// 恢复时,Flink 自动将持久化状态恢复到对应的 State 对象
// 在此处可以注册定时器等恢复后必要的初始化
}
}
四、StateBackend:快照的存储引擎
Flink 通过 StateBackend 抽象将快照的持久化与算子解耦。生产中最常见的三种:
4.1 三种 StateBackend 对比
# Flink 配置选择
state.backend: hashmap # 或 rocksdb, forst
# HashMapStateBackend (内存型)
# - 状态存储在 TaskManager JVM 堆内存
# - 速度最快,适合小状态 (< 100MB)
# - 快照增量保存到 FS
# EmbeddedRocksDBStateBackend (磁盘型)
# - 状态存储在本地 RocksDB (off-heap, 磁盘)
# - 支持超大状态 (TB 级别)
# - 默认增量快照,基于 RocksDB checkpoint
# ForStateBackend (Flink 1.16+ 新贵)
# - 基于 Forrest (共享文件 + 内存索引)
# - 更好地利用云原生存储
# - 减少 RocksDB 的写放大问题
4.2 RocksDB 增量快照的工程细节
RocksDB StateBackend 是生产环境的选择,其增量快照机制如下:
# RocksDB 增量 Checkpoint 流程抽象
class RocksDBIncrementalCheckpoint:
"""
第二次及之后的 checkpoint 流程:
1. 触发 RocksDB flush (将 MemTable -> SST files)
2. 将自上次 checkpoint 以来新增的 SST files 复制到 checkpoint 目录
3. RocksDB 自身的 WAL (Write-Ahead Log) 被截断
4. 上传 SST 文件到 DFS (HDFS/S3/OSS)
关键优化:
- SST 文件是 immutable 的,复制不会阻塞在线读写
- 增量 checkpoint 通常只新增少量 SST,速度比全量快 10-100x
- 通过硬链接避免 RocksDB compaction 触发时重复复制
"""
def save_checkpoint(self, checkpoint_id, db_path, remote_dir):
# 1. Flush memtable
rocksdb.flush(db_path)
# 2. 挑出自上次 checkpoint 以来的新 SST 文件
live_sst = self._list_live_files(db_path)
last_sst = self._get_last_checkpoint_files()
new_sst = live_sst - last_sst
# 3. 拷贝到 checkpoint 目录(利用 RefLink 避免物理拷贝)
for sst_file in new_sst:
src = os.path.join(db_path, sst_file)
dst = os.path.join(remote_dir, str(checkpoint_id), sst_file)
# 创建 RefLink(如果文件系统支持),否则硬链接
os.link(src, dst) # 或 subprocess.run(['cp', '--reflink', src, dst])
# 4. 上传元数据
self._upload_manifest(remote_dir, checkpoint_id, live_sst)
4.3 写放大:RocksDB 的隐形成本
RocksDB 的 LSM-Tree 架构存在写放大(Write Amplification)问题,在流处理场景下尤为突出:
写放大因子估算:
- Level 0: 每 1 字节写入 -> 1 字节在 L0
- Level 1: compaction 时 1 字节 -> 10-20 字节(取决于 tier size ratio)
- Level 2+: 逐级放大,最终可达 20-40x
工程缓解手段:
1. 使用 block-based table with bloom filter,减少无效读取
2. 调大 memtable size,减少 flush 频率
3. 使用 'universal compaction' 替代 'level compaction',
牺牲部分读性能换取更低写放大(适合写多读少的流处理场景)
4. FLI-ForSt (Flink 实验性) 使用 BTree 结构替代 LSM,
从根本减少写放大
五、Unaligned Checkpoints:突破对齐瓶颈
5.1 背压下的快照困境
在大背压场景(如 sink 写入外部系统变慢),Barrier 对齐会导致灾难性后果:
对齐模式的背压灾难:
Source -> Map -> [慢速Barrier等待] -> Sink(极慢)
情形:
输入通道1 Barrier 已到,开始缓存通道2、3 的数据
缓存数据会传到 Map,Map 的输入缓冲满 → 阻塞 Map
Map 阻塞 → Source 被反压 → 数据堆积
积累的 in-flight 数据越多,对齐完成后需要 flush 的缓存就越多
→ 进入恶性循环,checkpoint 耗时从毫秒级恶化到分钟级
5.2 Unaligned Checkpoints 的解法
Flink 1.2 引入的 Unaligned Checkpoints 将"等待 Barrier 对齐"变为"将 in-flight 数据也纳入快照":
def unaligned_checkpoint_flow():
"""
Unaligned Checkpoint 流程:
算子收到 Barrier-N 来自通道1:
- 不阻塞其他通道
- 立即开始快照:
a) 保存当前算子状态本身
b) 将各输入通道中已接收但未处理的数据(包括队列中的字节)
作为快照的一部分写入持久存储
- 快照传输与正常处理同时进行
这样做的好处:
- Barrier 不再被"缓存中的 pending 数据"阻塞
- 快照大小包含 in-flight 数据,但 checkpoint 时间更短
- 消除了背压对 checkpoint 时间的非线性影响
"""
pass
配置启用:
class Process:
def __init__(self, pid, channels):
self.pid = pid
self.state = {}
self.in_channels = {c.id: [] for c in channels if c.dest == self}
self.out_channels = {c.id: c for c in channels if c.source == self}
self.recorded_state = None
self.channel_recorded = set() # 已经收到marker并完成记录的通道
def initiate_snapshot(self):
"""快照发起者调用"""
self.recorded_state = self.state.copy()
for ch in self.out_channels.values():
ch.send(MarkerMessage(sender=self.pid, snapshot_id=uuid4()))
# 开始记录所有入站通道
for ch_id, msgs in self.in_channels.items():
msgs.clear()
def on_receive(self, msg):
if isinstance(msg, MarkerMessage):
self._handle_marker(msg)
else:
# 正常业务消息
self.process_message(msg)
# 如果某个通道正在记录,把该消息加入
for ch_id in self.in_channels:
if ch_id not in self.channel_recorded:
self.in_channels[ch_id].append(msg)
def _handle_marker(self, marker):
if self.recorded_state is None:
# 首次收到marker,记录自身状态
self.recorded_state = self.state.copy()
ch_id = marker.channel_id
self.channel_recorded.add(ch_id)
# 向所有出站通道广播marker
for ch in self.out_channels.values():
ch.send(MarkerMessage(sender=self.pid,
snapshot_id=marker.snapshot_id,
channel_id=ch.id))
else:
# 已经记录过,这个marker对应通道C的状态 = 临时记录的消息
ch_id = marker.channel_id
self.channel_recorded.add(ch_id)
# 通道C的状态 = 上次快照到本次marker之间收到的消息
def get_snapshot(self):
"""返回此进程贡献的全局快照片段"""
return {
'process_state': self.recorded_state,
'channel_states': {
ch_id: list(msgs)
for ch_id, msgs in self.in_channels.items()
if ch_id in self.channel_recorded
}
}0
监控指标:
checkpointAlignmentBuffered:对齐期间缓存的数据量(unaligned 模式下应为 0 或极小)checkpointDuration:快照总耗时,unaligned 模式下应稳定在秒级checkpointSize:快照大小,unaligned 模式下通常略大于 aligned(因为包含 in-flight 数据)
六、Savepoint:何时、为何以及如何
Checkpoint 是自动的、用于故障恢复的;Savepoint 是手动的、版本化的,用于作业版本升级和迁移。
class Process:
def __init__(self, pid, channels):
self.pid = pid
self.state = {}
self.in_channels = {c.id: [] for c in channels if c.dest == self}
self.out_channels = {c.id: c for c in channels if c.source == self}
self.recorded_state = None
self.channel_recorded = set() # 已经收到marker并完成记录的通道
def initiate_snapshot(self):
"""快照发起者调用"""
self.recorded_state = self.state.copy()
for ch in self.out_channels.values():
ch.send(MarkerMessage(sender=self.pid, snapshot_id=uuid4()))
# 开始记录所有入站通道
for ch_id, msgs in self.in_channels.items():
msgs.clear()
def on_receive(self, msg):
if isinstance(msg, MarkerMessage):
self._handle_marker(msg)
else:
# 正常业务消息
self.process_message(msg)
# 如果某个通道正在记录,把该消息加入
for ch_id in self.in_channels:
if ch_id not in self.channel_recorded:
self.in_channels[ch_id].append(msg)
def _handle_marker(self, marker):
if self.recorded_state is None:
# 首次收到marker,记录自身状态
self.recorded_state = self.state.copy()
ch_id = marker.channel_id
self.channel_recorded.add(ch_id)
# 向所有出站通道广播marker
for ch in self.out_channels.values():
ch.send(MarkerMessage(sender=self.pid,
snapshot_id=marker.snapshot_id,
channel_id=ch.id))
else:
# 已经记录过,这个marker对应通道C的状态 = 临时记录的消息
ch_id = marker.channel_id
self.channel_recorded.add(ch_id)
# 通道C的状态 = 上次快照到本次marker之间收到的消息
def get_snapshot(self):
"""返回此进程贡献的全局快照片段"""
return {
'process_state': self.recorded_state,
'channel_states': {
ch_id: list(msgs)
for ch_id, msgs in self.in_channels.items()
if ch_id in self.channel_recorded
}
}1
Savepoint 与 Checkpoint 的关键区别:
| 维度 | Checkpoint | Savepoint | |------|-----------|-----------| | 触发方式 | 自动(周期) | 手动 | | 生命周期 | 滚动保留(通常 3-5 个) | 永久,由用户管理 | | 格式 | StateBackend 特定 | 标准化、可移植 | | 恢复速度 | 通常更快 | 可能涉及算子状态合并 | | 跨版本兼容 | 有限 | 设计上支持(Flink 提供了工具迁移) |
生产中的最佳实践是:先用 Savepoint 停止作业,修改代码/配置后从 Savepoint 恢复,这样可以保证状态严格一致。
七、大规模生产环境的调优实战
7.1 Checkpoint 三要素调优
class Process:
def __init__(self, pid, channels):
self.pid = pid
self.state = {}
self.in_channels = {c.id: [] for c in channels if c.dest == self}
self.out_channels = {c.id: c for c in channels if c.source == self}
self.recorded_state = None
self.channel_recorded = set() # 已经收到marker并完成记录的通道
def initiate_snapshot(self):
"""快照发起者调用"""
self.recorded_state = self.state.copy()
for ch in self.out_channels.values():
ch.send(MarkerMessage(sender=self.pid, snapshot_id=uuid4()))
# 开始记录所有入站通道
for ch_id, msgs in self.in_channels.items():
msgs.clear()
def on_receive(self, msg):
if isinstance(msg, MarkerMessage):
self._handle_marker(msg)
else:
# 正常业务消息
self.process_message(msg)
# 如果某个通道正在记录,把该消息加入
for ch_id in self.in_channels:
if ch_id not in self.channel_recorded:
self.in_channels[ch_id].append(msg)
def _handle_marker(self, marker):
if self.recorded_state is None:
# 首次收到marker,记录自身状态
self.recorded_state = self.state.copy()
ch_id = marker.channel_id
self.channel_recorded.add(ch_id)
# 向所有出站通道广播marker
for ch in self.out_channels.values():
ch.send(MarkerMessage(sender=self.pid,
snapshot_id=marker.snapshot_id,
channel_id=ch.id))
else:
# 已经记录过,这个marker对应通道C的状态 = 临时记录的消息
ch_id = marker.channel_id
self.channel_recorded.add(ch_id)
# 通道C的状态 = 上次快照到本次marker之间收到的消息
def get_snapshot(self):
"""返回此进程贡献的全局快照片段"""
return {
'process_state': self.recorded_state,
'channel_states': {
ch_id: list(msgs)
for ch_id, msgs in self.in_channels.items()
if ch_id in self.channel_recorded
}
}2
7.2 RocksDB 生产调优参数
class Process:
def __init__(self, pid, channels):
self.pid = pid
self.state = {}
self.in_channels = {c.id: [] for c in channels if c.dest == self}
self.out_channels = {c.id: c for c in channels if c.source == self}
self.recorded_state = None
self.channel_recorded = set() # 已经收到marker并完成记录的通道
def initiate_snapshot(self):
"""快照发起者调用"""
self.recorded_state = self.state.copy()
for ch in self.out_channels.values():
ch.send(MarkerMessage(sender=self.pid, snapshot_id=uuid4()))
# 开始记录所有入站通道
for ch_id, msgs in self.in_channels.items():
msgs.clear()
def on_receive(self, msg):
if isinstance(msg, MarkerMessage):
self._handle_marker(msg)
else:
# 正常业务消息
self.process_message(msg)
# 如果某个通道正在记录,把该消息加入
for ch_id in self.in_channels:
if ch_id not in self.channel_recorded:
self.in_channels[ch_id].append(msg)
def _handle_marker(self, marker):
if self.recorded_state is None:
# 首次收到marker,记录自身状态
self.recorded_state = self.state.copy()
ch_id = marker.channel_id
self.channel_recorded.add(ch_id)
# 向所有出站通道广播marker
for ch in self.out_channels.values():
ch.send(MarkerMessage(sender=self.pid,
snapshot_id=marker.snapshot_id,
channel_id=ch.id))
else:
# 已经记录过,这个marker对应通道C的状态 = 临时记录的消息
ch_id = marker.channel_id
self.channel_recorded.add(ch_id)
# 通道C的状态 = 上次快照到本次marker之间收到的消息
def get_snapshot(self):
"""返回此进程贡献的全局快照片段"""
return {
'process_state': self.recorded_state,
'channel_states': {
ch_id: list(msgs)
for ch_id, msgs in self.in_channels.items()
if ch_id in self.channel_recorded
}
}3
7.3 监控告警体系
class Process:
def __init__(self, pid, channels):
self.pid = pid
self.state = {}
self.in_channels = {c.id: [] for c in channels if c.dest == self}
self.out_channels = {c.id: c for c in channels if c.source == self}
self.recorded_state = None
self.channel_recorded = set() # 已经收到marker并完成记录的通道
def initiate_snapshot(self):
"""快照发起者调用"""
self.recorded_state = self.state.copy()
for ch in self.out_channels.values():
ch.send(MarkerMessage(sender=self.pid, snapshot_id=uuid4()))
# 开始记录所有入站通道
for ch_id, msgs in self.in_channels.items():
msgs.clear()
def on_receive(self, msg):
if isinstance(msg, MarkerMessage):
self._handle_marker(msg)
else:
# 正常业务消息
self.process_message(msg)
# 如果某个通道正在记录,把该消息加入
for ch_id in self.in_channels:
if ch_id not in self.channel_recorded:
self.in_channels[ch_id].append(msg)
def _handle_marker(self, marker):
if self.recorded_state is None:
# 首次收到marker,记录自身状态
self.recorded_state = self.state.copy()
ch_id = marker.channel_id
self.channel_recorded.add(ch_id)
# 向所有出站通道广播marker
for ch in self.out_channels.values():
ch.send(MarkerMessage(sender=self.pid,
snapshot_id=marker.snapshot_id,
channel_id=ch.id))
else:
# 已经记录过,这个marker对应通道C的状态 = 临时记录的消息
ch_id = marker.channel_id
self.channel_recorded.add(ch_id)
# 通道C的状态 = 上次快照到本次marker之间收到的消息
def get_snapshot(self):
"""返回此进程贡献的全局快照片段"""
return {
'process_state': self.recorded_state,
'channel_states': {
ch_id: list(msgs)
for ch_id, msgs in self.in_channels.items()
if ch_id in self.channel_recorded
}
}4
八、云原生时代的 Checkpoint 演进
8.1 Changelog State Backend:将快照变成流
Flink 1.15 引入的 Changelog State Backend 代表了一种新思路:不再将快照视为离散时间点的产物,而是将状态的每一次变更追加写入一个 changelog,周期性地将 changelog 截断(compact)。
class Process:
def __init__(self, pid, channels):
self.pid = pid
self.state = {}
self.in_channels = {c.id: [] for c in channels if c.dest == self}
self.out_channels = {c.id: c for c in channels if c.source == self}
self.recorded_state = None
self.channel_recorded = set() # 已经收到marker并完成记录的通道
def initiate_snapshot(self):
"""快照发起者调用"""
self.recorded_state = self.state.copy()
for ch in self.out_channels.values():
ch.send(MarkerMessage(sender=self.pid, snapshot_id=uuid4()))
# 开始记录所有入站通道
for ch_id, msgs in self.in_channels.items():
msgs.clear()
def on_receive(self, msg):
if isinstance(msg, MarkerMessage):
self._handle_marker(msg)
else:
# 正常业务消息
self.process_message(msg)
# 如果某个通道正在记录,把该消息加入
for ch_id in self.in_channels:
if ch_id not in self.channel_recorded:
self.in_channels[ch_id].append(msg)
def _handle_marker(self, marker):
if self.recorded_state is None:
# 首次收到marker,记录自身状态
self.recorded_state = self.state.copy()
ch_id = marker.channel_id
self.channel_recorded.add(ch_id)
# 向所有出站通道广播marker
for ch in self.out_channels.values():
ch.send(MarkerMessage(sender=self.pid,
snapshot_id=marker.snapshot_id,
channel_id=ch.id))
else:
# 已经记录过,这个marker对应通道C的状态 = 临时记录的消息
ch_id = marker.channel_id
self.channel_recorded.add(ch_id)
# 通道C的状态 = 上次快照到本次marker之间收到的消息
def get_snapshot(self):
"""返回此进程贡献的全局快照片段"""
return {
'process_state': self.recorded_state,
'channel_states': {
ch_id: list(msgs)
for ch_id, msgs in self.in_channels.items()
if ch_id in self.channel_recorded
}
}5
8.2 与 Kubernetes 的集成
在 Kubernetes 上运行 Flink 时,容错策略需要结合 K8s 的 Pod 生命周期:
class Process:
def __init__(self, pid, channels):
self.pid = pid
self.state = {}
self.in_channels = {c.id: [] for c in channels if c.dest == self}
self.out_channels = {c.id: c for c in channels if c.source == self}
self.recorded_state = None
self.channel_recorded = set() # 已经收到marker并完成记录的通道
def initiate_snapshot(self):
"""快照发起者调用"""
self.recorded_state = self.state.copy()
for ch in self.out_channels.values():
ch.send(MarkerMessage(sender=self.pid, snapshot_id=uuid4()))
# 开始记录所有入站通道
for ch_id, msgs in self.in_channels.items():
msgs.clear()
def on_receive(self, msg):
if isinstance(msg, MarkerMessage):
self._handle_marker(msg)
else:
# 正常业务消息
self.process_message(msg)
# 如果某个通道正在记录,把该消息加入
for ch_id in self.in_channels:
if ch_id not in self.channel_recorded:
self.in_channels[ch_id].append(msg)
def _handle_marker(self, marker):
if self.recorded_state is None:
# 首次收到marker,记录自身状态
self.recorded_state = self.state.copy()
ch_id = marker.channel_id
self.channel_recorded.add(ch_id)
# 向所有出站通道广播marker
for ch in self.out_channels.values():
ch.send(MarkerMessage(sender=self.pid,
snapshot_id=marker.snapshot_id,
channel_id=ch.id))
else:
# 已经记录过,这个marker对应通道C的状态 = 临时记录的消息
ch_id = marker.channel_id
self.channel_recorded.add(ch_id)
# 通道C的状态 = 上次快照到本次marker之间收到的消息
def get_snapshot(self):
"""返回此进程贡献的全局快照片段"""
return {
'process_state': self.recorded_state,
'channel_states': {
ch_id: list(msgs)
for ch_id, msgs in self.in_channels.items()
if ch_id in self.channel_recorded
}
}6
8.3 前瞻:基于共识的 Checkpoint
最新的研究方向将分布式快照与 Raft/Paxos 等共识算法融合。典型代表是:
- Flink 的 Pulsar-based coordination:利用 Pulsar 的持久化日志作为 checkpoint 元数据的共识层
- RiseLight / 学术界的 Consensus-based Snapshot:通过 leader election 驱动的全局快照协议减少对中心 coordinator 的依赖
这些探索的目标是:将单 JobManager Coordinator 的单点故障风险消除,同时保持快照的一致性。
九、总结
分布式快照算法经历了从理论到工程的优雅演化:
- Chandy-Lamport (1985):理论基石,证明了在异步系统中获取一致性全局快照的可能性
- Lamport 的 Lazy Snapshot 优化:降低 marker 传播的延迟开销
- Flink ABS (2015-):将快照从"停止-保存-恢复"改进为"异步屏障+非阻塞"
- Unaligned Checkpoints (Flink 1.12):解决高背压场景下的快照阻塞问题
- Changelog State Backend (Flink 1.15):从离散快照到连续变更流的范式转换
- ForSt / 新型 LSM 优化:针对流处理负载重新设计存储层
对于流处理工程师而言,理解这一演进路径不仅是知识积累,更是在面对生产故障时做出正确决策的基础——当你需要在一个每秒处理百万级事件、状态达 TB 级的集群上实现秒级 RPO 时,选择哪类 StateBackend、是否开启 unaligned checkpoint、如何配置超时参数,这些决策都根植于对快照算法本质的理解。
实战建议:如果你今天正在构建一个生产级别的流处理Pipeline,请记住三个黄金原则:
1. 增量快照永远是大状态的必须——全量快照在 TB 级状态下的 IO 和延迟是不可接受的
2. 宁可调整 checkpoint 间隔,也不要关闭 checkpoint——单次允许的 checkpoint 间隔上限由你对延迟容忍度和故障时丢失窗口的 SLA 决定
3. 监控 checkpoint 耗时增长率而非绝对值——checkpoint 耗时的非线性增长往往是状态膨胀或存储 IO 瓶颈的早期信号

发表评论 取消回复