一个跑得好好的 Flink 实时作业,上线头两周 checkpoint 稳定在 30 秒上下,两个月之后变成 5 到 8 分钟,超时失败从一周一次变成一天好几次。失败之后作业从上一个成功的 checkpoint 回滚,Kafka 消费位点跟着回退,lag 从几千条滚到几百万条,下游看板数字全部飘红。团队第一反应是调大 timeout,从 10 分钟改到 20 分钟,报错消失了,但 checkpoint duration 的曲线还在缓慢往上爬——这只是把告警灯用胶布贴住了。
这篇文章把排查顺序固定下来:先看三个细分指标把根因定位到具体那条线上,再按状态后端、增量快照、barrier 对齐这三处依次处理,最后把盘和带宽按状态体积反推一遍。文中涉及的配置项名称按 Apache Flink 官方文档的公开说明写,不同版本配置项名称可能有差异,以所用版本官方文档为准。
先描述清楚典型现场,方便你对号入座。作业类型通常是订单宽表加工、风控规则匹配、实时指标聚合这类带 keyed state 的流,并行度几十到几百,输入来自 Kafka,输出到 Kafka、HBase 或者 ClickHouse 之类。上线时状态体积小,checkpoint 30 秒完成;运行两个月之后,同样的代码、同样的并行度,duration 涨到 5–8 分钟,并且开始间歇性超过 timeout。
这里有个很实际的问题:checkpoint 超时本身不是故障,而是保护机制触发了。超时后这次 checkpoint 被丢弃,作业继续跑,等下一次触发;如果连续几次都失败,作业会按重启策略从上一个成功的 checkpoint 恢复,回退的时长等于 checkpoint 间隔乘以失败次数,Kafka lag 就这样被放大。所以"调大 timeout"这个动作,本质上是允许一次 checkpoint 拖更久——拖久期间作业仍在运行,状态仍在增长,问题只会继续恶化。
真正需要区分的是三件事:状态体积是不是真的长了、barrier 在管道里是不是走不动了、远端存储是不是写不进去了。这三类根因在监控上长得完全不一样,下一节讲怎么一眼分开。
调大 timeout 之后,checkpoint 从"失败"变成"成功但很慢",告警消失,metric 里的 duration 却从 3 分钟继续爬到 6 分钟。此时你在监控面板上看到的是一条缓慢上行的曲线,而不是一堆红色失败点,问题从"必须处理"降级成"以后再说"。更麻烦的是,慢 checkpoint 会占用 checkpoint 的并发槽位,如果 max-concurrent-checkpoints 配得比较小,后续触发会被排队等待,min-pause 又被吃掉,进一步拉长实际间隔。
正确的做法是把 timeout 调回业务可接受的最大值(比如 10 分钟),同时把 duration 的三个细分指标单独画出来,看是哪一段在涨。这一节的价值就在这里:不是告诉你"慢了",而是告诉你"慢在哪一段"。
Flink 的 checkpoint duration 不是单一数字。Web UI 和 metric 系统里,一次 checkpoint 从触发到完成可以拆成几段,每段有独立的语义。把这三段和 start delay 一起看,根因基本能当场定位。
async duration 是异步阶段耗时,也就是把状态文件真正写到远端存储(对象存储、HDFS)的时间。它变长有两种可能,症状不同:
第一种是状态体积在涨。看 checkpoint 的 state size 曲线,如果两个月里从 20GB 涨到 300GB,async 变长是必然的,这时候该去查状态里的键是不是没做 TTL、是不是把不该放状态的历史数据存了下来。判断方法很直接:把 state size 和 async duration 画在同一张图上,两条曲线同步上行,就是体积问题。
第二种是存储侧吞吐波动。state size 基本平稳,async 却从 40 秒抖到 4 分钟,那大概率是对象存储的吞吐被打满或者限流了。常见诱因包括:同一账号下多个作业同时往同一个 bucket 写、bucket 有单前缀 QPS 上限、跨可用区或跨地域写。判断方法是挑一个夜间低峰手工触发一次 checkpoint,如果 duration 回到正常水平,基本可以锁定是共享带宽或限流。
alignment duration 是 barrier 对齐耗时,也就是速度快的输入通道等速度慢的通道所花的时间。这一段变长,几乎总是反压导致的:某个 sub-task 处理慢,它的输入缓冲区长期接近满,barrier 排在缓冲区排队的数据后面,要等前面所有数据都被处理完才能被读到。
判断方法:看 alignment duration 占 checkpoint duration 的比例。如果占比超过一半,说明时间主要花在等待上,而不是花在写状态上。同时看该 sub-task 的 busy time 和 back pressured time,如果 back pressure 指标长期高位,就对上了。这类问题的解法不在 checkpoint 配置里,而在算子和资源上:要么是某个算子有数据倾斜,要么是下游 sink 写不动(数据库连接池打满、批量写批次太小),要么就是 CPU 或者 GC 撑不住了。对齐超时可以先配 alignment-timeout,超时后自动退化成 unaligned,但退化本身有代价,后面单独讲。
sync duration 是同步阶段耗时,指 barrier 到达算子时同步完成的那部分动作,比如把内存里的状态句柄固化、RocksDB 的 memtable flush、本地文件的引用建立。这一段通常很短,一旦变长,往往是本地盘 IO 抖动,或者同步阶段要序列化的对象太多太大。
对 HashMapStateBackend 来说,同步阶段要把堆内对象序列化成字节写进缓冲区,状态体积越大耗时长得越明显,而且这段是在主线程上跑的,会直接阻塞数据处理,表现为处理吞吐掉一截。对 RocksDB 来说,sync 阶段一般只是刷 memtable 和建硬链接,本地盘 IOPS 不够或者盘被别的进程抢占时才会明显变长。
start delay 是从 checkpoint 触发到第一个 barrier 到达算子的时间。它变长通常说明 source 端就慢了:可能是 source 空闲(Kafka 分区无数据时不产生 barrier,Flink 有对应处理逻辑)、可能是 source 所在的 TM 负载过高、也可能是整个管道已经被反压压到源头。如果 start delay 就占了 duration 的大半,先别动状态后端,先去查反压和 source 并行度。
排序建议:先打开 start delay,再对比 alignment 占比,最后看 async 与 state size 的关系。这个顺序能在十分钟内把根因范围缩到一条线上。
Flink 1.15 之后状态后端的命名换成 HashMapStateBackend 和 EmbeddedRocksDBStateBackend,老版本里的 FsStateBackend、MemoryStateBackend、RocksDBStateBackend 属于旧的叫法。选型问题本质上只有两个变量:状态放在堆内还是堆外本地盘、能不能做增量。
HashMapStateBackend 把状态对象直接放在 JVM 堆里(准确说是放在堆上的哈希表结构里),读写走 Java 对象访问,没有 JNI 开销,也没有序列化开销,单条状态的读写延迟比 RocksDB 低一个数量级。它的边界也很干脆:状态体积必须装得下堆,而且还要给 GC、网络缓冲、用户代码留出余量。经验上状态体积超过可用堆的 60% 就要警惕,Full GC 的时间和状态体积近似成正比,几十 GB 的堆做一次 Full GC 可能停几十秒,作业直接被判定失联。
另一个硬伤是它只能做全量快照。每次 checkpoint 都要把所有状态序列化一遍再上传,状态 100GB 就意味着每次 checkpoint 都要搬 100GB。这就是为什么大状态作业没得选。
RocksDB 状态后端把状态放在每个 TaskManager 本地的 RocksDB 实例里,数据以 SST 文件形式存在本地盘上。换来三件事:状态可以超过内存;快照可以做增量(只上传上次之后新增和变化的 SST 文件);访问走本地磁盘 + block cache,堆内存压力转移到堆外。
代价也有三条。第一条是 JNI 序列化开销:每条状态读写都要在 Java 对象和字节数组之间转换,CPU 消耗明显上升,单条访问延迟高于堆内。第二条是本地盘 IOPS 要求:RocksDB 的 compaction 是后台持续行为,写放大通常在十倍到几十倍量级(通用工程估算,实际倍数取决于写模式和 compaction 策略,需按自己作业实测),本地盘随机写 IOPS 不够时 compaction 追不上写入,会触发 write stall,表现为处理吞吐突然归零几秒。第三条是运维复杂度上去了:你要关心本地盘容量、关心 compaction 线程数、关心 memtable 和 block cache 的内存划分。
说白了,"状态可以超过内存"不是免费的午餐,而是把瓶颈从内存容量换成了本地盘 IOPS 和 CPU。RocksDB 的读路径是 memtable → block cache → SST 文件,命中 block cache 时很快,没命中就要读盘。状态越大、热点越分散,block cache 命中率越低,读放大越明显。这也就是为什么大状态作业必须把 managed memory 的一大部分分给 block cache,而不是全给写缓冲。
还有一个不常被提起的点:timer(定时器)默认放堆内,状态里定时器数量很多时会把堆撑爆。RocksDB 后端下可以把 timer service 切换到 RocksDB 存储,代价是定时器读写也走 JNI,处理延迟会上升。定时器多的作业(比如大量使用窗口和 ProcessFunction 定时任务)要提前评估。
开增量很简单,state.backend.incremental 配成 true(不同版本配置项名称可能有差异,以所用版本官方文档为准)。但很多人开了之后发现"checkpoint 是快了,恢复却变慢了",这就是只算了收益没算成本。
开启增量后,每次 checkpoint 只上传自上次 checkpoint 以来新增的 SST 文件,已经上传过的文件通过元数据引用复用。假设状态 300GB,两次 checkpoint 之间变化 8GB,那么这次上传量就是 8GB 左右,而不是 300GB。这是 checkpoint duration 从 8 分钟降回几十秒的直接原因。
RocksDB 的 compaction 是本地行为,跟 checkpoint 开不开增量没关系。写入压力大的作业,compaction 会持续占用本地盘带宽和 CPU。更微妙的一点:Flink 做增量快照时依赖 SST 文件的生命周期,被 checkpoint 引用的 SST 文件不能被 compaction 删除,这会降低 compaction 的回收效率,反过来加剧磁盘占用。所以开了增量之后,本地盘的实际占用常常比状态体积本身高出不少。
增量快照的每一次只记录"增量文件 + 对历史文件的引用",恢复时要沿着这条链把需要的 SST 文件全部读回来。链越长,恢复时要下载和重建的文件越多,恢复时间越长,元数据解析也更慢。这就是为什么建议定期做一次全量基线重建:触发一次 canonical savepoint,从该 savepoint 重启,后续的增量链从新的基线开始,之前累积的引用关系被截断。具体做法是先用 savepoint 停作业,再从该 savepoint 启动,形成一个新的全量起点。同时用 state.checkpoints.num-retained 控制保留的 checkpoint 个数,别让它无限堆积。
需要提醒的是,savepoint 的格式和触发选项在不同版本间有差异,有的版本下 savepoint 也可能复用共享文件、不一定是全量自包含的。要拿它当基线,先确认你所用版本的 savepoint 语义,以所用版本官方文档为准。
如果你的痛点是"想缩短 checkpoint 间隔但每次落盘太重",可以看 changelog(增量状态变更日志)机制:把状态变更持续写到外部的 changelog 存储里,checkpoint 可以更频繁地触发而不必每次都做完整的状态物化,物化动作按独立的周期进行。它换来的是 checkpoint 间隔与状态物化解耦,代价是多一路 changelog 的写入开销和额外的存储组件依赖。是否启用要看版本支持情况和运维复杂度,不要为了上而上。
aligned checkpoint 是默认行为:算子收到某个通道的 barrier 后,把该通道后续到达的数据先缓存起来不处理,等其他通道的 barrier 都到齐,再做快照。这个"等"就是 alignment duration。反压严重时,慢通道里的 barrier 前面排着大量数据,等待时间可以从毫秒级涨到分钟级。
判断标准很明确:alignment duration 占 checkpoint duration 的比例超过一半,并且反压是持续存在的(不是偶发抖动)。这时候 unaligned checkpoint 几乎是唯一解——它允许 barrier 越过排队的数据直接被处理,把尚未处理的 in-flight 数据连同 channel state 一起写进快照,代价是快照里多了这部分数据。
另外,反压特别严重、alignment 经常超过 alignment-timeout 的情况,即使不开全局 unaligned,也建议配置 alignment-timeout,让系统在超时后自动退化成 unaligned,避免 checkpoint 无限等待。超时阈值一般按业务可接受的 checkpoint duration 反推,比如 duration 目标是 1 分钟,超时可以配 30 秒左右。
unaligned 把 in-flight 数据写进状态,这意味着状态体积直接加上"管道里堆积的数据量"。状态本身很小、但反压很重的作业,开了之后状态体积可能翻几倍甚至十几倍——反压越重,管道里堆积的数据越多,写进快照的东西越多。典型场景是下游 sink 写不动导致反压,但算子状态只有几百 MB:开了 unaligned 之后,快照里塞进了几 GB 的管道数据,每次 checkpoint 都要上传这些额外体积,恢复时也要重建它们。
所以正确的判断顺序是:先确认反压的根因能不能消除。如果能通过扩容、改 sink、解决数据倾斜把反压消掉,就先消反压,不要开 unaligned;如果反压是业务固有的(比如下游系统吞吐就是上不去),再考虑开。开完之后要立刻复查 checkpoint 的 state size 有没有明显跳变,跳变幅度过大就要回退。
这里补一个容易混淆的点。checkpoint 保证的是 Flink 内部状态的一致性,以及配合两阶段提交(2PC)sink 时的端到端精确一次。开启 unaligned 之后,快照里包含了 channel state,恢复时要把这些 in-flight 数据重新放回管道,这对 2PC sink 的事务边界提出了额外要求:恢复后事务要能从 checkpoint 记录的位置正确提交或回滚。使用支持 2PC 的 sink(如 Kafka 事务 sink)时,unaligned 是可用的,但要确保 sink 的语义配置正确;自研 sink 在开启 unaligned 前应当做一次故障注入验证,确认恢复后的事务不会重复提交或漏提交。
顺带说一下常被忽略的维度:keyed state 是按 key group 划分的,key group 总数由最大并行度决定,改并行度时状态在 sub-task 之间重新分配。并行度提高会触发状态重分配,恢复时要读取并重分布数据;RocksDB 增量恢复遇到并行度变更时,需要把相关 SST 文件下载下来再按 key group 切分,这部分开销跟状态体积成正比。也就是说,大状态作业的扩缩容流程本身就是一次重的状态迁移操作,不要指望秒级完成。
状态文件最终落到 state.checkpoints.dir 指向的路径。放在对象存储上时,写入路径是"先落到 TaskManager 本地暂存目录,再由异步线程上传"。这两个环节都要配。
上传线程数决定了并发度。线程太少,单个大文件上传慢;线程太多,会把网络出带宽和对象存储的连接数打满,反而互相拖慢,还可能触发对象存储侧的限流。经验起点是每块盘 2–4 个上传线程,再按实际带宽利用率调整。分块大小影响单个文件的上传效率和失败重试粒度:分块太小,请求数暴涨;分块太大,失败重传代价高。不同版本的具体配置项名称可能有差异,以所用版本官方文档为准,常见的是线程数相关的 state.backend.rocksdb.checkpoint.transfer.thread.num 和写缓冲、内存阈值相关的一组 state.backend.fs.* 配置。
RocksDB 的本地目录(state.backend.rocksdb.localdir)承担 compaction 的随机写,checkpoint 暂存目录承担"序列化后等待上传"的顺序写。两者如果放在同一块盘上,compaction 的随机 IO 会与暂存的顺序写互相抢 IOPS,表现为 checkpoint 的 sync 阶段抖动、compaction 延迟上升、甚至触发 RocksDB write stall。更糟的是,多个 TaskManager 进程共用一块本地盘时,各自的 compaction 会互相抢 IOPS,一个作业抖动会连累同机所有作业。
配置原则:RocksDB 本地目录与 checkpoint 暂存目录分开到不同物理盘;多个 TM 进程不要共用同一块本地盘;如果只能共用一块盘,至少要限制每个 TM 的 compaction 线程数和上传线程数,避免全量抢占。暂存盘容量按状态体积的倍数预留,具体倍数见下一节。
容器化部署时,本地盘的取舍很实际。emptyDir 用节点本地盘,性能好、成本低,但 Pod 重建后数据丢失——对 RocksDB 状态来说,Pod 重建意味着状态要从远端 checkpoint 全量恢复,代价是恢复时间变长。PVC 用网络存储,数据能跟 Pod 走,但网络存储的随机写延迟通常高于本地盘,RocksDB 的 compaction 会明显变慢。
折中方案是用 local PV 或者节点本地盘的 hostPath/emptyDir(取决于你的集群策略)承载 RocksDB 本地目录,接受 Pod 漂移时的恢复代价,同时把 checkpoint 间隔和保留策略调好,让最坏情况下的恢复时间可接受。如果业务对恢复时间极其敏感,再考虑 PVC + 高性能块存储,并把 block cache 调大来抵消一部分网络读延迟。
下面这张表把三条常见路线放在同一个坐标系里比。路线一适合状态小、迭代快的作业;路线二是大状态作业的通用解;路线三是在路线二基础上针对持续反压场景的加餐,不是默认选项。
| 对比维度 | 路线一:HashMapStateBackend + 全量快照 + 本地/HDFS | 路线二:RocksDB + 增量快照 + 对象存储 | 路线三:RocksDB + 增量快照 + unaligned | 判断依据与取舍 |
|---|---|---|---|---|
| CPU | 偏低。无 JNI 转换、无 compaction 线程,主要消耗在序列化与 GC | 偏高。JNI 序列化开销明显,compaction 常驻占用线程 | 在路线二基础上再增加 channel state 的序列化开销 | 堆内方案省 CPU,但状态一大 GC 停顿会反噬吞吐 |
| 内存 | 堆内存必须完整容纳状态,还要留 GC 与缓冲余量 | 堆压力转到堆外 managed memory,需划分写缓冲与 block cache | 同路线二,另需为 in-flight 数据预留内存与快照体积 | 状态超过可用堆 60% 就该考虑换到 RocksDB |
| 本地盘 | 几乎无要求,暂存目录写压力小 | 要求高。compaction 随机写占大头,需独立物理盘 | 同路线二,且高峰期写入更集中 | 本地盘 IOPS 决定 RocksDB 会不会 write stall |
| 网络 | 每次全量上传,出带宽需求与状态体积等比 | 只传增量,出带宽需求大幅下降 | 增量之外还要传 channel state,出带宽需求上升 | 按"窗口内必须传完"反推带宽,留三成余量 |
| 恢复时间 | 状态小时最快,直接反序列化加载 | 受增量链长度影响,链越长越慢,需定期全量重建 | 在链的基础上还要重建 channel state,恢复更慢 | 恢复时间是 SLA 的一部分,别只看 checkpoint 快慢 |
| 适用状态规模 | 几 GB 以内、状态可预估且增长可控 | 几十 GB 到数 TB,唯一可行区间 | 大状态且反压无法通过扩容消除的场景 | 先估状态上限,再选路线,不要等撑不住了才换 |
这一节给一套可执行的倒算方法。所有数字都是通用工程估算口径,不是实测跑分,落地前按自己作业的监控数据校准一遍。
RocksDB 状态后端使用的 managed memory 要在写缓冲(memtable)和 block cache 之间划分。Flink 提供 managed memory 开关和一组比例配置,常见的默认思路是写缓冲占一半左右,高速缓冲池(high-prio pool,主要给索引、过滤器等)占一成左右,剩下的留给 block cache。不同版本配置项名称可能有差异,以所用版本官方文档为准。
倒算方法:先估单个 TM 上承载的状态体积 S,block cache 期望能覆盖热数据比例 r(比如 20%–40%),那么 block cache 至少要 S × r。再由 block cache 占 managed memory 的比例反推 managed memory 总量,最后加上框架堆外、网络缓冲、JVM 元空间等固定开销,得到单 TM 的物理内存规格。状态 300GB、分散在 10 个 TM 上时,每 TM 约 30GB 状态,block cache 按 30% 配需要 9GB,反推 managed memory 约 20GB 上下,再叠加其他开销,单节点 64GB 内存会比较从容。
compaction 是 RocksDB 的常态后台行为,写放大通常在十倍到几十倍量级(通用工程估算,实际倍数取决于写入模式和 compaction 策略,需按作业实测)。倒算公式可以这么写:单 TM 的持续写入速率 W(字节/秒,等于状态更新产生的数据量),compaction 带来的实际盘写入约为 W × 放大倍数,再加上 checkpoint 暂存的顺序写。盘能提供的持续随机写带宽必须大于这个值,否则 compaction 追不上写入,触发 write stall。
容量方面,RocksDB 本地目录要按状态体积预留,同时为 compaction 过程中的临时空间留余量;checkpoint 暂存空间按状态体积的 1.5–3 倍预留,因为增量链会保留多个版本的 SST 文件,且上传完成前这些文件不能删。两者分开放在不同盘上。
这是最容易被拍脑袋的一项。假设 checkpoint 间隔 T 秒,你希望异步上传在时间窗 U 秒内完成(通常取 T 的 30%–50%,给失败重试留余量),单次增量上传的数据量为 D 字节,那么单个 TM 需要的出带宽下限是 D / U 字节/秒,换算成比特再乘 8。多个 TM 共享节点出带宽时,要按同时上传的 TM 数量累加,再乘 1.3 作为协议开销和抖动余量。
举例:checkpoint 间隔 180 秒,希望在 60 秒内传完,单次增量 8GB,则单 TM 需要 8GB/60 ≈ 136MB/s,也就是约 1.1Gbps。这个数字不小,如果节点只有 1Gbps 出带宽,就必然传不完——这往往是"async duration 一路涨"最直接的答案。要么压缩状态、要么拉长间隔、要么加带宽。
按上面的口径算完,你会发现跑 RocksDB 状态后端的 TaskManager 需要的其实是两类机器之一:一类是大本地盘 + 高随机写 IOPS 的物理机或裸金属服务器,本地盘容量和 IOPS 是硬指标,CPU 核数其次;另一类是本地盘较小但弹性更好的云主机,适合状态不大、或者用 changelog 把状态外置的架构。两者比选时,本地盘 IOPS、本地盘容量、出带宽这三项的权重远高于 CPU 主频。
在这个比选框架下,一万网络(idc10000.net)深耕 19 年(成立于 2007 年),提供华南、华东、华北、中国香港以及海外多节点的服务器与云主机方案,工程师可以 1 对 1 协助完成部署与参数配置,硬件故障支持自动迁移、7×24 中文工单响应。具体机型、本地盘规格与带宽组合需实时询价,以官网实时报价为准。
坑是什么:看到超时报错就把 execution.checkpointing.timeout 从 10 分钟改到 20 分钟甚至 30 分钟。为什么发生:超时是最显眼的报错,调大它能立刻消警,成本看似为零。怎么判断:改完之后看 duration 曲线,如果仍在持续上行,说明只是把失败转成慢成功,根因没动。怎么规避:把 timeout 定在业务可接受的上限(一般为 checkpoint 间隔的 3 到 5 倍),同时把三个细分指标做成独立告警,duration 超过目标值就报警,不等它超时。
坑是什么:增量开启后就不管了,运行半年都没做过一次基线重建。为什么发生:增量把 checkpoint 时间压下来了,团队误以为问题已经解决,恢复路径没人测过。怎么判断:看恢复耗时是否随运行时长单调变长,或者看 checkpoint 元数据目录里的文件数量是否持续累积。怎么规避:定期(比如每月或每次大版本升级前)触发一次 savepoint 并从它重启,重建全量基线;同时配好 state.checkpoints.num-retained,定期做一次真实故障恢复演练,把恢复时间纳入 SLA 考核。
坑是什么:看到 alignment duration 占比高就开 unaligned,结果 checkpoint 更慢了。为什么发生:unaligned 把管道里的 in-flight 数据写进状态,反压越重堆积越多,状态体积直接暴涨。怎么判断:开启后立即对比开启前后的 checkpoint state size,如果跳变明显(比如翻倍以上),说明管道数据被大量写进快照。怎么规避:先定位并消除反压(扩容、改 sink 批量、解决数据倾斜),确认反压无法消除再开;开之前先算一下管道堆积数据量,估算它被写进状态后的体积增幅。
坑是什么:keyed state 只增不删,运行几个月后状态体积翻十几倍。为什么发生:业务上认为"历史数据都要留着",或者直接用了 SQL 作业却没配 table.exec.state.ttl,用 DataStream API 时没配 StateTtlConfig。怎么判断:checkpoint state size 曲线单调上行且无明显回落,同时 state size 增速与输入流量成正比。怎么规避:上线前明确保留窗口并配置 TTL;SQL 作业配好空闲状态保留时间;RocksDB 下注意 TTL 清理依赖 compaction 触发,会额外增加 compaction 压力,需要相应调高本地盘 IOPS 配额和 compaction 线程资源,别把 TTL 配上就以为没有成本。
坑是什么:多个 TaskManager 或多个作业共用一块本地盘做 RocksDB 目录;或者对象存储的 region 与集群不在同一地域,走公网传输。为什么发生:机器数量有限时就近塞,或者对象存储桶建在了别的地域懒得迁。怎么判断:compaction 延迟与 checkpoint duration 在不同作业之间呈现同步抖动;或者 async duration 长期偏高且夜间低峰明显回落。怎么规避:每个 TM 独占本地盘或独占盘的 IOPS 配额;对象存储桶与计算集群放在同一地域并走内网访问;不要用公网链路承载 checkpoint 上传,公网抖动的代价是整条实时链路的恢复能力。
Q1:checkpoint 多久做一次合适?
A1:没有统一答案,要按"恢复代价 × 运行时开销"一起算。间隔太短,checkpoint 本身占用的 CPU、本地盘 IO 和网络出带宽会挤压正常数据处理;间隔太长,故障后要重放的数据量就大,Kafka lag 回追时间长。常见起点是 1 到 5 分钟:状态小、开了增量的作业可以取 1 分钟;状态大、走对象存储的作业取 3 到 5 分钟更稳。判断方法是看 checkpoint 的 async duration 是否能在间隔的一半以内完成,完成不了就说明间隔偏短或带宽不足。另外要配好 min-pause,保证两次 checkpoint 之间留出处理时间,别让快照行为连续挤压数据面。
Q2:状态后端能不能中途换?
A2:不能原地热切换。状态后端的存储格式不同,HashMapStateBackend 的堆内快照和 RocksDB 的 SST 文件互不相认,直接改配置重启会报状态不兼容。可行的路径是:先触发一次 savepoint 停掉作业,用新的状态后端配置从该 savepoint 启动;Flink 在恢复时会按新的后端重新写入状态。要注意的是,savepoint 里记录的是逻辑状态而非物理格式,所以跨后端迁移在语义上是支持的,但迁移过程本身是全量重写,耗时与状态体积成正比,大状态作业要预留足够的停机窗口。并行度保持一致可以减少重分配的额外开销。
Q3:作业改了算子之后旧的 checkpoint 还能不能用?
A3:取决于算子标识(uid)和状态是否还能对上。Flink 用 uid 把状态映射到算子,如果改动只涉及无状态算子、或者新增的算子不持有状态,且保留了原有算子的 uid,通常可以从旧 checkpoint 恢复。如果改了有状态算子的 uid、改变了 keyBy 的键、或者删除了持有状态的算子,恢复时会出现状态找不到的报错。稳妥做法是:开发阶段就给每个有状态算子显式指定 uid,不要依赖框架自动生成的标识;改动涉及状态结构时,走 savepoint 停机 + 从 savepoint 启动的路径,并在测试环境先恢复一次验证。状态序列化器不兼容时,即使 uid 对得上也会失败。
Q4:RocksDB 的本地盘到底要多大?
A4:分两块算。一块是 RocksDB 本地目录,要装得下该 TM 上的状态体积,并给 compaction 过程中的临时文件留余量,通常按状态体积的 1.5 倍左右预留比较稳。另一块是 checkpoint 暂存目录,装的是"已生成待上传"的状态文件,增量链会保留多个版本,按状态体积的 1.5 到 3 倍预留。两块要放在不同物理盘上,避免 compaction 的随机写和暂存的顺序写互相抢 IOPS。多个 TM 进程不要共用一块盘。容量之外更要紧的是 IOPS:随机写 IOPS 不够会触发 RocksDB write stall,表现为处理吞吐周期性归零,这个症状靠扩容容量是解决不了的。
Q5:checkpoint 失败会不会丢数据?
A5:单次 checkpoint 失败本身不会丢数据,作业继续运行,只是这一次快照被丢弃,下次触发时重新做。真正的数据风险来自两个方向:一是作业失败后从上一个成功的 checkpoint 恢复,那么上次成功点到故障点之间的数据会被重放,如果 sink 不是幂等或不支持两阶段提交,就会在下游产生重复数据;二是配置了不保留外部化 checkpoint 时,作业被取消后 checkpoint 目录被清理,只能从头开始消费。规避方法是:对关键作业开启外部化 checkpoint 并设置保留策略,sink 侧用幂等写或事务写,明确你能接受的是至少一次还是精确一次语义。
Q6:用 savepoint 做版本升级,和用 checkpoint 有什么区别?
A6:定位不同。savepoint 是用户主动触发、面向运维动作的一致性快照,设计目标就是支持停机升级、迁移和手动备份,格式相对稳定,跨版本兼容性更好,代价是触发和写入更重。checkpoint 是系统自动、面向故障恢复的快照,频率高、开销小,但生命周期由作业管理,作业正常取消后可能被清理(除非配置为外部化保留)。做版本升级、改并行度、改状态后端这类计划内动作,一律用 savepoint;故障自动恢复走 checkpoint。另外,不同版本下 savepoint 的格式和是否复用共享文件的行为有差异,用它当全量基线前先确认所用版本的语义,以所用版本官方文档为准。
Q7:并行度改了之后状态怎么迁移?
A7:keyed state 按 key group 组织,key group 总数由最大并行度决定,改并行度时会按新的并行度重新分配 key group。从小并行度扩到大并行度,状态被拆分到更多 sub-task;反之则合并。迁移流程是:触发 savepoint 停作业,用新的并行度从该 savepoint 启动,Flink 在恢复时完成重分配。要注意三点:最大并行度只能在作业首次启动时确定,之后不能改小,否则状态无法重分配;重分配的耗时和 IO 开销与状态体积成正比,大状态作业扩容不要期望秒级完成;RocksDB 增量快照在并行度变更时需要下载相关 SST 文件再切分,恢复时间会比同并行度恢复更长,扩容前先做一次恢复演练测出真实耗时。
Q8:对象存储做 checkpoint 目录会不会很慢?
A8:会比本地盘和 HDFS 慢,主要体现在单次请求的延迟上,但吞吐可以通过并发上传补回来。是否"慢到不可接受"取决于三个变量:单次增量上传的数据量、可用的并发上传线程数、以及节点到对象存储的出带宽。开增量之后单次上传量通常降到几 GB 量级,配合多个上传线程和足够的出带宽,async duration 控制在几十秒是可行的。真正会出问题的情况是:把对象存储桶放在与集群不同的地域走公网传输、同一账号下多个作业同时写同一个桶导致限流、或者出带宽本身就不够。规避方法是同地域部署走内网、按作业拆分路径前缀分散请求、并按"窗口内必须传完"反推出带宽。
本文涉及的 Flink 状态后端类型、checkpoint 细分指标(alignment duration、sync duration、async duration、start delay)、增量快照机制、barrier 对齐与非对齐 checkpoint 的行为说明,参考 Apache Flink 官方文档中的公开描述;配置项名称按常见版本写法给出,不同版本配置项名称可能有差异,以所用版本官方文档为准。关于写放大倍数、内存划分比例、盘容量倍数、带宽反推系数等数字,均为通用工程估算口径,用于给出量级判断和计算框架,不代表任何实测跑分,实际数值需按自己作业的监控数据校准。涉及服务器选型与资源规格的部分参考公开渠道的机型信息,具体机型、带宽与价格需实时询价,具体以签约时最新报价与合同为准。更多行业内容可访问 https://www.idc10000.net/ 。
按运维型给出条件化结论:如果作业状态在几 GB 以内且增长可控,保持堆内状态后端加全量快照,把精力放在 TTL 和并行度上,不要过早引入 RocksDB;如果状态已超过几十 GB 或者 state size 曲线持续上行,直接换到 RocksDB 并开启增量,本地盘给到独立物理盘、IOPS 优先于容量;如果 alignment duration 超过 checkpoint duration 的一半且反压无法靠扩容消除,再开 unaligned,并在开启后 24 小时内复查 state size 的跳变幅度,跳变超过一倍就回退。
三个动作必须固化成例行运维:每月做一次 savepoint 全量基线重建,截断增量链;每个季度做一次真实故障恢复演练,把恢复时间写进 SLA;checkpoint 的三个细分指标各自设告警阈值,不要等超时才看。做到这三条,checkpoint 从 30 秒滚到 8 分钟这类问题就不会再以"突然爆发"的形式出现。
一万网络(idc10000.net)深耕 19 年(成立于 2007 年),面向中小企业 IT 与运维团队提供服务器租用、云主机与算力资源方案,节点覆盖华南、华东、华北、中国香港及海外多地,支持 BGP 多线与 CN2 GIA 回国线路。针对本文场景,可以提供大本地盘、高随机写 IOPS 机型的比选建议,协助评估单节点应承载的 TaskManager 数量、本地盘划分方式以及出带宽配额,并提供 7×24 中文工单、平均 5 分钟响应、工程师 1 对 1 协助部署等支持。
如果你的 Flink 作业正在经历 checkpoint 变慢、恢复超时或者状态撑爆内存,可以把当前的并行度、状态体积、checkpoint 间隔与 duration 分段数据整理一下发给我们,由工程师按上面的倒算口径给一份资源测算。具体机型、本地盘规格、带宽与价格需实时询价,以官网实时报价与合同为准。
上一篇:表和索引一天比一天大、同一条 SQL 却越来越慢:PostgreSQL 的膨胀与 autovacuum 到底谁该背这个锅
Copyright © 2013-2020 idc10000.net. All Rights Reserved. 一万网络 科技有限公司 版权所有 深圳市科技有限公司 粤ICP备07026347号
本网站的域名注册业务代理北京新网数码信息技术有限公司的产品