KEEL · 龙骨 · A CURRICULUM FOR THE AI ERA
06 · 容量、积压与运维 — keel 龙骨
前面五章讲的是机制,这一章把这些机制收口成四件运维上真要做的判断:分区数怎么定、积压能不能追回来、盯哪些指标、什么时候该承认这台架构撑不住了。所有数字都来自本机实测。
前面五章讲的是机制,这一章把这些机制收口成四件运维上真要做的判断:分区数怎么定、积压能不能追回来、盯哪些指标、什么时候该承认这台架构撑不住了。所有数字都来自本机实测。
现场
故事线收尾:一批订单消息在 Kafka 里,消费端因为一次发布回滚变慢了,lag 开始涨。运维群里两种声音——「加消费者」和「再等等」。这一章给出算这个决策需要的数据。
积压处置的决策路径
flowchart TD
A["① 监控报警: 某消费组 lag 上涨"] --> B{"② 看分区级 lag<br/>是单分区还是全部?"}
B -- "单个分区" --> C["③ 查 key 倾斜或该分区消费者卡住"]
B -- "全部分区" --> D{"④ 消费者日志有异常/停顿吗?"}
D -- "有" --> E["⑤ 查单条处理耗时 / 重平衡频率<br/>是否超过 max.poll.interval.ms"]
D -- "没有" --> F{"⑥ 净速度 = 消费速率 - 生产速率"}
F -- "> 0" --> G["⑦ 会自行追平<br/>写出预计追平时间, 继续观察"]
F -- "<= 0" --> H{"⑧ 消费者数 < 分区数?"}
H -- "是" --> I["⑨ 加消费者, 最多加到分区数"]
H -- "否" --> J["⑩ 分区数已成瓶颈<br/>扩分区 / 扩集群 / 限流上游"]
C -. "确认是倾斜就拆 key 或改分区器" .-> J
style A fill:#fff4e6,stroke:#e8590c,color:#1a1a2e
style B fill:#fff9db,stroke:#f08c00,color:#1a1a2e
style D fill:#fff9db,stroke:#f08c00,color:#1a1a2e
style F fill:#fff9db,stroke:#f08c00,color:#1a1a2e
style H fill:#fff9db,stroke:#f08c00,color:#1a1a2e
style C fill:#e8f0fe,stroke:#3b5bdb,color:#1a1a2e
style E fill:#e8f0fe,stroke:#3b5bdb,color:#1a1a2e
style G fill:#e6fcf5,stroke:#0ca678,color:#1a1a2e
style I fill:#e6fcf5,stroke:#0ca678,color:#1a1a2e
style J fill:#ffe3e3,stroke:#c92a2a,color:#1a1a2e
路径 ⑥ 是全章的核心:积压不是「有多少条」,是「净速度正不正」。很多人卡在「1 亿条太多了」的直觉上,其实只要消费速率大于生产速率,它就是个时间问题。
实测吞吐与延迟
ProducerPerformance,200000 条 × 200 字节,acks=all,三副本:
不限速: 200000 records sent, 119260.6 records/sec (22.75 MB/sec),
627.05 ms avg latency, 1036.00 ms max latency,
640 ms 50th, 1004 ms 95th, 1024 ms 99th, 1035 ms 99.9th
限速5万: 200000 records sent, 49813.2 records/sec (9.50 MB/sec),
5.03 ms avg latency, 225.00 ms max latency,
3 ms 50th, 27 ms 95th, 41 ms 99th, 51 ms 99.9th
同一台机器、同一份配置,两次结果差得很远。不限速那次的高延迟不是 Kafka 慢,是客户端排队:生产速率超过 broker 消化速率,批全堆在缓冲区里等,延迟自然被拉长。限速档才是这台机器在「可持续速率」下的真实延迟画像——p99 41 ms。
这件事有个直接推论:「吞吐」和「延迟」必须分开测。任何人给你一个「Kafka 单分区能到 xx 万 TPS」的数字,你都要问清楚是在什么延迟约束下测的。同样一台本机,不计延迟能到 11.9 万条/秒,要求 p99 < 50 ms 就只能到约 5 万条/秒。
消费侧 ConsumerPerformance(单消费者读回 20 万条):
data.consumed.in.MB MB.sec data.consumed.in.nMsg nMsg.sec rebalance.time.ms fetch.time.ms
38.1813 10.8965 200180 57128.99 3268 236
单消费者稳态约 5.7 万条/秒(10.9 MB/秒)。注意 rebalance.time.ms = 3268:真正 fetch 只花了 236 ms,3.27 秒花在加入组上(第 04 章讲过)。一个消费者刚启动时,「吞吐」这个指标是不准的,启动开销会稀释掉它。
积压 1 亿条要追多久
用上面的实测值算(脚本输出见 09-throughput.txt):
生产(不限速) 119261 条/秒 = 22.75 MB/秒
生产(限速档) 49813 条/秒 = 9.50 MB/秒
消费(单消费者) 57129 条/秒 = 10.90 MB/秒
场景A 停止写入, 单消费者追 1 亿条: 1750 秒 = 29.2 分钟
场景B 边写边追(生产 5万/秒), 净追平 7316 条/秒: 13669 秒 = 3.80 小时
场景C 不限速生产继续写, 净速度 -62132 条/秒 -> 追不上, 积压只会涨
三个场景对应三种完全不同的处置:
- A:上游已经停了(发布回滚、上游宕机)。29 分钟能追平,加消费者只会让下游数据库更难受,这时候的正确动作是等。
- B:上游还在写,但比你消费得慢,净速度 +7316 条/秒。3.8 小时能追平。可以加消费者把这个时间压短,也可以不加。
- C:上游的生产速率高于你单消费者的消费速率。这种情况加一个消费者也要先算清楚加几个——本机一个消费者 5.7 万/秒,要压过 11.9 万/秒的生产至少得 3 个消费者(而且受分区数限制,第 04 章)。
注意这些数字是「本机单消费者」的量级,真实集群的绝对值会差很远(磁盘、网络、消息大小都不同),但算法不变:先测出消费速率和生产速率,再算净速度。不要用别人的数字替你做决定。
分区数怎么定
分区数是这门课里最难给「标准答案」的参数,因为它同时被四件事拉扯:
- 并行上限:一个消费组的最大并行度就是分区数。分区数 = 3 意味着最多 3 个消费者在干活。定分区数时,要按「未来一年内消费者实例的最大数量」来估,不是按当前。
- 单分区顺序:分区内才有顺序保证。把一个 key 的消息拆到多个分区,就等于放弃了这个 key 的顺序(第 03 章)。需要严格顺序的业务实体数量,是分区数的下界——比如「一个商户一个分区」时,分区数不能少于商户数。
- 元数据与文件开销:每个分区在每台持有副本的 broker 上都是一组文件。本机实验里 6 MB 数据配 1 MB 段就产生了 6 个段、18 个索引/快照文件。分区数 × 副本数 × 段数就是文件总数,分区过多的集群 broker 启动会很慢。
- 故障恢复时间:一个分区只能有一个 leader,broker 挂掉后它上面所有分区的 leader 都要重选,重选要并发做。分区太多时故障恢复会变慢。
一个可操作的起点:分区数取「峰值消费者实例数」和「需要保序的 key 基数」里的较大者,再留一点余量;然后压测验证。不要一上来给 1000 个分区——本机这个规模下,3 个分区就能打满这台机器的 IO。
分区数能增不能减(减少要重灌 topic),这是定分区数时要「宁可估大」的原因。
要盯的指标
按「谁出了问题」分组,比按组件分组更好用:
| 症状 | 该看的指标 | 具体含义 |
|---|---|---|
| 消费变慢 | 分区级 records-lag-max、poll-interval |
滞后是不是只压在个别分区 |
| 写入变慢 | 生产者 record-retry-rate、request-latency-avg |
重试率上升往往先于失败 |
| 副本异常 | UnderReplicatedPartitions、OfflinePartitionsCount |
前者是「副本数不够」,后者是「分区没 leader」 |
| 磁盘将满 | 各 log.dirs 可用空间、单分区大小 |
写满会直接让 log dir 下线 |
| 重平衡频繁 | rebalance-rate、rebalance-latency-avg |
频繁重平衡通常指向消费者处理超时 |
| 事务堆积 | LSO 与 LEO 的距离、未完成事务数 | 距离长期不为 0 说明有长事务卡着 |
UnderReplicatedPartitions 和 OfflinePartitionsCount 的区别值得背下来:前者是「副本还在,只是某个副本落下了」(第 02 章的 ISR 2/3),业务通常无感;后者是「Leader: none,这个分区读写全停」(第 02 章的 orders-06),是真正的可用性事故。
一次把三台机器上的坑都踩了一遍
本机是 Windows,索引文件用 mmap 打开。这个组合暴露了三件同源的事故,全部真跑到:
删 topic(README 里警告过,这次没做):目录重命名失败 → log dir offline → broker 自杀。
缩副本的 reassign:把
orders-06的副本从[1,2,3]改成[3],a1/a2 需要删本地副本目录,报java.nio.file.AccessDeniedException: ...\orders-06-0 -> ...\orders-06-0.8a1f910b...-delete Shutdown broker because all log dirs in ...\kafka-data\a1 have faileda1、a2 双双退出,controller 仲裁只剩 a3,元数据一起冻结。
retention 删段:
orders-12的一个段过期,删除任务先删.log、再改名.index,改名失败:Error while deleting segments for orders-12-0 in dir ...\kafka-data\a1 java.nio.file.FileSystemException: ...\00000000000000000000.index -> ...\00000000000000000000.index.deleted: 另一个程序正在使用此文件,进程无法访问。 Uncaught exception in scheduled task 'kafka-log-retention' Shutdown broker because all log dirs in ...\kafka-data\a1 have failed删除只做了一半(
.log变成了.log.deleted、.index还在),a1 又被带走。重启后 earliest offset 确实从 0 前进到了 2165,删除生效了,代价是没了一个 broker。
这三条是同一个机制:Windows 不允许重命名/删除被 mmap 打开的索引文件。前两条我在第 02、05 章已经用过它们的后果(Leader: none、仲裁丢失);第三条是这章要强调的——retention 不是「设了就没事」,它是一条会主动改文件系统的运维路径。
Linux 上这个问题不存在(unlink 打开的文件是允许的)。所以如果你的生产是 Linux,这三条不会遇到;但如果开发机是 Windows,本地跑 Kafka 复现 retention 行为时,别用短 retention 去等自动删除,否则 broker 会莫名其妙挂掉——这次的现场就是这样来的。
retention 什么时候真的生效
还有一个不直观的点:orders-12 的 retention.ms=15000,但数据写完之后等了整整 5 分钟才开始删。原因是删除动作由 broker 的定时任务 log.retention.check.interval.ms(默认 300000 ms = 5 分钟)驱动,它跟 topic 的 retention.ms 是两个层级:
retention.ms(topic 级)说「数据保留多久」;log.retention.check.interval.ms(broker 级)说「多久检查一次要不要删」。
所以实际删除时间 = retention.ms + 最多一个检查周期。设了 15 秒的 retention,最坏情况下要等 5 分钟才删。把检查周期调小会更频繁地扫目录,大集群上是要付代价的,通常不建议低于 1~2 分钟。
另外,正在被未完成事务引用的段不会删(deleteHorizonMs,第 05 章)。一个卡住的长事务能拖住一个分区末尾段的清理。
什么时候该换架构
Kafka 不是所有「消息」问题的答案。出现下面这些信号时,要承认它不合适:
- 需要单条消息的延迟投递或 TTL。Kafka 没有原生定时消息,retention 是层级而非单条。RocketMQ 的延迟级别、RabbitMQ 的 per-message TTL 更直接(见 《可扩展性》05 章)。
- 需要复杂路由。按内容路由到不同下游、一个消息要投多个队列,这落在 exchange / topic 匹配这类模型上。
- 需要「消费即确认、确认即删除」的强语义。Kafka 的位点是按分区推进的,做不到单条级别的确认与隔离,只能自己搭死信 topic。
- key 热点把单个分区打满。一个分区只能落在一台 broker 的一个目录上,某个超级热点 key 无论怎么扩分区都只在它命中的那一个分区里,这时候要在业务层拆 key 或者把热点单独拎出来。
- 需要跨行/跨消息的全局顺序且量还很大。单分区能保顺序但有写上限;多分区能扩但没全局顺序。这个二选一有时需要换 RocketMQ 的 Queue 或者自己用数据库做序列。
反过来,如果只是「写入量大、要回放、要多个下游各读一遍」,Kafka 是对的。
生产边界
- 教学替身 vs 真实依赖:本机是单机三节点、单块盘、无真实网络,吞吐数字只代表这台机器的量级,不能直接当容量规划输入。真实集群的分区和副本分布在多台机器上,还会有机架感知、跨机房复制这些本机不具备的维度。本机也无法演示「分区数打满后 broker 启动变慢」这类规模效应。
- 要盯的指标:见上面的指标表;最小集合是分区级 lag、
UnderReplicatedPartitions、OfflinePartitionsCount、各盘可用空间、消费端poll-interval。 - 失败策略:磁盘告警阈值要留够一个检查周期的余量;
log.retention.check.interval.ms不要设得比 retention 期望的精度还大;分区数定好后如果要改,改动方案按「重灌 topic + 双写迁移」准备,别指望在线缩分区。
动手
- 对你自己的一个 topic 跑一次
ProducerPerformance,分别在「不限速」和「限速到你能接受的延迟档」下测,记录两组吞吐/延迟数字。 - 用消费端实测速率算:如果你现在的 lag 是 500 万条、上游仍在以 2 万条/秒写入,单消费者(按本机 5.7 万/秒)要多久追平?如果加两个消费者呢?
- 断言题:给你一个 12 分区、6 消费者的组,消费速率总和是 30 万条/秒,上游是 35 万条/秒——列出你能做的三件事,并说出各自的代价。
自测
- 为什么「不限速」吞吐测量的 p50 延迟能到 640 ms,而限速档只有 3 ms?这两个数字里哪一个才能代表「Kafka 慢不慢」?
- 「净速度」为正和为负,分别对应什么处置?为什么净速度为正时加消费者可能是错的?
- 分区数太少和太多分别会怎样?「需要保序的 key 基数」为什么是分区数的下界?
UnderReplicatedPartitions和OfflinePartitionsCount分别在什么情况上涨?为什么前者业务常无感而后者是事故?- 设了
retention.ms=1h,为什么数据不会在第 3601 秒就消失?还有哪些因素在拖后实际删除时间?
本课到这里结束。回到 课程导读 复查章节地图,或者回到 《可扩展性》05 章 把 Kafka 放回选型语境里再看一遍。