Kafka 消费调优作战手册
下游消费峰值 = 生产 10×,跨云 50ms,从未压测。这份手册把「客户场景讲清 + Kafka 参数吃透 + 现场怎么做」压成可背记、可对答、可对比的一份。所有参数默认值来自 Kafka 4.0 官网自动生成配置表(curl 直抓、脚本解析、64/64 零缺失)+ apache/kafka 源码逐字核对 + AWS MSK 官方文档 + 东京真机 POC。
① 升机型能救「扇出型 10×」,救不了「单组追赶型 10×」——一个分区在一个 group 内只能有一个消费者,这是协议不是资源。
② 跨云 50ms 下,consumer 默认的 64KB 接收缓冲把单连接天花板压到 1.31 MB/s(同云同配置天花板 ~130 MB/s,根本不构成约束)。不改这一个参数,其他都是白调。
③
retention=1天 + auto.offset.reset=latest(默认)= 积压超一天就静默丢交易数据,不报错、不告警,只有一行 INFO 日志。今天就改。① 客户场景与消费方式(背景基线)
先把客户的部署形态、数据流、消费方式讲到能对答如流。这些是所有建议的前提。
部署与数据流拓扑
| 维度 | 事实 | 含义 |
|---|---|---|
| 平台 | AWS MSK 托管(ZooKeeper 模式,zookeeper.session.timeout.ms=18000 佐证) | 内核 sysctl 客户/我们都动不了;只能靠 MSK 暴露的 broker 配置 + OS 默认 |
| 双写架构 | 业务与数据双写:一个业务集群把数据写进 MSK,业务本身也写到 MSK 这一侧 | 多数写入同云同区,不跨云 |
| 跨云消费 | 一个集群存在跨云消费:阿里云侧消费者跨专线读 MSK,RTT≈50ms(已实测) | 跨云的是消费者,不是副本。副本同步走 AWS region 内 AZ 间链路(~1-2ms) |
| 消费端形态 | 下游是应用程序 consumer group,不是 Spark / 大数据并行读 | 受「一分区一消费者」协议约束,并行度 = 分区数 |
| 扩容偏好 | 纵向升机型,不加节点(怕加节点触发 rebalance 搬数据) | 这个判断部分正确(见 ⑨),正好是 AWS 官方首推的 Option 1 |
auto.create.topics.enable=false(生产必须关,很多客户没关)· RF=3 + min.insync.replicas=2(黄金搭配基础)· broker socket buffer=-1(MSK 上是最优解,见 ④)· 「怕加节点触发数据迁移」的判断(AWS 官方 Option 3 明确 CPU>70% 不建议做 reassignment)。先肯定这四点,是最快的信任建立。核心矛盾
若是 MSK 侧(机型/带宽/磁盘)→ 升机型;若非 MSK 侧 → 大概率是分区数不足 + 消费并行度不足 + 跨云 fetch 未调优。
现有 broker 参数清单(客户实际值)
| 参数 | 客户值 | 判决 | 备注 |
|---|---|---|---|
auto.create.topics.enable | false | 保持 | 生产环境应关 |
default.replication.factor | 3 | 保持 | 3-AZ 标配 |
min.insync.replicas | 2 | 保持 | 配 RF=3+acks=all 黄金组合 |
num.io.threads | 16 | 调至=vCPU | 官方公式是全量 vCPU 不是一半;升机型必须同步调 |
num.network.threads | 5 | 调至=vCPU/2 | = MSK 默认,客户没动过。大扇出压的正是它 |
num.partitions | 1 | 危险默认 | 漏写 --partitions 就静默单分区。逐 topic 核实 |
num.replica.fetchers | 2 | 调至 4 | = MSK 默认,没动过;跨云/多分区/大消息需更大 |
replica.lag.time.max.ms | 30000 | 保持 | = 默认,且 MSK 封顶 30000,调不大 |
socket.receive.buffer.bytes | -1 | 保持 -1 | 回落 OS tcp_rmem 自动调节,MSK 上最优 |
socket.send.buffer.bytes | -1 | 保持 -1 | 同上(明确告诉客户:这个不要动) |
socket.request.max.bytes | 100MB | 保持 | 对 10MB 消息有 10× 余量 |
unclean.leader.election.enable | true | 改 false | 官网默认 false,是 MSK 非分层默认值。交易数据必改 |
log.retention.ms | 1 天 | 关键 topic 3 天 | 消费 lag 超 24h 会静默丢数据 |
message.max.bytes | ~10MB | 保持 | 但必须连锁调 4 个参数(见 ⑥) |
unclean=true、num.network.threads=5、num.replica.fetchers=2、min.insync.replicas=2(3-AZ)、auto.create=false 全部 = MSK 非分层默认值。说明客户是继承了 MSK 默认,不是主动决策。→ 改 unclean=false 不是否定谁的决策,是「补一个默认值审计」。而且客户已改过 num.io.threads/retention/socket buffer 三项,证明改配置的流程和权限都通。② 「10×」的两种语义 —— 整场沟通的分水岭
同一句「下游消费是生产的 10 倍」,有两种完全不同的语义,对应两套完全相反的解法。不先分清,后面所有调优都可能调错方向。
| 语义 A:扇出(fan-out)10× | 语义 B:单组追赶(catch-up)10× | |
|---|---|---|
| 含义 | 同一份数据被 ~10 个 consumer group 各读一遍 | 一个 group 积压后,需以 10× 生产速率消化 lag |
| 瓶颈在哪 | broker 网络出带宽 + page cache 命中率 + EBS 冷读 | consumer group 并行度(分区数 × 实例数) |
| 分区数需求 | 不因扇出增加(每组各需 1× 并行度) | 必须 ≈ 10 倍于稳态所需分区数 |
| 升机型有用吗 | ✅ 直接有效 NIC/EBS/内存都随机型涨 | ❌ 完全无效 单分区只能一个消费者,机型再大也不快 |
| 加分区有用吗 | 基本无关 | ✅ 唯一解 |
BytesOutPerSec / BytesInPerSec 的实际比值。证据 2:
kafka-consumer-groups.sh --list 数出的 group 个数。比值 ≈ group 数 → 扇出型(走升机型);单个 group 就要跑 10× → 追赶型(升机型无效,必须加分区+实例)。
③ 三大地雷(开场就该主动点出,显专业)
客户 message.max.bytes ≈ 10MB,但 replica.fetch.max.bytes 官网默认只有 1 MiB,consumer max.partition.fetch.bytes 默认也是 1 MiB。
现代 Kafka(≥0.10.1 KIP-74)有「首个 batch 超限也返回」的兜底,所以不会永久卡死,但会退化为「每次 fetch 只拿回一个大 batch」。在 50ms RTT + 消费者背 12 个分区时,产生 head-of-line blocking → ~12× 效率损失,这与「消费跟不上」的症状高度吻合。复制侧退化会让 follower 落后 → 间歇掉出 ISR → 概率性触发 minISR=2 生产被拒。
关键操作约束:replica.fetch.max.bytes 是 read-only,不能热改,必须滚动重启(12 broker ≈ 2-3 小时)。明天要调得先跟客户约重启窗口,不能当场承诺。
另注意:客户 message.max.bytes=10485880 比 replica.fetch.response.max.bytes 默认 10485760 多 120 字节 → 每条满载消息都走 fallback 路径。这个边界重合是「我们逐字节核对过你的配置」的证据。
官网默认是 false。允许落后于 HW 的副本在 ISR 全空时当选 leader → 新 leader 日志比旧的短 → 其他副本按新 leader 截断 → 已 ack 给 producer 的数据永久丢失,且 consumer 的 committed offset 可能越界触发重置。整个过程自动、静默、凌晨发生,日志里只有一行记录。
好消息:Update Mode = cluster-wide,可动态改,无需重启。且 false 不等于放弃可用性——kafka-leader-election.sh --election-type unclean 可逐分区人工授权。真正的对比是:丢数据这个决定,是配置默认值做的,还是你们的值班工程师做的?
10× 峰值下积压超过 1 天 → broker 删除过期 segment → committed offset 落在已删日志段外 → 下次 fetch 返回 OFFSET_OUT_OF_RANGE → latest(默认)直接跳到队尾,中间所有未消费的交易数据永久丢失。
精确说法(会被内行追问):客户端会打一条 Fetch position ... is out of range, resetting offset,但级别只有 INFO,极易被淹没。危险在于不抛异常、不中断消费、默认监控不告警。这是本次评估最高危、最容易被忽略的一条。
处置:retention.ms 是 topic 级 动态参数,kafka-configs.sh --alter --entity-type topics 即时生效、不重启、只增磁盘。交易 topic 建议 3 天;auto.offset.reset 改 none(宁可 consumer 挂掉报警,不要静默丢数据)。
④ 跨云 50ms 的数学(BDP)—— 要能在白板上推
这是「升机型也救不了跨云消费慢」的第一性解释,也是明天最有冲击力的一段。
带宽时延积BDP(带宽时延积)= 带宽 × RTT
单 TCP 连接吞吐上限 = min(接收窗口, 发送缓冲) / RTT
Kafka consumer 的 receive.buffer.bytes 默认 = 65536 (64KiB) ← 显式正数,不是 -1
→ 单连接上限 = 64KiB / 50ms = 1.31 MB/s
→ 200B 消息 ≈ 6,550 msg/s;1KB ≈ 1,280 msg/s;10KB ≈ 128 msg/s
| receive.buffer.bytes | 50ms 单连接上限 | 200B 消息 | 1KB | 10KB |
|---|---|---|---|---|
| 65536(默认 64KB) | 1.31 MB/s | ~6,550/s | ~1,280/s | ~128/s |
| 1 MiB | 20.97 MB/s | ~105k/s | ~20.5k/s | ~2,050/s |
| 4 MiB | 83.9 MB/s | ~420k/s | ~82k/s | ~8,200/s |
| 8 MiB | 167.8 MB/s | ~840k/s | ~164k/s | ~16,400/s |
-1(交给 OS) | 由 tcp_rmem[2] 决定 | 推荐:client 端保持 -1 + 调大 EC2 tcp_rmem 上限 | ||
② 不要说「跨云比同云慢 100 倍」——同云同配置天花板 ~130 MB/s(0.5ms RTT),高到根本不构成约束。诚实表述是「跨云是天花板从不存在变成低于业务需求」。具体倍数按客户真实消息大小现场算。
receive.buffer.bytes=8MiB 但 EC2 net.core.rmem_max 只有 208KB(很多发行版默认)→ 内核静默截断到 rmem_max,不报错、不告警,客户以为调了其实没生效。而且一旦显式 setsockopt(SO_RCVBUF),Linux 的接收缓冲自动调节就被关闭。验证:
ss -tim dst <broker-ip> 看 rcv_space 是否长过 64KiB(注意 ss 显示的是你设值的 2 倍,别误判)。正解:client 端
receive.buffer.bytes 保持 -1(走自动调节)+ 把 EC2 net.ipv4.tcp_rmem[2] 调到 32MiB。两条路二选一,混着来就是白调。第二个天花板:fetch 深度 = 1(但与用户处理并行)
consumer 对同一 broker 同一时刻只允许一个 in-flight fetch(AbstractFetch 硬规则)。但精确说法是:同一 broker 的 fetch 深度=1,但 fetch 与用户处理是并行的——poll() 在返回本批前就发出下一轮 fetch。所以别说「fetch 完全不能流水线」,要说「单连接吞吐 ≈ 单次响应大小 / max(RTT+broker时间, 单批处理时间)」。默认下 TCP 窗口(1.31MB/s)比单分区 fetch 限额(20.97MB/s)低 16 倍 → 必须先调 receive.buffer.bytes,再调 fetch 参数才有意义,顺序反了没效果。
跨云 EC2 sysctl(客户自己的机器,完全可控)
/etc/sysctl.d/99-kafka-client.conf# 前提:客户端 Kafka receive.buffer.bytes 保持 -1 才走自动调节
net.ipv4.tcp_rmem = 4096 131072 33554432 # max=32MiB,50ms 下支撑 ~5Gbps 单流
net.ipv4.tcp_wmem = 4096 65536 33554432
net.core.rmem_max = 33554432 # 兜底:防显式设值被静默截断
net.core.wmem_max = 33554432
net.ipv4.tcp_window_scaling = 1 # 必须为 1,否则窗口物理封顶 64KiB
net.ipv4.tcp_moderate_rcvbuf = 1 # 自动调节总开关,必须为 1
# 教科书级判据:iperf3 -P 1 慢 而 -P 8 快 ⇒ 100% 是接收窗口不足
client.rack 匹配不到任何 rack → 回落到读 leader,50ms 一点没省。它真正有用的是同 region 内跨 AZ 消费(省 AZ 间流量费 + 省 ~1ms),对客户那些同云集群值得开。这个区分要讲清,别一句「没用」把好东西也否了。⑤ 分区 = 消费并行度的硬上限
Kafka 的 group 协议:一个 topic-partition 在同一 consumer group 内,同一时刻只被分配给恰好一个 consumer 成员(因为 __consumer_offsets 里每个 (group,topic,partition) 只有一条 committed offset,两个成员并发消费会互相覆盖)。所以:
一个 group 内「有效工作的 consumer 实例数」 ≤ 该 group 订阅的分区总数
实例数 < 分区数 → 每实例背多个分区(正常)
实例数 = 分区数 → 1:1,最大并行度
实例数 > 分区数 → 多出的实例分到 0 个分区,纯空转(--members 里 #PARTITIONS: 0)
auto.create.topics.enable=false 挡不住有人 kafka-topics.sh --create 时漏写 --partitions → 静默回落 num.partitions=1,不报错。明天第一条命令就该 grep 'PartitionCount: 1' 审计。10× 峰值需要多少分区(现场演算)
所需分区数 N ≥ ceil( R_peak / t_single ) × 1.5 再向上取整到 broker 数(12)的倍数
R_peak = 该 group 峰值消费速率(msg/s 与 MB/s 都算,取更严)
t_single = 单实例单分区实测能力(必须实测,见现场压测方案)
1.5 = 安全系数(GC / rebalance / 慢分区 / 倾斜余量)
| 示例(P=50k msg/s, t_single=5k) | 语义 A:10 组各读 1× | 语义 B:单组追赶 10× |
|---|---|---|
| 单组 R_peak | 50,000 msg/s | 500,000 msg/s |
| 裸算分区数 | 10 | 100 |
| × 1.5 安全系数 | 15 | 150 |
| 对齐 12 broker | 24 | 144(12×12)或 156(12×13) |
| 需要 consumer 实例数 | 10~24/组 | ≥ 100 |
分区数完全不是约束:36×RF3÷12 = 每 broker 仅 9 个副本,AWS 对 4xlarge+ 的推荐上限是 4000。客户完全不必在 36 停手,144~288 都很轻。
加分区的三个代价(客户「愿意加分区」前必须知道)
- 老数据不重分布:加分区只影响新写入 → 「加了分区 lag 没降」是正常现象 → 加分区是预防措施,不是急救措施,必须在峰值到来前做完。
- 破坏 key→partition 映射 最重要:默认 partitioner 是
toPositive(murmur2(key)) % numPartitions(不是 hashCode)。16→36 后同一 key 落点改变 → 新消息在新分区、旧消息在老分区未消费完 → 两分区间无顺序保证。对交易/撮合/账务这是正确性事故,不是性能问题。缓解:排空后扩 / 新 topic 双写 / 自定义 Partitioner(hash%12288固定模数再映射)。 - 触发一次 consumer group rebalance:默认
RangeAssignor是 eager = stop-the-world。用CooperativeStickyAssignor可降到只影响迁移的分区。
⑥ 大消息约束链(客户 10MB 消息,全链必须一起放大)
官方约束链producer max.request.size ≥ 10MB ← 不满足:produce 直接失败
broker message.max.bytes = 10MB ← 客户已设
broker replica.fetch.max.bytes ≥ 10MB ← 不满足:复制退化→ISR收缩→生产被拒
broker replica.fetch.response.max.bytes ≥ 32MB ← 客户值超默认 120B,走 fallback
consumer max.partition.fetch.bytes ≥ 10MB ← 不满足:消费退化,lag 单调涨
consumer fetch.max.bytes ≥ max.partition.fetch.bytes
RecordTooLargeException,消息文本其实很明确)。所以要先确认 kafka-clients 版本。max.request.size 和 max.partition.fetch.bytes 的值。」若这两个是默认 1MiB 而消息 10MB → 这是已经在发生的生产事故,且完全解释「消费跟不上」。30 秒定位问题,比任何压测都快。min(fetch.max.bytes, 分区数 × max.partition.fetch.bytes)。36 分区 × 12MiB = 432MB → 必须同步加 JVM -Xmx;fetch.max.bytes 设成明确天花板(如 100MB)当安全阀;加分区必须同步加实例数(让每实例分到的分区数不增长)。⑦ 可靠性黄金四角(交易类客户专题)
持久性不是一个参数,是一个四元乘积。任何一个因子是 0,整个乘积就是 0。
acks=all × min.insync.replicas=RF-1 × unclean.leader.election=false × 幂等
客户: ❓未知 ✅ =2 ❌ 当前 true ❓未知
| 错误组合 | 为什么错 | 后果 |
|---|---|---|
acks=1 + minISR=2 | minISR 仅对 acks=all 生效 | 🔴致命 自以为双副本保证,leader 挂即丢已 ack 数据。minISR=2 是纯装饰 |
acks=all + minISR=1 | ISR 可收缩到仅 leader | 🔴致命 抖动后静默退化为 acks=1 |
acks=all+minISR=2+unclean=true | 非 ISR 副本可当选并截断日志 | 🔴致命 前两项做对也照样丢已 ack 数据 ← 客户当前状态 |
| 四角齐全但未开幂等 | append-then-fail 窗口(NOT_ENOUGH_REPLICAS_AFTER_APPEND) | 🟠高 重复消息 |
acks 未知——若是 acks=1,minISR=2 就是纯装饰),而且留了一个后门(unclean=true)。」必须当场问出:producer 的
acks 实际值 + kafka-clients 版本(3.0+ 默认 all + 幂等 true,2.x 默认 acks=1)。⚠️ 3.x 隐式禁用陷阱:显式留了 acks=1 而没显式设 enable.idempotence → 幂等被静默禁用,只打一行日志。看启动日志实际生效值,不看配置文件。acks=1 只是选择不等复制完成就返回,复制该花的带宽和 IO 一分没省,你只是放弃了那个等待换来的保证。」延迟明显变差 = ISR 里有慢副本 = 同一个根因正在拖高所有消费者的端到端延迟(HW 卡在最慢 ISR 成员)。__consumer_offsets 的 replication.factor 默认一直是 3(别说「Apache 默认 1」,那是本地 quickstart 值,会被内行纠正)。但它只在首次创建那一刻生效——从老版本升上来的集群可能是 RF=1/2。若 RF=1,承载它的 broker 挂 → 全集群消费位点丢失。一条命令审计:kafka-topics.sh --describe --topic __consumer_offsets | head -1。参数速查 · Broker 端
默认值全部来自 Kafka 4.0 官网自动生成配置表。dyn=可动态改无需重启(cluster-wide);ro=read-only 必须滚动重启。
| 参数 | 官网默认 | 改动性 | 客户当前 | 建议 | 一句话理由 / 调错后果 |
|---|---|---|---|---|---|
unclean.leader.election.enable | false | dyn | true | false | ISR 全空时允许落后副本当 leader → 丢已 ack 数据。交易类必改 |
replica.fetch.max.bytes | 1 MiB | ro | 未知🔴 | ≥12MiB | <消息大小 → 复制退化 → ISR 收缩 → 生产被拒 |
replica.fetch.response.max.bytes | 10485760 | ro | 未知🔴 | ≥32MiB | 客户消息超默认 120B,每条走 fallback |
message.max.bytes | 1048588 | dyn | ~10MB | 保持 | 已设对,但必须连锁调整个链 |
num.io.threads | 8 | dyn | 16 | =vCPU | 官方=全量 vCPU;太小则 RequestHandlerAvgIdle<20%,请求排队 |
num.network.threads | 3 | dyn | 5 | =vCPU/2 | 大扇出压的是它;只加它不加 io.threads → queue saturation 更差 |
num.replica.fetchers | 1 | dyn | 2 | 4(分步) | 跨云/多分区/10MB 消息需更大;总线程=本值×(broker-1);盯 HeapMemoryAfterGC |
replica.lag.time.max.ms | 30000 | ro | 30000 | 保持 | MSK 封顶 30000,调不大;治因不治症(加 fetchers + 修 fetch.max.bytes) |
socket.receive/send.buffer.bytes | 102400 | ro | -1 | 保持 -1 | -1=交给 OS tcp_rmem 自动调节;改具体值 = 关自动调节 + 被黑盒 rmem_max 静默截断 = 净负收益 |
min.insync.replicas | 1 | dyn | 2 | 保持 | 配 RF=3+acks=all 黄金组合 |
num.partitions | 1 | ro | 1 | 12 或 24 | 兜底值。漏写 --partitions 就永久串行瓶颈。搭机型升级重启一起改 |
log.retention.ms | 7天(hours=168) | dyn | 1天 | 关键topic 3天 | lag 超 24h → 数据被删 → 静默丢。topic 级动态改 |
queued.max.requests | 500 | MSK锁死 | 500 | 不可改 | MSK 不在允许列表;队列深度不可调 → io.threads 是唯一旋钮 |
参数速查 · Consumer 端(本次头号战场,客户值全部未知,必须问)
| 参数 | 官网默认 | 建议(跨云 50ms) | 一句话理由 / 调错后果 |
|---|---|---|---|
receive.buffer.bytes | 65536(64KB) | 4~8MiB 或 -1 | 跨云头号杀手 64KB/50ms=1.31MB/s。超 rmem_max 内核静默截断 |
max.partition.fetch.bytes | 1 MiB | ≥12MiB | 大消息必需 默认对 10MB 消息 = 退化,直接解释消费跟不上。调大要加 -Xmx |
fetch.max.bytes | 52428800(50MiB) | ≤55MiB 才有意义 | broker 侧 fetch.max.bytes 默认 57671680(55MiB) 会截断,设 100MB 无效 |
fetch.min.bytes | 1 | 1MB | 默认=有一字节就返回。⚠️ 有积压时此项收益很小(broker 早有数据),治的是稳态 |
fetch.max.wait.ms | 500 | 100~200 | 与 fetch.min.bytes 成对;设 5000 → 空闲时 P99 延迟 5 秒 |
max.poll.records | 500 | 大消息降 10~50 | Streams 覆盖成 1000 别混淆。× 单条处理耗时 < max.poll.interval.ms,否则死亡螺旋 |
max.poll.interval.ms | 300000(5min) | 先别动 | 处理慢被踢的活性检测;先降 max.poll.records / 优化逻辑 |
session.timeout.ms | 45000 | 60000 | ⚠️新版已 45s(非旧版 10s)。硬校验 heartbeat < session;官网建议 ≤session/3 |
heartbeat.interval.ms | 3000 | 5000 | ≥session 则启动抛 IllegalArgumentException |
enable.auto.commit | true | false | 交易数据改手动批量提交。每条 commitSync 跨云=20 msg/s;每批一次=数千 |
auto.offset.reset | latest | none | 最高危 latest+retention1天=积压超期静默丢数据(仅一条 INFO 日志) |
partition.assignment.strategy | [Range, CoopSticky] | 只用 CoopSticky | 默认实际生效 Range=eager=stop-the-world。升级必须两轮滚动,不能一步切 |
group.instance.id | 无(动态) | 设置(静态) | 滚动重启零 rebalance。两实例同 ID 会互相踢 |
isolation.level | read_uncommitted | 视是否用事务 | 用事务必设 read_committed;长事务钉住 LSO 会造成假 lag |
参数速查 · Producer 端
| 参数 | 官网默认(新版) | 建议 | 一句话理由 |
|---|---|---|---|
max.request.size | 1 MiB | ≥12MiB | P0第一问 默认对 10MB 消息 → produce 直接抛 RecordTooLargeException |
acks | all(3.0+)/1(2.x) | all | 只有 acks=all 时 broker 才检查 minISR。决定 minISR=2 真假保护 |
enable.idempotence | true(3.0+) | true | 幂等开启时 max.in.flight≤5 仍保序。看启动日志实际值,不看配置文件 |
max.in.flight.requests.per.connection | 5 | ≤5 | 为「保序」设 1 是反的(50ms 下封顶 20 请求/s)。开幂等就是为了设回 5 还保序 |
compression.type | none | lz4/zstd | 跨云强烈建议,直接降跨云带宽与延迟 |
linger.ms | 5 | 跨云 20~100 | 攒批;新版默认已是 5(旧版 0) |
batch.size | 16384 | 64~256KB | 跨云提高攒批效率 |
delivery.timeout.ms | 120000 | 与业务 RTO 对齐 | 配 retries=MAX 故障期透明重试;但下界=一次 leader failover 时长(与分区数成正比,实测) |
buffer.memory | 32MB | ≥128MB | 10MB 消息,32MB 只缓 3 条,send() 会阻塞 |
linger.ms/batch.size/分区数。别只甩那个 20。但明天一定要 grep 一遍这个配置。⑧ 现场分诊决策树
RequestQueueTime 大 = io.threads 不够(调参解决)· LocalTime 大 = 磁盘/pagecache miss(升 EBS)· ResponseQueueTime 大 = network.threads 不够(调参解决)· ResponseSendTime 大 = 网络/socket 窗口(修客户端窗口)。跨云客户端会让 ResponseSendTime 天然偏高——这正是 BDP 受限在 broker 侧的镜像信号,看到它就不用再争「是不是 MSK 的问题」。⑨ 三种 rebalance 彻底分清(客户误解的核心)
客户「怕 rebalance 所以不加节点」,把三件不同的事混成了一件。分清之后扩容选项会一下子变多。
| ① 分区副本重分配 reassign-partitions | ② preferred leader 选举 leader-election | ③ consumer group rebalance | |
|---|---|---|---|
| 发生层 | Broker 之间(服务端) | Broker 之间(仅角色切换) | Consumer 实例之间(客户端) |
| 搬数据吗 | ✅ 搬(唯一搬数据的) | ❌ 不搬 | ❌ 不搬 |
| 耗时 | 分钟~小时 ∝ 数据量/throttle | 毫秒~秒 | eager: 秒~几十秒(stop-the-world) |
| 可限流/取消 | ✅ --throttle / --cancel | — | ❌ 只能靠 assignor 降影响 |
| 加节点触发? | ⚠️ 不自动触发(新broker空的,要人工reassign) | ✅(新broker有副本后) | ✅ |
| 加分区触发? | ❌ 不触发(新分区直接空建) | ✅ | ✅ 唯一代价 |
| 升机型触发? | ❌ 不触发 | ✅(滚动后 5min 内自动回归) | 通常不触发 group rebalance |
rebalance-rate-per-hour,不是 0 就先修这个。」⑩ 动作清单(红黄绿灯)
- 关键 topic
retention.ms: 1天→3天(topic级动态,先确认磁盘有 3× 余量) - 关键 topic
unclean.leader.election.enable=false(topic级动态,2分钟,升级前最后一道保险) - 建 4 条 CloudWatch 告警:CPU60% / TrafficShaping / UnderMinIsr / EstimatedMaxTimeLag 12h·18h
- 动态调
num.io.threads→=vCPU,num.network.threads→=vCPU/2(先 io 后 net) - 动态调
num.replica.fetchers2→4(分步,盯 HeapMemoryAfterGC) - 只读审计脚本 + perf-test 二分法(TrafficShaping 非零时先确认再跑)
- ⚠️ 加分区是唯一不可回滚的绿灯——先问「有 key 吗?依赖有序吗?」
replica.fetch.max.bytes→ ≥12MiB read-only 需滚动重启 12broker≈2-3hreplica.fetch.response.max.bytes→ ≥32MiB(同重启窗口)- broker 默认
num.partitions1→12(搭机型升级重启,零额外成本) - 客户端 consumer/producer 参数(
max.partition.fetch.bytes/receive.buffer.bytes/max.request.size/auto.offset.reset)+ 相应加 JVM heap partition.assignment.strategy→ CooperativeSticky(两轮滚动升级,需 clients≥2.4)
- 加分区(有 key 且依赖有序的 topic)→ 必须走排空后扩/新topic双写/自定义Partitioner
kafka-reassign-partitions.sh重分布老数据 → 必须--throttle,单次≤10分区,CPU>70% 不做,完成后必须--verify清限流(忘了 = 限流永久残留拖慢正常复制)- 升机型(滚动 2-3h)→ 三条硬性前置检查见下
--under-replicated-partitions 必须为空才动下一台。② 不存在 RF<3 的 topic(否则单台下线即触及 minISR)。③ 先把关键 topic unclean=false 打上(零重启,消除升级窗口内丢数风险)。④ 升机型后必须同步调线程池——MSK 非分层集群死给 8/5,不随 vCPU 涨,升机型不调线程池 = 白花钱。queued.max.requests(MSK 不让改)· 调大 replica.lag.time.max.ms(已在 30000 天花板)· 只加 network 不加 io.threads(queue saturation 更差)· 用 100-CpuIdle 当 CPU 利用率(iowait/irq/steal 不计入)· 指望 client.rack 解决阿里云跨云 · MM2 给交易数据做跨云容灾(结构性破坏 EOS)· 回退 acks=all 到 acks=1 降延迟。⑪ 必须问客户的清单(拿不到就是盲调)
- producer
max.request.size是多少?(默认 1MiB + 10MB 消息 = produce 直接失败,已在发生的事故) - consumer
max.partition.fetch.bytes是多少?(默认 1MiB → 消费退化,直接解释消费跟不上) - broker
replica.fetch.max.bytes是多少?(默认 1MiB → ISR 收缩 → 生产被拒,最高危)
BytesOut/BytesIn 比值 + group 数(定 10× 语义)· topic 总数/各分区数/单分区清单/RF<3 清单 · 每组实例数 · 峰值 msg/s + 平均消息大小 · 跨云集群哪些 topic + 专线带宽 · 生产端指定 key 吗?依赖有序吗?(决定加分区能否随手做)unclean=false(先关键 topic,零重启)· 是否接受 retention 提到 3 天(磁盘换安全)· 是否接受 auto.offset.reset=none(需业务方配合处理异常)· kafka-clients 版本号 · MSK 监控级别当前哪一档⑫ Observer/Learner 前瞻(加分项,不是主菜)
是什么:给 Kafka 加第三种副本状态——全量同步数据但永不进 ISR。四个性质全是 Kafka 规则的推论:不拖 HW / 不拦 acks / 不算 minISR / 不被选主。5 个 hook 点,ZK 模式 ~60 行 Scala,核心 gate canAddReplicaToIsr(KIP-497@2.7 引入,所以 2.7+ 才支持)。真机验证:2.7→4.3 共 20 build × S1-S8 全绿;observer 在最慢 AZ 时 acks=all 仍 2.04-2.35ms(零拖累);晋升 ≤10s 零重启零数据搬迁。
findPreferredReadReplica 一字不差):就近读的候选集显式过滤了不在 ISR 的副本 → observer 永不进 ISR → 阿里云消费者依旧回落读 leader,50ms 一点没省。要做「跨云本地只读」还差第 6 个 hook,当前项目没实现、ROADMAP 也没有。「这条路设计上通、但我们还没验证,我不给你没验证过的数字。」——这个诚实点反而是加分项。FAQ · 现场问答弹药
Q:unclean=true 是 MSK 默认的,AWS 自己都这么设,为什么改?
AWS 托管默认值服务所有负载,选了可用性优先。但 Apache 上游从 0.11 起默认就是 false,且 AWS 自己在 MSK 分层存储集群的默认配置里已经把它设成 false——AWS 在新形态下的取舍已经反转。你们是交易类,属于该反转的那一类。这不是 MSK 的 bug,是默认值和业务画像不匹配。而且你们已改过 io.threads/retention/socket buffer,改配置的流程本来就通。
Q:改成 false 会不会停服?
会,但只在「某 partition 全部 ISR 同时不可用」这个具体条件下。⚠️诚实补一句:minISR=2 不保证 ISR 始终有 2 个成员(它只在 ISR<2 时拒绝写入)——有慢副本时 ISR 可能只剩 leader,丢 1 个 AZ 就可能触发。所以代价取决于 ISR 健康度,这也是为什么要同时建 UnderReplicated/UnderMinIsr 告警。而且 false 不是禁止 unclean,是把它从「系统凌晨自动静默执行」改成「人工显式授权执行」。
Q:跨云太慢,在阿里云也建个 Kafka 用 MirrorMaker 同步?
分两种:性能目的(降带宽+本地消费)→ 对的,标准解法,跨云只传 1× 本地扇出 10×。交易数据容灾目的→ 请不要,数学上不可能:目标 offset 由目标 leader 重分配(两个 offset 空间,只能近似 translate);MM2 崩溃后从源位点重放(默认 at-least-once,实测重放 20000 条);PID/sequence/事务 marker 都不透传。KIP-618 只解决目标 produce 侧不重复,offset 空间不同它解决不了。
Q:我们升级到 3.x 了,幂等应该默认开着吧?
不一定,高频陷阱。3.x 客户端里若配置文件显式留了 acks=1(从 2.x 继承)而没显式设 enable.idempotence → 幂等被静默禁用,只打一行日志。验证:在应用启动日志里搜 idempotence 看实际生效配置,不看配置文件。
Q:加分区之后为什么 lag 没降?
加分区只对新写入生效,老数据一条都不重分布(官方原文 "Kafka will not attempt to automatically redistribute data")。旧积压还在老分区里靠原来的消费者慢慢啃,新空分区的消费者在空转,短期内有效并行度反而更低。加分区是预防措施,不是急救措施。另外扩分区后 consumer 最长 5 分钟(metadata.max.age.ms)才感知,这期间「新分区 lag 快速上涨」是假故障。
Q:为什么不能把 max.in.flight 设成 1 来保序?
反了——开幂等就是为了能设回 5 还保序。机制:broker 对非连续 sequence 直接拒收 + 客户端按 seq 重排后重投(不是 broker 缓存 5 个 batch 帮你排序,那 5 个缓存是用来识别重复的)。50ms RTT 下 in-flight=1 把单连接请求速率钉在 ~20/s(但要说清是 20 请求/s 不是 20 条/s)。明天一定要 grep 一遍这个配置。
官方数字备查(被追问时报得出来)
| 数字 | 出处 |
|---|---|
| 每 broker 推荐 4000 分区(含副本,4xlarge+);update 操作上限 6000 | AWS MSK Best Practices |
| 每集群 200,000 分区 | Apache Kafka 1.1 官方博客 |
| CPU 硬线 60%(CpuUser+CpuSystem);reassign 不建议 >70%;单次 reassign ≤10 分区 | AWS MSK Best Practices |
| 机型升级 10-15 分钟/broker(滚动,clean shutdown);加 broker 数必须是 AZ 数整数倍 | AWS MSK 文档 |
consumer 默认:fetch.min.bytes=1 fetch.max.wait.ms=500 max.partition.fetch.bytes=1MiB max.poll.records=500 session.timeout.ms=45000 receive.buffer.bytes=65536 auto.offset.reset=latest | Kafka 4.0 官网 consumer_config |
broker 默认:unclean=false min.insync.replicas=1 message.max.bytes=1048588 num.partitions=1 replica.fetch.max.bytes=1MiB(ro) replica.lag.time.max.ms=30000(ro) num.replica.fetchers=1 | Kafka 4.0 官网 kafka_config |
producer 默认(新版):acks=all enable.idempotence=true linger.ms=5 max.in.flight=5 delivery.timeout.ms=120000(均 3.0+ / KIP-679) | Kafka 4.0 官网 producer_config |
默认 partitioner = toPositive(murmur2(key)) % numPartitions;默认 assignor = [Range, CoopSticky] 实际生效 Range=eager | BuiltInPartitioner.java / 官网 |
__consumer_offsets RF 默认 3(一直是 3,非 1),"should not change after deployment" | Kafka 4.0 官网 |
Kafka 不支持减少分区数(InvalidPartitionsException) | ops.html + ReplicationControlManager |
开场话术(三句话建立信任,顺序不能变)
auto.create.topics.enable=false(很多客户没关);RF=3 + minISR=2(黄金搭配,基础比很多客户好);socket buffer 保持 -1(MSK 上最优解,网上那些『调到 1MB 提升吞吐』是给自管 Kafka 写的,你们没被误导);还有『怕加节点触发数据迁移』这个判断是对的——AWS 官方明确 reassignment CPU>70% 不要做,你们选纵向升机型正好是官方首推的 Option 1。」--list 加一眼 CloudWatch 就能分清。我们先花两分钟把这个分清,再谈任何调优。」message.max.bytes 设了 10MB,但 producer 的 max.request.size、consumer 的 max.partition.fetch.bytes、broker 的 replica.fetch.max.bytes 官方默认都只有 1MB。只要有一个没跟着放大——第一个让 produce 直接失败,第二个让消费吞吐掉 12 倍(可能就是『消费跟不上』的直接原因),第三个更严重会让生产被拒。前两个要看代码,第三个我现在就能查。先查第三个。」② 跨云 50ms 下 consumer 默认 64KB 接收缓冲把你们钉在 1.31 MB/s,不改这一个参数其他都白调。
③ retention=1天 + auto.offset.reset=latest = 积压超一天就静默丢交易数据,不报错不告警。今天就改。
数据源:Kafka 4.0 官网自动生成配置表(curl 直抓 + 脚本解析 64/64 零缺失)· apache/kafka 源码逐字核对 · AWS MSK 官方文档 · 东京真机 POC(ap-northeast-1, m7g, 2.7→4.3 共 20 build S1-S8)。
本文档不含任何客户名/公司名/账号/IP/主机名等敏感信息。研究细节见 research/01~06 + 00-official-defaults-verified.md。