给 Flink 选机器,最容易犯的错是按「CPU 核数 + 内存 G 数 + 硬盘 T 数」这套通用清单去下单。这套清单对 Web 服务、对数据库、对批处理都还算好用,但放到 Flink 上就偏了:Flink 的稳态负载里,CPU 往往不是最紧张的那一项,真正先把机器压垮的是状态——它落在哪儿、有多大、多久要被整体快照一次、快照期间磁盘还能不能腾出手来干正事。一台 32 核的机器如果挂的是机械盘、状态后端又选了 RocksDB,它的实际吞吐很可能还不如一台 16 核配 NVMe 的机器。
本文只讲硬件与配置选型这一层,不虚构任何报价。凡涉及预算的地方,一律按「预估价格」处理,并标注以咨询为准;能引用的公开标价只有一万网络官网明示的那几档,且需以官网实时价为准。下面这些判断来自实际排障里反复出现的同一批原因,顺序也是按「先看状态后端、再看 Checkpoint、再看网络与内存、最后才轮到 CPU」排的。
先把本篇的几个关键判断摆出来:
先把「状态」这个词说清楚。Flink 里的状态就是算子自己记住的东西:一个按用户 ID 做累计的算子,要记住每个用户当前的累计值;一个做窗口聚合的算子,要记住窗口里还没触发的那批数据;一个做双流 join 的算子,要记住两边各自还没匹配上的记录。这些东西不落在 Kafka 里,也不落在外部库里,它就在跑任务的这台机器的进程里(或者这块盘上)。业务跑了三个月,状态就有三个月那么大——它跟你的数据量不是一回事,跟你的吞吐也不是一回事。
状态后端(state backend)决定的就是这批东西的物理归宿。Flink 现在主要有两个选择:HashMapStateBackend 把状态当成 Java 对象放在 JVM 堆里,EmbeddedRocksDBStateBackend 把状态序列化后写进内嵌的 RocksDB,而 RocksDB 的文件就落在 TaskManager 的本地目录里。这两个选择的差别不是「快一点慢一点」,而是整台机器的资源模型会被改写成完全不同的形状。
选了 HashMapStateBackend,你的内存条有多大,状态上限就基本有多大,而且这个上限还要扣掉堆里其他开销、扣掉网络缓冲区、扣掉框架自身。选了 RocksDB,内存条不再是硬约束,但它换来的代价是:每一次状态读写都要走一次序列化加一次本地磁盘访问,于是磁盘的随机 IOPS 成了吞吐的分母。同样一份业务,换一个后端,机器的瓶颈从内存条挪到了硬盘,采购清单当然也要跟着改。
这就是为什么「先定状态后端,再定机器」这个顺序不能反。反过来的典型后果是:照着 CPU 核数和内存下单,机器到位了才发现状态后端没定,或者定了一个跟硬件不匹配的,然后上线两周开始频繁 OOM 或者 Checkpoint 超时,再回头换盘换配置——租用的机器中途换硬件的成本,比当初多花一点选对配置要高得多。这里插一句实在话:如果你的状态总量还在几个 GB 以内、业务刚起步,别急着上 RocksDB,Heap 后端跑起来省心太多;状态过百 GB 之后再谈 RocksDB,才是划算的。
还有一个经常被忽略的前提:状态大小是会涨的,而且涨的速度常常超出预估。业务方说「我们大概几十 GB」,通常是按当前一个月的量说的。等到窗口从 1 小时改成 24 小时、或者保留期从 7 天改成 30 天,状态可能是三倍五倍地涨。所以选型阶段算出来的容量,最好按「当前估算值 × 2」去下单,或者在一开始就把盘的规格留成可扩展的。
HashMapStateBackend 的工作方式非常直白:状态就是 JVM 堆里的 Java 对象,读写都在内存里完成,不落盘、不序列化。它的好处是访问延迟极低、没有 compaction 这种周期性抖动、调优参数少,跑小状态时几乎不用管。它的坏处同样直白:堆有多大,状态最多就有多大,超了就是 OutOfMemoryError,而且这个 OOM 往往发生在业务高峰期,不是在你做测试的时候。
算这笔账的时候,不能拿内存条的 G 数直接当状态容量。TaskManager 的进程内存要被切成好几块:框架堆内外各有一份固定开销,网络缓冲区要占一块,JVM 元空间要占一块,JVM 额外的堆外开销(线程栈、直接缓冲、GC 自身开销这些)也要占一块,剩下的才轮到任务堆与托管内存。也就是说,一台 64GB 内存的机器给 TaskManager 进程设了 48GB,状态能用的远不到 48GB,可能只有二十几 GB,剩下的全是必要的固定成本。
更要命的是对象本身的膨胀系数。你的逻辑状态可能只有 10GB,但一旦变成 JVM 里的 Java 对象(每个对象有对象头、有引用、有 HashMap 的桶数组和负载因子留白),实际占用翻倍是很常见的,翻到 2 到 3 倍也不稀奇。这也是为什么很多人「按 1:1 配内存」上线之后没多久就 OOM——他算的是序列化后的字节数,JVM 花的是对象形态的字节数。
那么什么情况下该用 HashMapStateBackend?我的判断口径是:单 TaskManager 上承载的状态,序列化后不超过几 GB、且未来半年看不到明显增长,就用它;窗口较短、状态 TTL 清理及时、算子不做超长周期的去重,也适合它。它的另一个优势是没有磁盘 IO 尖峰,Checkpoint 期间虽然要把状态整体刷到远端,但稳态运行时磁盘压力很小——这一点对那些盘不怎么样的机器反而是好事。
反过来,下面这些情况请直接绕开它:需要按天甚至按周保留的去重状态、双流 join 里等待匹配的长周期状态、超大窗口的聚合、以及任何「状态总量已经摸到内存条一半」的场景。有人会说可以开增量堆快照来缓解,那确实能缓解快照压力,但缓解不了堆本身的容量上限——Heap 后端的问题不是快照慢,是装不下。
换到 EmbeddedRocksDBStateBackend,状态从 JVM 堆里搬出来,序列化后写进 RocksDB 的 SST 文件,文件落在 TaskManager 的本地目录。这样一来,状态的容量上限从内存条变成了本地盘的容量,从几十 GB 一下放宽到 TB 级。听起来是彻底解决了内存问题,但代价是每一次状态访问都要付出序列化与反序列化的 CPU 成本,以及一次可能命中也可能不命中的本地磁盘读。
RocksDB 是个 LSM 树结构的存储引擎。它写入时先写 WAL(预写日志,用来保证崩溃后可恢复),再写内存里的 memtable,memtable 写满就刷成一个 SST 文件落到盘上,后台线程再定期把多层 SST 做 compaction 合并。这一套机制的顺序写性能很好,代价是两件事:一是写放大,一份逻辑数据可能被反复重写若干次;二是 compaction 是周期性触发的,触发时会在短时间内吃掉大量磁盘 IO 和 CPU,这就是那个著名的「一会快一会慢」的根源。
所以 RocksDB 后端对磁盘的要求是硬性的:必须是 NVMe。SATA SSD 在顺序读写上还行,但 RocksDB 的随机读(尤其是不命中 block cache 时去 SST 文件里翻)和 compaction 期间的混合读写,需要的是高 IOPS 与低且稳定的延迟。企业级 NVMe 的随机读延迟在几十微秒量级,机械盘是毫秒量级,差两个数量级——这不是「慢一点」,是吞吐模型整个变了。用机械盘跑 RocksDB 大状态,你会看到的现象是 Checkpoint 永远在超时边缘、反压周期性出现、业务曲线像锯齿。
RocksDB 还要吃一块内存,这块内存来自 Flink 的托管内存(managed memory)。每个 RocksDB 实例(大致对应一个算子的一个并行实例)都有自己的 memtable 与 block cache,托管内存给得越足,读命中率越高、刷盘次数越少。给少了的后果非常直观:memtable 频繁刷盘、block cache 命中率掉到很低,磁盘 IO 直接翻倍,而磁盘正是这条链路上最贵的资源。
顺带说一个调优点:RocksDB 后端建议开增量 Checkpoint。全量 Checkpoint 每次都把整个状态重传一遍,状态一上 TB,快照本身就成了负担;增量模式下只上传自上次以来变化的那部分 SST 文件,快照时间和存储占用都会大幅下降,代价是恢复时要按链条回放,且保留策略要留够份数,别把中途的某一份删了导致链条断掉。
把两个后端放在一张表里对照,硬件清单该往哪个方向偏就一目了然了。表里不列价格,因为状态后端的选择本身不产生费用,产生费用的是它倒逼出来的内存与磁盘规格;整机预算属于估算范畴,按预估价格处理,以咨询为准。
| 状态后端 | 状态存放位置 | 主要瓶颈 | 容量与配置口径 | 适合的状态规模 |
|---|---|---|---|---|
| HashMapStateBackend | JVM 堆内的 Java 对象,不落本地盘 | 堆内存容量与 GC 停顿;对象膨胀后实际占用约为序列化大小的 2–3 倍 | 按「状态 × 2~3 倍膨胀 + 网络缓冲 + 元空间 + JVM 开销」反推进程内存;磁盘只需容纳快照中转,NVMe 非必需 | 单 TM 几 GB 以内、短期看不到翻倍增长的中小状态 |
| EmbeddedRocksDBStateBackend | 序列化后写入本地盘上的 RocksDB SST 文件 | 本地盘随机 IO 与 compaction 的周期性 IO 尖峰;托管内存不足会放大刷盘 | 本地盘按「状态 ×(保留份数 + 1)+ compaction 余量」估,再留 30% 以上;托管内存按每 GB 状态 32–64MB 估;必须 NVMe,快照目录与状态目录分盘 | 几十 GB 到 TB 级、长周期去重与超长窗口的大状态 |
Checkpoint 是 Flink 容错的核心动作,简单说就是给所有算子的状态拍一张全局一致的快照,存到一个可靠的地方,出错时从最近的快照恢复。它相关的参数很多,但真正决定「能不能稳定跑」的是三个:执行间隔(interval)、两次之间的最小停顿(minPauseBetweenCheckpoints)、以及单次超时(timeout)。这三个参数设不对,机器再好也会出问题。
第一个要算清楚的关系是:间隔必须明显大于一次 Checkpoint 的实际耗时。假设你的一次快照从触发到全部算子确认完成平均要 40 秒,而间隔设的是 30 秒,那么会出现什么情况?上一次还没结束,下一次的触发条件已经满足,系统只能排队或者并发跑,快照在队列里越堆越多,恢复点永远落后于当前进度。表现出来的现象就是「快照一直在跑、完成率看着还行、但一旦故障恢复就丢十分钟数据」。我的经验口径是间隔至少留出单次耗时 2 到 3 倍的余量,并且要按峰值耗时算,不是按平均耗时算。
第二个是最小间隔,也就是 minPauseBetweenCheckpoints。这个参数的作用是:一次快照完成后,强制系统休整一段时间再允许下一次触发。它的价值在于给磁盘留一口气——快照写完,本地盘上还有 compaction 要跑、还有正常的状态读写要处理,如果立刻又来一次快照,磁盘就在两件事之间被反复撕扯。很多人只设了 interval 没设这个,结果就是快照一个接一个地压着跑,磁盘永远没有喘息窗口,超时率居高不下。
第三个是超时。超时设得太短,正常的大状态快照会被判失败,失败后从头再来,反而更慢;设得太长,一次真出问题的快照会长时间挂着占用资源。合理的做法是先实测:在真实数据量下连续观察几十次快照的耗时分布,取 P99 左右的数值,再往上留 50% 左右的余量作为超时值。别拍脑袋写 10 分钟,也别照抄别人的配置。
还有一个常被忽略的参数是并发快照数。默认值是 1,也就是同一时刻只允许一个快照在跑。有些团队为了「加快」把它调大,结果两个快照同时往磁盘和对象存储上怼,IO 压力翻倍,两个都变慢。除非你的磁盘和上游存储都确认有余量,否则保持 1 就好。另外,如果作业里有部分任务已经结束(典型如批流混合场景),「任务完成后是否继续开启 Checkpoint」这个开关也要按需要确认,默认是关闭的,某些场景反而需要打开。
最后提醒一句判读方法:看 Checkpoint 是否健康,别只看「成功率」。成功率 100% 但平均耗时从 20 秒慢慢爬到 5 分钟,这是典型的缓慢劣化,往往是状态增长了或者磁盘开始吃紧。真正该盯的是三个趋势:单次耗时曲线、对齐时间(alignment,反压会让对齐时间暴涨)、以及每次快照的持久化数据量。这三个指标放在同一个面板上看,比任何单一成功率都有用。
Checkpoint 的存储位置是个独立的决策,和状态后端不是一回事。常见两条路:一是最常见的直接写远端可靠存储(HDFS、对象存储之类,通过插件接入),二是先落到本地盘再异步上传到远端,也就是开启本地恢复(local recovery)后,TaskManager 在本地保留一份最近的快照副本。两条路的硬件含义完全不同。
直接写远端的话,本地盘的压力主要来自快照期间读取本地状态并序列化上传这一段,压力是周期性、短时的。这条路的好处是本地盘容量不用为快照预留,坏处是恢复时要从远端把整个状态拉回来,TB 级状态的恢复时间会很难看,而且恢复时所有 TM 同时拉,容易把存储和网络一起打满。
开本地恢复的话,每个 TaskManager 在自己的本地盘上保留一份最近的快照副本,故障时优先从本地加载,恢复速度能快一个量级。但代价是本地盘要为这份副本预留容量——保留几份,就要多留出几份状态大小的容量。这是一笔明确的硬件账,不是免费的性能提升。
这里必须点破一个常见误解:Checkpoint 失败最常见的原因,不是网络,不是对象存储限速,而是本地盘 IO 打满。原因是快照期间 TaskManager 要做三件事——继续处理正常流量(状态读写)、把当前状态刷成快照文件(大量读加大量写)、以及 RocksDB 后台的 compaction。这三者全部落在同一块盘上。如果这块盘同时还是系统盘、还跑了日志,那基本就是必挂。所以排查 Checkpoint 超时的第一步,应该是去看快照期间本地盘的利用率与队列深度,而不是去查带宽。
我的建议很明确:大状态作业把本地恢复打开,同时给快照单独一块盘。分盘之后,正常状态访问的 IO 和快照的 IO 物理隔离,compaction 的尖峰不会直接撞上快照的写入窗口。如果因为机型限制只能共用一块盘,那就必须把快照的并发度压到最低、把最小间隔拉长,并且把日志目录挪走——这是退而求其次的做法,不是首选。
反压(backpressure)说的是下游处理不过来,压力沿着数据链路往上传,最后把源头也拖慢。很多人看到反压的第一反应是「网络带宽不够,加机器」,这个判断十次里有八九次是错的。Flink 内部的算子间传输走的是 TaskManager 之间的网络连接,瓶颈通常不在网卡速率,而在网络缓冲区的内存量与下游的处理能力。
网络缓冲区是从 TaskManager 的网络内存池里划出来的,它决定的是「在途数据能堆多少」。缓冲区不足的表现不是带宽跑满,而是每个算子的输入池(inPoolUsage)长期接近满、输出池(outPoolUsage)也接近满,吞吐直接塌方。默认配置下网络内存大约占总 Flink 内存的百分之十,并且有上下限约束;当你的并行度高、算子链路长、单个记录的体积又大时,这个比例很可能不够,需要往上调。
判断反压位置的方法其实很直接:在监控界面上逐段看各算子的繁忙程度与输入池使用率,找到第一个「输入池满、自身又不是 CPU 吃满」的算子,瓶颈就在它或者它的下游。如果那个算子本身 CPU 已经吃满,那是算力问题;如果 CPU 空闲但输入池满,那大概率是它调用了外部系统(数据库、HTTP 接口、缓存)在同步等待,瓶颈在外部依赖;如果它下游的算子系统指标正常但这一段卡住,那才是缓冲区与网络栈的问题。
网络栈本身也有讲究。高吞吐场景下,网卡的多队列与中断分布会直接影响性能:如果所有队列的中断都压在同一个 CPU 核上,那个核会成为隐形瓶颈,表现出来就是「整机 CPU 才用了三成,吞吐就上不去了」。这种情况下需要检查中断分布是否均衡、网卡队列数是否与可用核数匹配。这类问题在租用的物理机上尤其值得提前确认,因为虚拟化层的网卡配置有时候不给你调。
所以遇到反压,正确的排查顺序是:先看有没有外部系统调用在同步阻塞,再看网络缓冲区是否吃紧,再看本地磁盘是不是被 compaction 拖住,最后才考虑加并行度或者加机器。跳过前面几步直接扩容,最常见的结果是吞吐只涨了一点点,而成本涨了一整台机器——因为真正的瓶颈还在原地。
要给 Flink 配机器,得先理解 TaskManager 的内存是怎么被切分的。这一层不清楚,下单时很容易出现「内存条很大但状态装不下」或者「内存看着够却不停 OOM」这类怪事。TaskManager 的进程内存大致分成:框架堆内外开销(固定给框架用的一小份)、任务堆、托管内存、网络内存、JVM 元空间、JVM 额外开销。这里面托管内存和网络内存是 Flink 自己管的堆外部分,其余归 JVM 管。
托管内存(managed memory)是这里面的关键一块。它是 Flink 托管给具体后端用的一块堆外内存:RocksDB 后端拿它做 memtable 与 block cache,部分算子拿它做排序与哈希表。默认大约占 Flink 总内存的四成。这个比例不是神圣不可改的——用 RocksDB 且状态大时可以往上调,用 Heap 后端且没有排序类算子时可以往下调,把内存让给任务堆。
slot 数(numberOfTaskSlots)决定的是一个 TaskManager 里能并行跑多少个任务槽。需要说清楚的是,slot 只隔离「任务」,不隔离内存:同一个 TM 里的所有 slot 共享这个 TM 进程的托管内存。所以 slot 数设大了,每个 slot 分到的托管内存就少了,而 RocksDB 是按实例(大致按 slot 粒度)各自开 memtable 与 block cache 的——slot 越多,每个实例分到的越少,刷盘越频繁,磁盘压力越大。这是一个典型的「看着免费其实很贵」的旋钮。
JVM 元空间与堆外开销是最容易被忽略的两项。元空间放的是类加载信息,规模取决于你作业里带了多少依赖——把一堆框架打进一个巨大的 fat jar 里,元空间需求会明显上升,默认 256MB 有可能不够用。JVM 额外开销覆盖线程栈、直接缓冲、GC 自身的数据结构等,默认是进程内存的一成,并且有上下限。这两项配小了的表现不是「状态装不下」,而是进程直接被系统杀掉或者抛出奇怪的本地内存错误,排查起来很费时间。
把这些加在一起,得到一个实用的下单口径:先估状态总量,反推托管内存需求,再按托管内存占比反推 Flink 总内存,最后加上框架、网络、元空间、JVM 开销这几项固定成本,才是 TaskManager 进程内存的数值;进程内存再乘以一台机器上要跑几个 TM 进程(通常 1 到 2 个),加上操作系统与其他进程留的余量(建议至少 4GB),才是内存条的规格。下面两节把这个口径展开成能直接代入的算式。
「机器有多少核,并行度就设多少」这句话在 Flink 里是错的,而且是个代价不小的错。并行度决定的是作业被切成多少个并行实例在跑,它和 CPU 核数之间应该保持一个合理比值,而不是相等。原因是每个并行实例都不是纯计算:它要维护自己的状态、要做序列化、要跑自己的那一份 RocksDB 实例(含后台 compaction 线程)、要占用自己的网络连接与缓冲区。
核数少而并行度高,后果是上下文切换频繁、每个实例分到的 CPU 时间片太碎,吞吐不升反降,延迟还会抖。核数多而并行度按核数设满,后果是实例数量爆炸:RocksDB 实例数量跟着涨,每个实例的 memtable 和 block cache 都被摊薄,compaction 线程也在抢 CPU,GC 的堆虽然小但对象 churn 更快——最后表现出来就是 GC 停顿变长、磁盘 IO 变乱、整体吞吐还不如把并行度砍一半。
我的经验口径是:先按「每 slot 一个核左右」起步,也就是并行度大致等于可用核数的一半到三分之二,留出余量给 GC 线程、网络线程、RocksDB 后台线程和操作系统;然后用压测往上调,看吞吐曲线到了哪个点开始变平,那个点就是这台机器的真实上限。别把上限当成起点。还有一点容易被忘:CPU 核数里如果有超线程,那部分核的算力不能按物理核算,按 1.2 到 1.4 倍折算比较稳妥,别当成两倍。
另外,算子之间的并行度不必整整齐齐一致。源头消费 Kafka 的那一段,并行度通常受限于 topic 的分区数——分区只有 12 个,并行度设成 24 也是浪费,其中 12 个实例空转。而下游重计算的算子(比如做复杂解析、做模型推理)可以多给一些并行度。按段设并行度比全局一刀切要省资源,代价是中间会多一次数据重分布(shuffle),这部分的开销也要算进去。
最后说 GC。Heap 后端大堆(几十 GB 级别)建议用 G1 并合理设停顿目标,别用默认配置硬扛;RocksDB 后端因为状态不在堆里,堆压力反而小,GC 问题通常不严重,主要压力挪到了堆外的托管内存与磁盘。选后端的时候,其实也顺手把 GC 调优的方向定了——这也是为什么「先定状态后端」这一步不能省。
本地盘是 Flink 机器上最容易被低估的一项。很多人按「数据量」估盘,结果状态涨起来之后盘先满了,而盘满的后果非常严重:RocksDB 写不进去、Checkpoint 直接失败、任务大规模重启。正确的估法不是按数据量,是按状态大小加快照保留策略来算。
先把本地盘拆成两个逻辑区域。第一区是 RocksDB 本地目录,放的是当前正在用的状态文件,容量大约是状态大小乘以一个写放大系数——compaction 期间新旧 SST 会同时存在,经验上按 1.5 倍估比较稳,状态更新非常频繁的场景按 2 倍估。第二区是 Checkpoint 本地目录(开了本地恢复才有),容量是状态大小乘以保留份数;如果用的是增量快照,第二份及以后只存增量,容量可以往下打折,但别打太狠,链条越长增量累积越多。
两个区算完之后,还要统一加水位余量。理由很实在:文件系统层面需要余量,compaction 的峰值占用比稳态高,而且状态会涨。我建议至少留 30%,状态增长预期明显的场景留 50%。这一条不是保守,是因为盘满了之后没有任何优雅降级的余地。
举个能直接代入的例子。假设单个 TaskManager 上承载的状态是 200GB,Checkpoint 保留 2 份,开增量快照。第一区:200GB × 1.5 = 300GB,加 30% 余量 ≈ 390GB。第二区:200GB × 2 = 400GB,增量按七折 ≈ 280GB,加 30% 余量 ≈ 364GB。两块盘分别按 480GB(或上一档)与 512GB 去选。如果因为机型限制只能共用一块盘,那就要按 300GB + 400GB = 700GB、再加 30% ≈ 910GB 去选,也就是 1TB 起步,并且要接受快照与 compaction 抢 IO 这个事实。
还有几个盘相关的细节值得写进验收清单。一是盘的类型必须是 NVMe,不接受机械盘,SATA SSD 只适合状态很小或者纯 Heap 后端的场景。二是系统盘、日志目录、状态目录、快照目录最好分开,至少系统盘要独立——日志不打理起来会非常快,一个没做轮转的日志目录几周吃满系统盘是常有的事。三是别忘了 inode 与文件数量:RocksDB 会产生大量 SST 文件,虽然单文件不小,但数量多了之后文件系统的元数据操作也会变慢,选文件系统和做监控的时候把这一项带上。
这一节给一个能直接代入的算账方法,方向是从状态大小反推机器内存。先声明口径:这是估算,用于选型阶段定配置区间,按此口径估算,实际以压测为准,不同业务的序列化密度、访问热点、compaction 压力差异很大,最终要以真实负载下的实测数据校正。
第一步,算每个并行实例分摊到的状态。设总状态 S,keyed 算子的并行度 P,那么每个实例大致是 S / P。比如总状态 600GB、并行度 20,每个实例 30GB。注意这里是「大致」:key 的分布不可能完全均匀,热点 key 会让某些实例明显偏高,估算时按平均值再乘 1.2 到 1.3 的经验系数更稳。
第二步,算托管内存。RocksDB 实例需要的托管内存用一个系数 m 表示,即每一 GB 状态配多少 MB 托管内存。经验区间是每 GB 状态 32 到 64MB,状态访问热点集中、需要高 block cache 命中率的取上限,访问分散、以写入为主的可以取下限。取中位 48MB/GB 的话,一个 30GB 状态的实例大约需要 30 × 48MB ≈ 1.44GB 托管内存。
第三步,按 slot 数放大到一个 TaskManager。设一个 TM 有 k 个 slot,则该 TM 需要的托管内存是 k × 1.44GB。k = 4 时约 5.6GB。再按托管内存占比反推 Flink 总内存:占比取默认 0.4 的话,Flink 总内存约 14GB。最后加上 JVM 元空间(默认 256MB,jar 大的往上调)、JVM 额外开销(约一成的进程内存,取 1GB 量级)、框架开销(256MB 量级),TaskManager 进程内存大致落在 16GB 上下。
第四步,反推内存条。一台机器跑 1 个 TM 进程就是 16GB 加系统余量,跑 2 个就是 32GB 加系统余量,再留 4GB 给操作系统与监控代理。按上面的数字,64GB 内存的机器跑两个 16GB 的 TM 进程是很宽裕的;如果状态继续涨到 TB 级,优先加内存条到 128GB 或 256GB,而不是把 slot 数往上堆——堆 slot 会摊薄每实例的托管内存,反而更糟。
反过来也算一遍,用来验证已有机器能扛多少状态。以一台 64GB 内存的机器为例:给 TM 进程 48GB,扣掉框架、网络、元空间、JVM 开销后,Flink 总内存约 45GB,托管内存按 0.4 算约 18GB;设 slot 8 个,每 slot 分到 2.25GB 托管内存;按 48MB/GB 系数反推,每个 slot 能承载约 46GB 状态,全机约 368GB 状态。也就是说,这台 64GB 的机器在 RocksDB 后端下大致能撑 350GB 到 400GB 的状态量(已含水位的部分不应计入)。这个数字和你按「内存条 G 数」直觉猜的量级应该不一样——因为状态根本不在内存条里。
为什么坑:RocksDB 的状态访问是随机读为主的模式,命中不了 block cache 时就要去 SST 文件里翻,机械盘的随机读延迟是毫秒量级,NVMe 是几十微秒,差两个数量级。再叠加 compaction 的混合读写,机械盘会直接把吞吐压到不可用的水平,而且表现出来是「周期性卡顿」而不是「均匀变慢」,非常难排查。
怎么避:选型阶段就把 NVMe 写进硬性要求,不接受机械盘,也不接受来路不明的「高性能云盘」而不给 IOPS 指标。上架后用 fio 做 4K 随机读写实测,绕过页缓存、跑够时间,看延迟的 P99 而不是平均值。如果预算确实紧张,宁可缩小状态(缩短 TTL、缩小窗口、把部分状态外置到 KV 存储)也不要用慢盘硬扛。
为什么坑:托管内存是每个 RocksDB 实例的 memtable 与 block cache 的来源。给少了之后,memtable 频繁刷盘产生大量小 SST、触发更多 compaction,block cache 命中率掉下来又让每一次读都落到磁盘上。结果是磁盘 IO 成倍上升,而磁盘正是这条链路上最先到顶的资源。
怎么避:按上一节的口径先算出托管内存需求,再用 managed.fraction 把比例调到匹配;如果 slot 数比较多,优先减 slot 而不是硬加内存。压测时盯 block cache 命中率与 memtable 刷盘次数这两个指标,命中率长期偏低就是托管内存不够的直接信号。
为什么坑:并行度等于核数,意味着实例数最多、RocksDB 实例最多、网络连接最多、GC 对象 churn 最快。每个实例摊到的托管内存变少,compaction 线程互相抢 CPU,GC 停顿变长。最后往往是 CPU 利用率看着上去了,吞吐却没涨,延迟反而更抖。
怎么避:从核数的一半左右起步,用压测找吞吐曲线的拐点,曲线变平的位置就是上限而不是起点。另外记得超线程不能当物理核算,按 1.2 到 1.4 倍折算;源头算子的并行度要对齐上游 topic 的分区数,超过分区数的那部分是纯浪费。
为什么坑:快照期间要干三件事:正常的状态读写、把状态刷成快照文件、RocksDB 后台 compaction。这三者全落在同一块盘上,IO 就被三件事分摊,谁都跑不快。这也是「Checkpoint 失败的头号原因是本地盘 IO 打满而不是网络」这句话的物理来源——你去查带宽、查对象存储限速,方向从一开始就错了。
怎么避:状态目录与快照目录物理分盘,各按各的容量口径估大小。分不了盘的话,就把并发快照数保持为 1、把最小间隔拉长给磁盘留出喘息窗口、把日志目录挪到别的盘,并且接受快照耗时更长这一事实。系统盘一定独立,别让日志有机会吃满它。
为什么坑:这是最典型的采购错位。吞吐上不去,第一反应是加核数,但 RocksDB 场景下瓶颈在磁盘 IO 与托管内存,加核数只是让更多线程去等待同一块盘。机器升级之后吞吐曲线几乎是平的,账单却涨了一截,然后得出「Flink 扩展性差」这个错误结论。
怎么避:升级前先定位瓶颈,方法很简单:压测期间同时看 CPU 利用率、磁盘利用率与 IO 等待、网络缓冲区使用率、GC 停顿。哪一项先到顶就升哪一项。如果磁盘先到顶,正确的动作是换更快的盘或者分盘,不是加核。
为什么坑:反压是结果不是原因。它可能来自下游调用外部系统时的同步等待、来自网络缓冲区不足、来自磁盘被 compaction 拖住、也可能只是某一个算子写得慢。不做定位直接扩容,等于把瓶颈原封不动地复制一份,成本翻倍而收益接近于零。
怎么避:按「外部依赖 → 网络缓冲区 → 本地磁盘 → 并行度」这个顺序逐项排除。找到第一个输入池满而自身 CPU 不忙的算子,从它开始查。确认瓶颈确实在算力之后再扩容,并且扩容后要重跑一遍压测确认曲线真的抬起来了。
这条线对应的是本文的主场景:状态总量在几百 GB 到 TB 级、用 EmbeddedRocksDBStateBackend、Checkpoint 频繁。它的硬件诉求非常明确——NVMe 本地盘(最好能分两块,一块给状态一块给快照)、大内存条(64GB 起步,状态大的上 128GB 或 256GB)、核数中等偏上(16 到 32 核,用来给并行度与后台线程留余量),以及无虚拟化开销的独占资源。
为什么推荐裸金属而不是云主机:Flink 对磁盘延迟和 IO 稳定性极度敏感,而虚拟化层的磁盘与网络会引入一层不可控的抖动,compaction 的尖峰撞上宿主机的邻居噪声时,你会看到无法解释的延迟毛刺。裸金属没有这一层,fio 测出来的数字就是你实际能拿到的数字。一万网络深耕 IDC 19 年(成立于 2007 年),裸金属这条线是自营机柜、最快 1 分钟上架,磁盘规格和是否支持多块 NVMe 分盘这类细节可以直接找销售确认——这一步别省,Flink 选型最怕的就是「盘是共享的」这件事下单后才发现。
运维侧的理由同样实在。跑生产的 Flink 集群最怕硬件故障后的恢复时间:7×24 中文工单、平均 5 分钟响应,硬件故障 10 分钟内自动迁移,这几条对实时任务的意义比纸面参数大——状态恢复期间整个链路是停的,早十分钟恢复就是少丢十分钟数据。另外免费的系统盘每日 3 份快照、30 秒回滚,在改配置改出问题的时候能救一次。
另一类常见形态是 Flink 负责实时特征与数据流转,链路末端接一个模型推理(实时风控打分、实时推荐、实时图像审核都属于这一类)。这时候瓶颈会分裂成两半:前半段还是典型的 Flink 状态与磁盘账,后半段则完全是 GPU 显存与算力的账。这两半的硬件诉求不一样,不要指望一台机器同时最优。
如果推理只是轻量的小模型,可以用 AI 算力云的切片形态先跑起来,成本可控,业务量涨了再换整卡;如果推理是主力负载、延迟要求严格,那就该上独立的 GPU 物理机,把 Flink 集群和推理服务分开部署,中间用消息队列解耦。分开部署的另一个好处是扩容互不影响——Flink 因为状态涨要加盘加内存的时候,不用连带把 GPU 也一起换掉。
报价这块我只说官网明示的锚点:A100 40G ¥2800/月、T4 ¥900/月,这些是官网公开档位,实际以官网实时价为准、以下单时页面显示为准。至于整台 Flink 集群加推理节点的总预算,那属于需要按配置核的估算,一律按预估价格处理,以咨询为准,本文不给具体数字。国内节点之间如果涉及跨地域的链路(比如华南跑 Flink、华东跑推理),网络侧看 BGP 多线与 CN2 GIA 回国的低延迟能力,跨境场景再单独确认。
没有一刀切的数字,但有个好用的判据:看单个 TaskManager 上承载的序列化状态量,以及它的增长预期。在几 GB 以内且看半年看不到翻倍,Heap 后端完全够用,还省掉一堆 RocksDB 调优;一旦摸到十几 GB、或者业务方明确说窗口要拉长、去重周期要延长,就该换。真正的红线是「状态量接近可用堆的一半」——到这一步不换就是赌它不涨,而状态几乎总会涨。另外记住 Heap 后端的实际堆占用是序列化大小的 2 到 3 倍,别按 1 比 1 估。
取决于一次快照要跑多久,而不是取决于你想要多细的恢复粒度。正确做法是先实测:在真实数据量下连续跑几十次,拿到耗时的 P99;间隔至少取这个数值的 2 到 3 倍,同时把最小间隔设成一个合理值(比如耗时的一半),给磁盘留出喘息窗口。如果实测一次要 40 秒,间隔设 1 分钟就是不够的,会一直追不上。反过来,如果一次只要 5 秒,1 分钟的间隔就偏保守,可以往 30 秒调。这个参数必须跟着实测数据走。
查本地盘的 IO,不要先查网络。快照期间 TaskManager 要同时处理正常读写、刷快照文件、跑 RocksDB 后台 compaction,三者挤在同一块盘上就把 IO 打满了。具体看三件事:快照期间磁盘利用率与队列深度、compaction 是否正好撞上快照窗口、日志目录是不是也在同一块盘。确认是磁盘问题之后,解法优先级是:分盘 → 拉长最小间隔 → 把日志挪走 → 换更快的盘。只有在这些都排除之后,才去查对象存储的限速和网络。
大概率不该,至少不该马上加。反压是结果,要先定位原因。按顺序查四件事:第一,下游算子有没有同步调用外部系统(数据库、HTTP 接口),这类同步等待会让 CPU 看着很闲但吞吐上不去;第二,网络缓冲区是不是吃紧,看输入池与输出池使用率;第三,本地磁盘是不是被 compaction 拖住;第四,才轮到算力是否不足。四项都指向算力瓶颈了再加机器,而且加完要重跑压测确认曲线真的抬起来,别凭感觉。
slot 只隔离任务不隔离内存,同一个 TaskManager 里所有 slot 共享那份托管内存,所以 slot 越多每个实例摊到的越少。RocksDB 场景下 slot 太多会直接导致刷盘变频繁。经验做法是按「每 slot 一个核左右」起步,也就是并行度取可用核数的一半到三分之二,然后用压测找拐点。还有一个约束:源头算子的并行度不要超过上游 topic 的分区数,超出的部分纯属空转。记住超线程不能按物理核算,按 1.2 到 1.4 倍折算。
取决于状态规模与访问模式。状态很小(几 GB)、或者用 Heap 后端、或者读写以顺序为主,SATA SSD 勉强能用。但只要上了 RocksDB 且状态到几十 GB 以上,就不太行——随机读延迟和 compaction 期间的混合 IO 是 SSD 与 NVMe 拉开差距的地方,差的是一个数量级,不是百分比。如果预算卡得很死,我的建议是宁可先缩小状态规模(缩短 TTL、缩小窗口、把冷状态外置),也不要用慢盘硬扛大状态。
把 Flink 和推理分开部署,中间用消息队列解耦,不要挤在同一台机器上——前者的瓶颈在磁盘与托管内存,后者在显存与算力,混在一起两边都配不好。轻量小模型先用算力云的切片形态起步,成本可控;延迟要求严格或者是主力负载,就上独立 GPU 物理机。官网明示的锚点价是 A100 40G ¥2800/月、T4 ¥900/月,以官网实时价为准;整套集群的总预算属于估算范畴,按预估价格处理,以咨询为准。
两个信号。一是磁盘延迟开始出现无法解释的毛刺——虚拟化层的邻居噪声会和 RocksDB 的 compaction 叠加,导致吞吐曲线出现不可复现的抖动,这时候 fio 测出来的数字和你实际拿到的数字已经不是一回事了。二是状态规模进入 TB 级、需要多块 NVMe 分盘,而云主机的盘规格给不到。除此之外,如果任务只是中等规模、状态几十 GB,云主机的弹性反而更合适,按量扩缩容比提前买断划算。
本文不涉及任何虚构报价。文中出现的机器选型判断、容量估算口径与参数关系,来自 Flink 官方文档对状态后端与内存模型的定义、RocksDB 的 LSM 工作机制,以及实际排障中反复出现的同一批故障模式;容量与内存部分给出的是可代入的估算方法,按此口径估算,实际以压测为准。涉及预算的部分一律按预估价格处理,以咨询为准,不给出具体成交数字。
文中提到的公开标价(A100 40G ¥2800/月、T4 ¥900/月等)来自一万网络官网公示的产品页,属官网明示档位,需以官网实时价为准,实际以下单时页面显示为准;产品与配置细节以 https://www.idc10000.net/ 对应页面为准,具体以签约时最新报价与合同为准。服务承诺部分只写官网明示的内容:7×24 中文工单、平均 5 分钟响应、硬件故障 10 分钟内自动迁移、免费系统盘每日 3 份快照与 30 秒回滚、免费备案协助、5–20G 免费 DDoS 防护、自营机柜最快 1 分钟上架、BGP 多线与 CN2 GIA 回国。
给 Flink 选机器,把顺序定对了就成功一大半:先定状态后端,再算本地盘,再算托管内存与网络缓冲区,最后才轮到 CPU 核数。按这个顺序走下来,你会发现自己最终下单的配置,跟一开始凭直觉写出来的那张清单差别不小——通常核数更保守,而盘和内存大了一档。这不是保守,是因为 Flink 的钱本来就该花在这两处。
另外一个建议:别在选型阶段就追求一步到位。先用压测把这台机器的吞吐上限和状态上限测出来,把数字留档,后面状态涨了、要扩机器的时候,这份基线就是最有说服力的依据。Flink 的硬件账是动态的,三个月前合适的配置,三个月后可能就要调——把它当成一件需要定期复算的事,而不是一次性的采购动作。
Copyright © 2013-2020 idc10000.net. All Rights Reserved. 一万网络 科技有限公司 版权所有 深圳市科技有限公司 粤ICP备07026347号
本网站的域名注册业务代理北京新网数码信息技术有限公司的产品