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 条/秒 -> 追不上, 积压只会涨

三个场景对应三种完全不同的处置:

注意这些数字是「本机单消费者」的量级,真实集群的绝对值会差很远(磁盘、网络、消息大小都不同),但算法不变:先测出消费速率和生产速率,再算净速度。不要用别人的数字替你做决定。

分区数怎么定

分区数是这门课里最难给「标准答案」的参数,因为它同时被四件事拉扯:

一个可操作的起点:分区数取「峰值消费者实例数」和「需要保序的 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 打开。这个组合暴露了三件同源的事故,全部真跑到:

  1. 删 topic(README 里警告过,这次没做):目录重命名失败 → log dir offline → broker 自杀。

  2. 缩副本的 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 failed
    

    a1、a2 双双退出,controller 仲裁只剩 a3,元数据一起冻结。

  3. 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 + 最多一个检查周期。设了 15 秒的 retention,最坏情况下要等 5 分钟才删。把检查周期调小会更频繁地扫目录,大集群上是要付代价的,通常不建议低于 1~2 分钟。

另外,正在被未完成事务引用的段不会删(deleteHorizonMs,第 05 章)。一个卡住的长事务能拖住一个分区末尾段的清理。

什么时候该换架构

Kafka 不是所有「消息」问题的答案。出现下面这些信号时,要承认它不合适:

反过来,如果只是「写入量大、要回放、要多个下游各读一遍」,Kafka 是对的。


生产边界

动手

  1. 对你自己的一个 topic 跑一次 ProducerPerformance,分别在「不限速」和「限速到你能接受的延迟档」下测,记录两组吞吐/延迟数字。
  2. 用消费端实测速率算:如果你现在的 lag 是 500 万条、上游仍在以 2 万条/秒写入,单消费者(按本机 5.7 万/秒)要多久追平?如果加两个消费者呢?
  3. 断言题:给你一个 12 分区、6 消费者的组,消费速率总和是 30 万条/秒,上游是 35 万条/秒——列出你能做的三件事,并说出各自的代价。

自测

  1. 为什么「不限速」吞吐测量的 p50 延迟能到 640 ms,而限速档只有 3 ms?这两个数字里哪一个才能代表「Kafka 慢不慢」?
  2. 「净速度」为正和为负,分别对应什么处置?为什么净速度为正时加消费者可能是错的?
  3. 分区数太少和太多分别会怎样?「需要保序的 key 基数」为什么是分区数的下界?
  4. UnderReplicatedPartitions 和 OfflinePartitionsCount 分别在什么情况上涨?为什么前者业务常无感而后者是事故?
  5. 设了 retention.ms=1h,为什么数据不会在第 3601 秒就消失?还有哪些因素在拖后实际删除时间?

本课到这里结束。回到 课程导读 复查章节地图,或者回到 《可扩展性》05 章 把 Kafka 放回选型语境里再看一遍。

进入 keel 阅读