1. 为什么我坚持在项目排期里单独留出“上线检查”这一项
先说个我自己的经历。两年多前,团队第一次把 Flink 任务推到生产,当时觉得代码能跑、数据也出得来,就直接点了发布。结果上线当晚,Kafka 消费位点疯狂回退,因为 source 端的并发度和 partition 数没有对齐,部分 TaskManager 直接处于空转状态,真正干活的并行度只有预期的一半。后来通过监控发现,checkpoint 虽然一直显示成功,但实际间隔已经被拉长到十分钟以上,因为 RocksDB 状态后端的磁盘容量不够,在频繁做 compaction。那次事故让我意识到一件事:Flink 任务能不能稳定跑起来,和代码能不能跑通,完全是两套评估标准。
后来我养成了一个习惯,不管任务大小,在上生产之前,都要把一套固定的清单完整过一遍。这份清单不是从网上抄来的,是在多次线上故障、容量评估失误和版本升级事故里一点点补出来的。这篇文章就是把这份清单摊开来讲,每一项都配上我踩过的坑和判断依据。如果你正准备把一个 Flink 任务推向生产,或者你已经在生产环境里维护着 Flink 作业,这份清单里至少有一两项,值得你对照着自己的任务重新检查一遍。
先交代一下这份清单的适用范围。我尽量让它适用于大多数常见的 Flink 生产场景:流式 ETL、实时数仓同步、基于事件驱动的业务逻辑处理。如果你的场景是 Flink SQL 作业,或者用 Flink CDC 做 MySQL 到 ClickHouse 的同步,大部分条目同样适用,但在状态后端、序列化、精确一次语义这几个点上,需要做一些针对性的调整,这些我在下文中会单独提出来。
需要模型API调用? 免费领10W Token,多模型网关一键接入 Claude、DeepSeek 等主流模型。
2. 发布前七十二小时:资源、并行度与状态后端的容量账
检查清单里权重最高的部分,永远是资源与状态。很多任务上线后出问题,追溯到最后,都绕不开并行度设置不合理、状态大小预估错误、内存配置和实际负载不匹配这三件事。我不建议把资源检查放在临近发布的时间点才做,更合理的节奏是提前三天到一周,因为这一部分一旦发现问题,往往需要调整并行度或者修改状态后端的配置,随之而来的是重新压测和数据比对,时间成本比较高。
2.1 并行度不是越大越好,要和上游数据源的分区数对齐
Flink 任务的并行度,尤其是 source 端的并行度,有一个非常容易被忽略的约束:它需要和上游数据源的物理分区数对应起来。以 Kafka 为例,一个 Flink Kafka Source 的并行度,原则上不应该超过 topic 的 partition 数。我之前遇到的情况是,topic 只建了 8 个 partition,但任务配置了 12 的并行度,结果 4 个 subtask 一直处于 idle 状态,白白占着 slot,另外 8 个 subtask 承接了全部流量,单点热点问题立刻暴露出来。
再往后排查的时候,还要看 key 的分布。Kafka 的 key 如果高度集中在少数几个 key 上(比如按用户 id 哈希但头部用户流量巨大),即使 partition 数足够,这些热点 partition 对应的 subtask 也会成为瓶颈。生产上合理的做法是:先用上游 topic 的 partition 数作为并行度的基准值,然后结合峰值 TPS、单条消息大小、下游写入能力综合评估。我一般会在压测环境里测出单并行度(单个 subtask)能扛住的峰值吞吐,然后用线上预估的峰值吞吐除以单并行度能力,乘以一个 1.5 到 2 的冗余系数,这个结果就是相对安全的并行度。
2.2 RocksDB 状态后端:容量和磁盘 IO 是两笔账
如果你的任务启用了 Keyed State 或者用了 Flink SQL 的聚合、窗口、去重等能力,大概率会用到 RocksDB 状态后端。RocksDB 的内存配置比 Heap 状态后端复杂得多,不是简单地把 taskmanager.memory.managed.size 调大就行。RocksDB 实例本身会占用 JVM 堆外内存,同时在做 compaction 时需要额外的临时磁盘空间,这些成本和状态量的大小、写入频率直接相关。
有一个经验值可以参考:RocksDB 状态在磁盘上占用的空间,通常是状态逻辑大小的 3 到 6 倍。为什么会有这么大的放大系数?首先是内部多版本数据的保留,其次是 block cache 和 memtable 的缓冲,再加上 compaction 过程中的中间文件。我在生产环境吃过一次亏,当时预估状态是 10GB 左右,给每台机器挂了 60GB 的磁盘,觉得绰绰有余。结果运行了两周之后,Checkpoint 开始持续失败,日志里的报错是磁盘空间不足。后来统计发现,状态目录、checkpoint 临时目录、RocksDB 自身的 WAL 日志加在一起,峰值一度突破了 50GB。
所以评估磁盘时要同时算三块:RocksDB 的数据目录、Checkpoint 的存储目录(如果使用远程文件系统则看临时本地目录)、以及 JVM 的堆外内存和网络缓冲。每台机器的本地磁盘建议留有至少 20% 的余量,远程存储(比如 HDFS 或 S3)则要确认写入带宽能够支撑 checkpoint 周期性的全量快照。尤其要注意:如果你的 checkpoint 周期比较短、状态又比较大,每十分钟一次全量快照对带宽的冲击是很可观的。
2.3 TaskManager 内存参数别只调总量,要逐块核对
Flink 的 TaskManager 内存由多个区域组成:框架内存、任务内存、网络缓冲、托管内存、JVM 堆外内存。我见过太多只改 taskmanager.memory.process.size 的用法,内存总量倒是够了,但网络缓冲设得太小,在高吞吐场景下出现频繁的反压和 Full GC;或者框架内存和托管内存互相挤占,导致 RocksDB 的 block cache 不够,状态读写性能直线下降。
我给一个保守的配方,适合大多数流式计算任务:
- taskmanager.memory.process.size:根据单 slot 内存乘以并行度估算,但不要超过物理机内存的 80%。
- taskmanager.memory.task.off-heap.size:指向堆外,用于部分 connector 和算子,默认值通常够用,不必刻意调大。
- taskmanager.memory.network.size:如果你的任务有大量的 shuffle(比如 keyBy、窗口重新分区),这个值要适当增大。经验上,从默认值往上加 64MB 到 128MB 是比较稳妥的幅度。
- taskmanager.memory.managed.size:如果用了 RocksDB,这个值决定了每个 slot 能用的托管内存大小,直接影响 RocksDB 的 block cache 大小。推荐给到进程内存的 40% 到 50%。
这里必须强调一个容易误导的观察:很多时候你感觉任务内存不够用、频繁 GC,实际瓶颈根本不在堆内存,而是网络缓冲或者 RocksDB 的堆外内存与磁盘 IO 的相互作用。调内存之前,一定先看监控面板上的内存池使用率、GC 频率、磁盘 IO 三个指标一起的变化趋势,单看任何一项都容易误判。
3. 数据正确性校验:从序列化到精确一次语义的层层设防
很多人上线 Flink 任务时,只关心“数据能不能跑出来”,很少关心“数据跑出来是不是对的”。但生产环境的实时任务,数据正确性是最不能让步的底线。这一节我按数据在任务中流动的顺序,从序列化、时间语义到端到端一致性,挨个讲。
3.1 序列化体系:POJO、Avro、JSON 的坑各有不同
Flink 的序列化机制直接决定了状态会不会大面积膨胀、算子之间的数据传输效率高不高。生产上最常见的序列化坑有三类。
第一类,POJO 类型没有显式声明无参构造器和 getter/setter。Flink 的 POJO 序列化器要求类型是 public 的、有无参构造、字段可以通过 getter/setter 访问。一旦不满足,Flink 会降级到 Kryo 序列化。Kryo 序列化的数据体积比 POJO 序列化要大不少,而且序列化/反序列化的 CPU 开销更高。判断方法很简单,看监控里的序列化耗时,或者打开任务日志看是否出现 “The class is not a valid POJO type” 的警告。
第二类,Avro 和 JSON 的 schema 演进。如果你的上游数据模型会新增字段或修改字段类型,而 Flink 任务里用了固定 schema(比如用 Avro 的 SpecificRecord 而非 GenericRecord),那么一旦遇到新数据,整个作业可能直接报错。解决方案是尽量使用 GenericRecord + Schema Registry,或者在做 JSON 处理时保持字段的宽容度,例如用 JsonNode 而非强类型 POJO。
第三类,key 的类型不稳定。Keyed State 的 key 是序列化后存储在状态后端里的,如果同一逻辑 key 在不同时期被序列化成不同类型(比如之前是 String 后来变成 Tuple2<String, String>),那么状态会完全对不上,聚合结果也会出错。我曾经接手过一个实时指标任务,因为改动了 key 的封装方式,上线后统计数据错乱了整整一天,最后用离线数据反查才定位到这个问题。所以,key 的类型一旦定下来,生产环境中不要轻易改。
3.2 时间语义与水位线:处理时间还是事件时间,决定你看到的是哪一份数据
Flink 的时间语义不是实现层面的细节,而是业务逻辑层面的决策。处理时间简单但不准确,事件时间能还原真实发生顺序,但需要处理乱序和延迟。生产环境的实时数仓、风控和指标计算,几乎清一色要求事件时间。这里面的坑主要在 watermark 的生成策略。
常见的误区是把 watermark 延迟设得太小。如果上游 Kafka 消息的乱序是秒级的,而 watermark 延迟只给了 500ms,那么窗口触发时会有大量本该属于该窗口的数据还没到,它们会被归入下一个窗口或直接丢弃(取决于是否允许 late data)。反之,watermark 延迟设得太大,结果虽然准确了,但窗口输出的时间会明显滞后,“实时”的意义就被削弱了。
我建议的切入角度是:先统计线上数据从产生到进入 Kafka 的端到端延迟分布,取 P95 或 P99 的延迟值作为 watermark 延迟的基线,再乘以一个 1.2 到 1.5 的系数。如果你用 Flink SQL 的 GROUP BY 窗口,watermark 的本质作用没有变,但要额外注意 idle source 的处理——如果某个 source subtask 长时间没有新数据,watermark 会卡住,整个窗口的输出被无限期推迟。解决办法是给 source 设置空闲超时(idleTimeout),让空转的 source 不再阻塞 watermark 前进。
3.3 精确一次语义:Checkpoint、事务型 Sink 与重启恢复的配合
端到端的精确一次,不只是打开 checkpoint 就能做到的。它需要两件事同时成立:状态的 checkpoint 恢复是精确一次的(这部分 Flink 引擎自身可以保证),以及下游 Sink 是幂等的或事务型的(这部分取决于 connector 和表设计)。
用 Flink CDC 做 MySQL 同步到 ClickHouse,这是一个典型场景。MySQL CDC 本身可以保证读取 binlog 时基于 checkpoint 记录位点,重启后不会重复消费;但 ClickHouse 端如果用的是普通的 JDBC 批量写入,那在 Flink 任务重启或网络抖动时,就存在重复写入的可能。ClickHouse 本身不天然的强主键约束,其去重能力依赖 ReplacingMergeTree 等引擎和优化手段,因此生产上要么在 ClickHouse 表层面做幂等设计,要么接受最终一致的语义,在查询层做兜底去重。
另外要注意 checkpoint 本身不能太频繁。每秒钟做一次 checkpoint 是一种极端做法,它会让状态后端的磁盘 IO 压力变成常态。合理的频率是 30 秒到 5 分钟一次,具体看状态大小与恢复时间要求。如果状态已经有几十 GB,checkpoint 间隔太短会让每次做快照的时间远超预期,反而影响正常的流处理性能。我一般把 checkpoint 间隔和保存时间作为“恢复点目标 RPO”的平衡:接受最多丢失多少秒的数据,就对应地设置 checkpoint 间隔。
4. 上游与下游的联动风险:Kafka 位点、JDBC Connector 与 ClickHouse 写入
Flink 任务是嵌在整条数据链路中间的,上游一抖动,任务就反压;下游一抖动,任务就堆积或重试。这一节我们跳过 Flink 任务内核,专门看和上下游的对接细节。
4.1 Kafka Source 的提交策略:位点提交失败不等于消费失败
Kafka Source 的 checkpoint 机制和位点提交是两套逻辑。checkpoint 负责 Flink 内部状态的快照,而 Kafka 的 offset 提交(如果开启了 checkpoint)则是在 checkpoint 完成之后才执行的。很多人的误解是:只要 Flink 任务一直在消费数据,offset 就会持续提交。但如果你把 checkpoint 关了,Kafka 的位点提交逻辑就不受 Flink 管理了,任务重启后可能从旧的位点重新消费,造成大量重复。
生产上强烈建议开启 checkpoint,同时把 Kafka offset 的提交模式设置为 checkpoint 完成后才提交(默认行为)。还有一点容易被忽略:如果任务停止方式是普通的 cancel 而不是 savepoint 挂起,已消费但未提交的位点会丢失。所以维护任务时,能走 savepoint 就尽量走 savepoint,尤其是需要临时停机做版本升级的时候。
4.2 JDBC Connector 的批处理参数:批量大小和刷新间隔的平衡
用 Flink JDBC Connector 写关系型数据库时,两个参数至关重要:sink.buffer-flush.max-rows 和 sink.buffer-flush.interval。前者控制攒多少条再批量写入,后者控制在多长时间内强制刷新一次。
我踩过的一个坑是把 buffer-flush.max-rows 设得很大(比如 10000),但忽略了 buffer-flush.interval。在数据稀疏的时段,可能一分钟都没有攒够 10000 条,导致一批数据迟迟不写,延迟变得很大。反过来,数据密集时段,10000 条可能几秒钟就攒够了,频繁批量写入对数据库的压力又偏大。合理做法是把两个参数搭配起来:max-rows 根据单条记录大小和数据库写入能力来定,interval 则根据业务对延迟的容忍度来定。比如允许 10 秒延迟,interval 就设 10 秒,max-rows 设 5000。这样无论流量高低,都能在延迟和数据堆积之间找到平衡。
另外,JDBC Connector 重试逻辑要注意。数据库短暂不可用并不可怕,可怕的是任务在重试期间继续往内存里缓冲数据,一旦内存被撑爆,任务直接挂掉。有条件的话,建议在 Sink 前增加一个轻量的去重或限流逻辑,避免下游故障时任务自身也被拖垮。
4.3 ClickHouse 写入的常见阻塞点:批量大小、分区裁剪与 Mutations
如果你是用 Flink 往 ClickHouse 同步数据,除了前面说的幂等性问题,还有两个生产中的高频阻塞点。
第一,ClickHouse 的写入并不适合太小的批次。它的高性能建立在列式批量写入之上,如果每条都写,性能会极差。Flink 往 ClickHouse 写时,尽量走支持 HTTP 的批量插入接口,并且建议按照 ClickHouse 的分区键设计来组织写入数据的顺序,这样能避免严重的分区写入热点。比如按天分区,如果你的数据大量集中在某个时间段,写入压力会集中到个别分区 part 上,产生大量的 part 合并操作。可以通过在 Flink 侧做 keyBy + 分区键哈希来打散写入压力。
第二,ClickHouse 的 mutations(更新/删除)是重操作。如果业务上需要经常更新,而底层表又是普通的 MergeTree,更新操作会在后台产生大量数据重写,拖累查询和写入的整体性能。生产上我一般建议用 ReplacingMergeTree 配合版本字段来模拟更新,尽量把“更新”转换成“插入一条新版本数据”。同步链路里如果有一个字段能表示事件发生时间或版本号,就不需要真正执行 DELETE/UPDATE。这套思路在很多实时数仓里已经是默认方案了。
5. 部署形态的稳定性设计:K8s 环境里的 Flink 存活策略
随着 Flink 上 K8s 越来越普遍,“作业能跑”和“作业能在 K8s 里稳定跑”已经变成了两件事。K8s 的调度重试、资源抢占、网络抖动,都会影响 Flink 作业的可用性。这一节主要讲我在 K8s 上运行 Flink 的稳定性策略。
5.1 资源请求与限制不能一刀切,要给 Flink 留出“喘息空间”
K8s 的 requests 和 limits 设置,对 Flink 作业的影响非常大。如果你把 limits 设得和 requests 一样紧,TaskManager 在高峰期遇到内存抖动时,很容易被 OOMKilled。注意,Flink 的 JVM 堆外内存、网络缓冲和 RocksDB 临时空间,都会使得 JVM 实际的 RSS 占用明显高于堆内内存的配置值。你配置的 process.size 是 Flink 期望的总内存,但 K8s 看到的实际占用往往还有额外开销。所以在 K8s 环境里,我倾向于这样设置:
- requests:和 Flink 的 process.size 对齐,作为调度保障。
- limits:在 requests 基础上加 10% 到 15% 的余量,避免 JVM 因为堆外内存波动被误杀。
- JVM 参数里额外保留一部分 headroom(比如 -XX:MaxRAMPercentage 设置得比堆的期望值略高),减少 OOM 概率。
另一个常见的坑是 K8s 的节点资源碎片化。如果你的 TaskManager 需要 8GB 内存,而节点的可分配内存恰好只剩 7GB,调度器就会把 TaskManager 的调度卡住,整个作业无法扩容。解决方式有两个:一是给 Flink 的任务单独部署在专用的 NodePool 上,避免和其他应用抢占资源;二是使用合适的节点亲和性配置,尽量把 TaskManager 的 Pod 集中调度到满足内存需求的节点上。
5.2 故障恢复策略:从 Checkpoint 恢复到 Savepoint 恢复,语义完全不同
在 K8s 环境中,Pod 重启是常态。Flink 作业在 Pod 被杀后,会自动尝试从最近一次 Checkpoint 恢复。如果 Checkpoint 本身是好的,这个过程通常能跑通。但有一个前提:状态大小必须能够支撑快速恢复。如果状态达到几十 GB 甚至上百 GB,从 Checkpoint 恢复的时间可能长达几分钟甚至十几分钟,这期间任务处于空窗期,数据必然大量积压。
我建议在压测阶段就实打实地做一次“Pod 随机删除”演练,观察任务恢复需要多久,数据积压了多少。这个数字落在你业务的容忍范围内,才算真正合格。如果恢复时间超标,可以考虑的优化方向有两个:一是缩短 Checkpoint 间隔,减少每次恢复需要重放的数据量;二是从全量 Checkpoint 切换到增量 Checkpoint(RocksDB 原生支持增量 checkpoint),这能显著缩短保存和恢复的时间。
另外,如果你的运维流程里包含版本升级、SQL 逻辑变更,建议所有变更都走 Savepoint 而不是直接重启。Savepoint 是业务逻辑变更时用来恢复状态的更合适的手段,而 Checkpoint 更适合自动化故障恢复。这个区别一定要记得:线上临时恢复用 Checkpoint,版本迭代用 Savepoint。
5.3 Flink on K8s 的日志与监控,要提前定义好采集口径
K8s 环境里 Pod 是随时可以被销毁重建的,如果日志只存在 Pod 里,一旦 Pod 被重建,日志就丢了。所以上生产之前,日志必须走标准输出,配合采集器统一收集。Flink TaskManager 的日志尤其重要,反压、Checkpoint 失败、序列化异常这类问题,往往就藏在这些日志里。
监控方面,至少要把这几类指标接到你的监控平台:JVM 内存池使用率、GC 频率和耗时、Checkpoint 的 duration 和大小、Kafka 消费位点落后量(Lag)、反压率(busy 时间比例)、网络吞吐。这些指标是观察 Flink 作业健康度的六大维度。我实践下来,Lag 和反压率往往是最早暴露问题的两个信号,抖音评论区里很多“Flink 作业不稳定”的讨论,本质上就是这两个信号没被盯住,等到用户端感知到数据延迟时,问题已经发酵很久了。
6. 上线演练清单:我自己每次都逐项打勾的十二个问题
这一节,我把自己生产环境的整套检查条目拆成一个可以直接拿去用的清单。这不是理论推演,而是我在多个任务、多个集群上反复验证过的一套流程。每条我都标注了检查方式,方便你在自己的环境里复现。
6.1 启动与恢复专项检查
- 从空状态启动一次,确认启动过程中没有 UDF 初始化异常或连接器初始化失败。
- 从最近一次 Savepoint 启动一次,确认状态恢复成功,Kafka 消费位点和上游一致。
- 手动杀掉一个 TaskManager Pod,等待自动重启,记录恢复耗时和这段时间的 Kafka Lag 增量。
- 检查启动后 10 分钟内是否出现 Checkpoint 连续失败,失败原因是否为资源不足或磁盘空间不足。
- 确认任务重启后,下游数据没有出现明显的重复窗口或重复聚合(可以对比重启前后同一时间窗口的输出条数)。
6.2 数据质量与语义专项检查
- 准备一组已知结果的数据样本,在测试环境跑一遍,对比输入输出,确认聚合、过滤、维表关联逻辑正确。
- 构造乱序数据样本(延迟 1 分钟、5 分钟、30 分钟),验证 Watermark 策略是否把乱序数据归入了正确的窗口。
- 构造空输入场景(暂停发送消息持续 5 分钟以上),验证空闲 Source 不会让 Watermark 卡死、窗口不会无法触发。
- 检查所有状态清理 TTL 是否设置:没有 TTL 的无界 Keyed State 会在长时间运行后无限膨胀。
- 检查序列化报警:任务日志中出现任何 POJO 降级为 Kryo 的警告,都要在发布前解决,而不是“先这么跑”。
6.3 上下游与容量专项检查
- 核对 Kafka Topic 的 Partition 数、单条消息最大大小、预留的消费带宽。确保 Topic 扩容时 Flink Source 的并行度能够随之调整,而不是写死。
- 写下游时做一次削峰测试:人为把下游写入速度限制在正常值的一半,观察反压传导路径和任务恢复时间,判断是否有内存溢出或 Checkpoint 超时的风险。
这一套检查跑下来,快的话三四个小时,慢的话一整天。我负责任地讲,花在这上面的时间,和上线后凌晨三点爬起来处理故障的时间相比,性价比高太多。
7. 压测与故障演练的具体做法:用真实流量而非测试数据说话
在生产环境里做压测和故障演练,听起来奢侈,但恰恰是最省事的一条路。很多人只在测试环境用模拟数据跑一遍,觉得没问题就上了。但测试环境的流量模型、数据规模、机器规格和线上一比,差距往往是数量级的。模拟数据掩盖的问题,上线后第二天就还给你了。
7.1 压测数据怎么构造:三倍峰值流量是起步线
我建议的压测流量是三倍线上预估峰值。为什么是三倍?因为线上流量有毛刺,Kafka 的消费速率会因为下游抖动而积压,反压会让任务在一个时间段内承受远高于平均值的瞬时吞吐。三倍峰值流量能覆盖大多数突发情况,让任务暴露出在常规流量下看不出的问题(比如 RocksDB 的 compaction 风暴、网络缓冲耗尽、下游批量写入超时)。
压测时最好直接从 Kafka 拉取线上 Topic 的一个完整副本,不经过任何过滤地灌给任务,这样测出来的延迟和吞吐才是可信的。如果只能用模拟数据,至少要把数据量级调到线上规模的 1.5 倍以上,并且模拟真实的 Key 分布——尤其要制造热点 Key,因为真实线上永远是存在热点的。
7.2 故障演练的三个最小必要场景
故障演练不是把所有故障都演练一遍,那工作量太大也不现实。我认为有三个场景必须真实执行:
第一,Kafka 停止发送消息 10 分钟。这能检验两件事:任务会不会因为水位线卡住而导致下游输出中断,以及事件时间窗口能不能在无数据期间正常触发(如果 Watermark 配置了空闲 Source 处理,这个问题通常能解决)。
第二,下游 ClickHouse(或数据库)短暂不可用 5 分钟。停顿之后恢复连接,看任务是否会因为重试风暴导致内存暴涨,以及积压数据在恢复后能否平滑消费。
第三,随机 kill 一个 TaskManager Pod。这直接对应 K8s 里的节点重启、驱逐等真实事件。观察自动恢复时间、状态恢复的准确性、Kafka Lag 的积压趋势。如果恢复时间超过业务容忍值,就需要缩短 Checkpoint 间隔或启用增量 Checkpoint。
做完这三项演练,你对这个任务在真实故障下的表现大概心里有数了。很多团队不做演练并不是因为没有时间,而是担心演练会暴露问题、惹出麻烦。但你要想清楚:故障演练踢出的问题,顶多是让你加班一晚上;线上真实故障踢出的问题,可能是业务损失和全天候的紧急大会战。
7.3 压测通过的标准:不是“没崩”就完了
压测通过的标准一定要量化,不能是“没报错”“没崩”这种模糊描述。按照我常用的标准,至少满足以下四项:
- 在持续三倍峰值流量下,Checkpoint 连续成功率达到 100%(允许极少量的失败后自动重试成功,但不允许连续失败)。
- 端到端延迟(从 Kafka 消息时间戳到下游写入时间戳的差值)的 P95 稳定在业务要求的阈值内(比如 30 秒)。
- 反压率(BUSY 时间比例)在所有 subtask 上不超过 30%,且没有出现单一 subtask 长期 100% 的情况。
- 任务运行至少 2 小时,JVM 内存曲线平稳,Full GC 频率不超过每 10 分钟一次。
任何一条不满足,都不能算“压测通过”。很多时候你会发现,补一个并行度、调一段 Watermark 策略,比在监控面板上不停刷新等待奇迹更有效。
8. 上线后的黄金两小时与持续巡检节奏
上线不是终点,是运维长跑的开始。任务发布后的一段时间,数据链路、资源水位和状态增长速度都会经历一个从震荡到平稳的过程。这段窗口期的观察方法,决定了你能不能在上线初期就发现问题。
8.1 黄金两小时:重点盯哪几个监控指标
上线后的头两个小时,我基本不做别的事,只盯五个指标走线。
第一,Kafka Lag。刚上线时 Lag 通常比较高(因为要追历史积压),这是正常的。关键在于 Lag 是否在持续下降并收敛到接近零。如果 Lag 在半小时后仍然持平或者缓慢上升,说明消费速度跟不上生产速度,需要立刻排查。
第二,Checkpoint 间隔和大小。上线初期状态还没达到稳定的量级,Checkpoint 的间隔和大小会有一个爬坡过程。如果间隔在持续增大,比如从 1 分钟增加到 5 分钟仍不收敛,就说明 checkpoint 出现了性能瓶颈,常见原因是 RocksDB 磁盘 IO 饱和或网络带宽受限。
第三,反压率。反压可以发生在 source 端,也可以发生在 sink 端。通过逐层查看每个算子链的反压率,定位瓶颈是在消费、计算还是写出。
第四,GC 时间占比。如果 Full GC 时间占比超过 5%,需要立刻调整内存配置,这不是能靠“运行几天自动好转”解决的问题。
第五,黄金两小时内不要做任何变更。不要在刚上线的时候为了“优化”顺手调参数。任何调整,都至少等任务运行稳定四小时后再考虑。我之前遇到过,上线一小时后觉得 Lag 下降太慢,手动调低了 checkpoint 间隔,结果状态频繁快照把磁盘 IO 打满,反而引发了连锁故障。
8.2 首周巡检的重点方向:状态增长与数据对账
上线后的第一周,是验证状态模型和业务逻辑是否匹配的关键窗口。如果状态大小每天都在增长,而且没有看到明显的淘汰(由 TTL 清理或窗口关闭触发),就要警惕了。最典型的问题出在无界 Keyed State 上——如果某个 key 长期不更新、但状态里保留着大量历史数据,TTL 设置不当的话,状态会一直膨胀下去,最终拖垮整个任务。
数据对账是另一种发现问题的有效手段。做法是,选取一个有明确业务含义的指标,比如近一小时的订单金额总和或事件计数,同时从 Flink 输出结果和离线数仓计算结果分别取数,做一次比对。如果偏差持续存在,说明任务逻辑存在系统性错误,需要立刻定位。我实践下来,这种对账最好每天做一个轻量版本(全量太重),重点盯指标的变化趋势,而不是绝对值。
8.3 版本升级时最容易忽略的兼容性检查
版本升级是另一个高危场景。Flink 本身的小版本升级,通常被认为兼容性良好,但状态序列化格式、连接器 API 的行为、SQL Planner 的处理策略都可能发生细微变化。你没法保证一个在旧版本上正常跑了一年的任务,升到新版本后还能拿着旧状态的 Savepoint 无缝恢复。
所以版本升级一定要走“升级 + 从 Savepoint 恢复”的完整验证流程。先在测试集群用新版本启动任务,用线上 Savepoint 恢复状态,跑几个小时,对比计算逻辑和数据结果。完全一致之后再考虑生产变更。升级窗口尽量选在业务低峰,准备好回滚方案:一旦出现异常,立即降回旧版本并从上一个 Savepoint 恢复。
9. CDC 同步到 ClickHouse 场景的专项检查点
如果你正好是用 Flink CDC 做 MySQL 到 ClickHouse 的同步,这一节值得单独看。这个场景和通用的流式 ETL 有一些差异,尤其是 binlog 读取机制、DDL 同步和 ClickHouse 表模型的选择上。
9.1 MySQL CDC 的位点管理与表结构变更
MySQL CDC 依赖 binlog,所以任务能够读取的位点范围与 binlog 的保留时间直接相关。如果你的 MySQL 实例 binlog 只保留 24 小时,而 Flink 任务的 Checkpoint 保存了 3 天的状态,一旦任务需要从 3 天前的 Checkpoint 恢复,binlog 里已经没有对应的历史数据了,恢复必失败。
生产上建议做两件事:一是把 MySQL 的 binlog 保留时间调长,至少大于等于 Flink 任务的最大可恢复窗口(通常由 Checkpoint 保留策略决定);二是在恢复时优先使用最近一次的 Checkpoint,而不是硬着头皮恢复很早的 Checkpoint。如果你用的是全量 + 增量同步的 CDC Connector,还要额外注意全量阶段对源库的查询压力,尤其是在表数据量特别大时,需要评估锁和 IO 影响。
DDL 同步是另一个老大难。Canal 和 Flink CDC 在 DDL 事件的处理上并不完全一致,有些版本的表结构变更不会自动同步到 ClickHouse。你在 MySQL 里加了一个字段,ClickHouse 那侧如果不同步改表结构,后续数据写入会因为字段不匹配直接失败。解决方式比较朴素:建立一套 DDL 变更的审批和同步流程,每次变更先在 ClickHouse 执行对应的 ALTER,再让 CDC 任务继续跑。很多团队在这个问题上吃过亏,我也吃过,后来就写进了团队的手册里。
9.2 ClickHouse 表引擎选型:MergeTree 与 ReplacingMergeTree 的取舍
前面提到过,ClickHouse 不天然做去重,所以“MySQL 里一行数据更新后同步过来应该覆盖旧数据”这件事,ClickHouse 默认是不保证的。如果你只是用普通 MergeTree,那么 MySQL 更新一行数据时,Flink 会同步过去一行新数据,最终表里会有重复且版本不一致的多行。业务查询时如果不做任何处理,看到的就是脏数据。
实际生产里,我推荐在同步场景统一使用 ReplacingMergeTree,并且把版本字段设置为一个单调递增的值(比如 binlog 的 timestamp 或 MySQL 的 update_time)。这样 ClickHouse 在后台合并 part 时,会按版本字段保留最新一条。要注意的是,合并不是实时的,查询时如果想即时拿到最新版本,需要使用 FINAL 关键字或依赖查询时的 PREWHERE 过滤。这个语义上的“最终一致性”要在业务层面讲清楚,否则使用方很容易误以为只要写入成功就是实时可见。
9.3 两条链路的延迟基准:MySQL 到 Kafka 与 Kafka 到 ClickHouse 要分开测
很多人测“MySQL 到 ClickHouse 的同步延迟”,是把整条链路端到端放在一起测,一旦延迟超标,很难定位是上游 CDC 读取慢,还是下游 ClickHouse 写入慢。我建议把链路拆成两段来测:从 MySQL binlog 产生到该事件进入 Kafka 的延迟,以及从 Kafka 消息被消费到写入 ClickHouse 成功的延迟。
通常前者会稳定在毫秒级到秒级,后者会因为批量刷新策略出现秒级到十几秒的波动。如果两段分别都在合理范围内,但端到端延迟仍然明显偏高,那多半是中间有堆积(比如 Kafka topic 的消费者组 Lag 持续增长)。把两段指标接进同一张监控大屏,同步延迟的问题很快就能定位到具体环节。
10. 最后再分享三个我在生产环境里养成的习惯
这三条不属于清单里的任何一项,但恰恰是保障清单有效性的前提。第一个习惯是“变更必带回滚方案”。不管是调一个参数、改一段 SQL,还是升一个版本,动手之前先想清楚:出问题了怎么退回现在的状态?Flink 任务有 Savepoint 还好办,但更微小的变更(比如改 ClickHouse 表结构)同样要有预案。没有回滚方案的操作,我宁可多等半天,也不贸然执行。
第二个习惯是“监控告警规则里,永远有一项是关于 Checkpoint 失败的”。Checkpoint 失败是 Flink 任务健康度的先行指标,比延迟告警往往要早半小时暴露问题。我在生产环境把 Checkpoint 连续失败 3 次的告警设置为 P0 级别,收到的有效告警数量明显增加,但真正需要半夜起来处理的事件反而减少了,因为问题总在早期就被发现了。
第三个习惯是“每次线上事故,都要补进清单里”。这份清单不是静态的,每次踩坑之后,我会把对应的检查项补进去,并注明是哪一次事故促使我加了这一条。这样循环迭代下来,清单覆盖的盲区越来越少,团队里新同学拿到这份清单,也能很快建立对生产环境的敬畏感。
Flink 上生产这件事,本质上是在和不确定性打交道。代码能不能跑通,压测一下就知道;但能不能在凌晨三点依然稳如磐石,靠的是发布前那份不被省略的认真。希望这份清单能帮你少踩几个我踩过的坑。
