这些年因为业务需要,我没少跟 Kafka 打交道。从最早拿它当日志管道,到后来支撑核心交易链路,再到帮同事排查线上“消息延迟飙到几分钟”的诡异故障,一个体会越来越深:Kafka 的性能好,不是靠某一项黑科技,而是靠一套环环相扣的架构设计。很多人背了一堆 Kafka 面试题,知道分区、知道 ISR、知道零拷贝,但真到了线上出问题时,依然不知道怎么定位,原因就是把“原理”和“实战”拆开了。
这篇文章,我想把 Kafka 高性能架构设计这件事完整讲透。不堆概念,而是把它为什么快、快在哪里、牺牲了什么,以及网上高频出现的“Kafka 消息延迟高”“Kafka 如何延迟 30 分钟消费”“Kafka OOM”这些真实问题背后的原理讲清楚。适合正在学 Kafka 的开发者,也适合准备面试、或者正在被线上 Kafka 性能问题折磨的运维和后端同学。读完之后你再去看那些所谓的 Kafka 面试题,会发现答案其实都是相通的。
1. 从宏观架构看 Kafka 为何能快起来
1.1 核心组件与协作关系:消息系统界的“快递分拣中心”
先建立一个整体认知。Kafka 的架构可以用一个非常生活化的类比来理解——它就像一个大型快递分拣中心。
- Producer(生产者):就是各个营业网点,不断把包裹(消息)送进分拣中心。
- Broker(代理节点):就是分拣中心里的若干条分拣线,每台机器就是一个 Broker,负责接收、存储和分发包裹。
- Topic(主题):相当于不同的包裹类型,比如“生鲜件”“普通件”“国际件”。
- Partition(分区):这是 Kafka 最核心的设计。每个 Topic 会被拆成多个分区,你可以理解成同一种包裹被分配到多条分拣通道上并行处理。
- Consumer(消费者):就是各个区域的配送站,按需从分拣通道上取走包裹。
从技术角度看,一个 Kafka 集群由多个 Broker 组成;每个 Broker 上会承载若干个 Partition;每个 Partition 是一个有序的、不可变的消息日志,新消息只能追加写入。消费者通过维护自己的 Offset(偏移量)来记录消费位置,就像书签一样。
这套组件协作起来,核心路子就是:生产者把消息写到某个分区的日志尾部,消费者从自己记录的偏移量位置顺序读取日志。没有复杂路由,没有全局索引,数据流是一条笔直的管道。这是 Kafka 高性能的第一个基石——简单。
1.2 分区模型:并行吞吐的第一推动力
分区是 Kafka 弹性的根本。没有分区,一个 Topic 就只能在一台机器上顺序读写,吞吐上限就是单机磁盘和网卡的上限。有了分区,一个 Topic 的数据可以被拆到多个 Broker 上,读写可以并行;同一个分区内的消息有顺序,不同分区之间没有顺序约束。
分区数怎么定,是个经典问题。我在实际项目里的经验是:分区数不是越多越好。分区太多会带来几个副作用:文件句柄数量暴涨、Broker 端元数据变大、消费者 Rebalance 时间变长、单分区副本数量多时副本同步压力大。一般情况下,我会按“目标吞吐量 / 单分区吞吐量”粗算一个基数,再结合消费者实例数调整。比如目标吞吐 50MB/s,单分区实测能扛 10MB/s,那 5 个分区是底线;但为了让消费者并行度更充裕,我通常会再乘一个 1.5 到 2 的系数,同时把峰值流量考虑进去。
还要特别注意:分区数是 Topic 创建时指定的,虽然 Kafka 支持后续增加分区,但一旦增加,键值到分区的映射就会改变,原本有序的消息可能被打散。所以初期设计分区数时,要么按未来三到六个月的峰值预估,要么像很多大厂那样干脆建一个足够大的分区数,比如 24 或 48,用空间换未来扩展的灵活性。
2. 高性能的微观核心:磁盘顺序写与缓存命中的艺术
2.1 顺序写磁盘 vs 随机写:为什么 Kafka 敢用磁盘
很多第一次接触 Kafka 的人都会有一个疑问:明明大家都在吹 Redis 是内存操作所以快,为什么 Kafka 用磁盘还能达到百万级吞吐?
答案在于访问模式。传统消息队列或者数据库,数据的读写位置是随机的,每一次写入可能都要移动磁头寻道,机械硬盘随机写的性能惨不忍睹,通常只有每秒几百次 IOPS。而 Kafka 的设计是 append-only(只追加),消息永远写到日志文件的末尾,读的时候也尽量从末尾附近顺序读。顺序写对磁盘意味着什么?看一组实测数量级你就明白了:
| 访问模式 | 机械硬盘实测性能 | 说明 |
|---|---|---|
| 随机写 | 约 0.1 MB/s 级别(IOPS 制约) | 磁头频繁寻道,性能极差 |
| 顺序写 | 100 MB/s 以上 | 磁头几乎不移动,接近磁盘物理极限 |
| 随机读 | 类似随机写,性能很低 | 缓存未命中时非常痛苦 |
| 顺序读 | 可超过 100 MB/s | 配合预读机制,吞吐非常可观 |
即便换到 SSD,顺序写也远比随机写稳定和高效,而且对闪存寿命更友好。Kafka 正是把随机读写的场景硬生生变成了顺序读写,所以才能把磁盘用出接近内存的效果。
消息写到日志文件后,并不是每条都立刻刷盘。Kafka 允许操作系统把数据先放在 PageCache 里,由操作系统根据脏页比例和空闲内存情况统一刷盘。这里面有一个关键参数组合:log.flush.interval.messages(默认不限制)和log.flush.interval.ms(默认不限制),意思是默认完全交给操作系统管理。绝大多数场景下,让操作系统自己刷盘就是最优解,不要轻易去设这两个参数,强行频繁刷盘反而会杀掉吞吐。
2.2 PageCache:读写都走缓存,绕开物理磁盘
Kafka 高性能的第二个关键,是它对 PageCache 的极致利用。PageCache 是操作系统内核维护的一块内存缓存,用来缓存磁盘文件的内容。Kafka 读写消息时,数据其实都会先经过 PageCache。
生产者写入一条消息:数据写入操作系统的 PageCache 后,Kafka 就认为写入完成了(即使还没有真正落盘)。消费者读消息:如果消息刚写进来,大概率直接命中 PageCache,完全不需要访问物理磁盘。这就是 Kafka 在“写入-读取”时间差较小的场景下,性能表现堪比内存消息队列的根本原因。
这个设计带来一个很反直觉的结论:Kafka 的 Broker 进程本身不需要在 JVM 堆内缓存数据,数据缓存全部交给操作系统的 PageCache。这样做的好处非常明显:
- JVM 堆不需要被海量消息塞满,GC 压力大幅降低,避免了“堆越大、GC 越痛”的经典问题;
- 操作系统的 PageCache 管理策略经过几十年打磨,内存回收、预读、回写都比自己用 Java 写一套缓存要可靠得多;
- 当消费者追赶不上生产速度时,旧数据会自然被换出 PageCache,从磁盘读取,不会导致 Broker OOM。
我在线上见过很多 Kafka 集群,堆内存给个 6GB 到 8GB 就够了,剩下的操作系统内存都留给 PageCache。如果你发现 Kafka Broker 的 JVM 堆占用率居高不下,先别急着加堆内存,看看是不是用了什么花里胡哨的缓存组件,或者消费者拉取参数设置得太奔放。
2.3 零拷贝:数据从磁盘到网卡的“直达通道”
Kafka 高性能的第三个关键机制是零拷贝(Zero Copy),尤其是消费场景下的sendfile系统调用。
如果不用零拷贝,一个消息从 Broker 磁盘发送到消费者网络,要走这么一条路:磁盘 → 内核缓冲区 → 用户态应用程序(JVM)→ 内核 Socket 缓冲区 → 网卡。数据在内核态和用户态之间来回拷贝了多次,每次拷贝都有 CPU 开销和上下文切换开销。
Kafka 是怎么做的?如果消息在 PageCache 中或者需要从磁盘读取后直接发给消费者,它会调用sendfile,让数据在内核态直接完成“磁盘文件 → Socket”的传输,完全绕过用户态。CPU 不再负责数据的复制,只负责控制传输,网卡可以直接从内核缓冲区读数据发出去。这就是“零拷贝”的含义:不是不拷贝,而是不在用户态和内核态之间来回倒腾。
实际使用中还有一点值得注意:Kafka 的零拷贝主要适用于消息不需要解压、不需要转换的场景,比如消费者拉取原始消息。如果消息做了压缩,Broker 可能需要先解压再发送,那就没法完全零拷贝了。所以生产环境中是否开启压缩、用什么压缩算法,不只影响存储,也影响消费链路的 CPU 开销。关于压缩的选择,我后面会详细说。
3. 高可靠的代价与权衡:副本、ISR 与 ACK 机制
3.1 副本因子与水印机制:leader 挂了怎么办
光快还不够,作为一个消息中间件,数据不能随便丢。Kafka 的高可用依赖多副本机制:每个分区有多个副本(Replica),其中一个是 Leader,其余是 Follower。所有读写请求都由 Leader 处理,Follower 只负责从 Leader 拉取数据保持同步。副本因子(replication.factor)一般建议 3,也就是 1 个 Leader 加 2 个 Follower,这样挂一台机器依然能保证数据完整。
这里有个核心概念叫ISR(In-Sync Replicas),也就是与 Leader 保持同步的副本集合。它不是一个固定的列表,而是动态维护的。Kafka 通过replica.lag.time.max.ms来判断一个 Follower 是否“掉队”:如果 Follower 超过这个时间没有跟 Leader 同步最新的消息,就会被踢出 ISR。默认值是 30 秒,我一般不会调它,因为太敏感会把短暂的网络抖动误判为副本故障,太迟钝又可能让副本长期落后。
每个分区还有一个高水位(High Watermark)的概念,它表示 ISR 中所有副本都已经同步到的位置,消费者只能消费到高水位之前的消息。高水位的作用是避免消费者读到“未来数据”或更糟糕的、之后可能被回滚的数据。Kafka 还有一个 LEO(Log End Offset),表示每个副本日志中下一条待写入消息的偏移量。副本同步过程中,Leader 会定期把高水位广播给 Follower,Follower 根据高水位来决定哪些消息可以向消费者暴露。
3.2 ACK 级别:性能与可靠性的取舍
生产者通过acks参数控制消息“写成功”的定义,这是 Kafka 性能和可靠性之间最核心的一个旋钮:
| acks 值 | 行为 | 数据可靠性 | 性能影响 | 适用场景 |
|---|---|---|---|---|
| 0 | 发出去就不管,不等待确认 | 最低,可能丢消息 | 最高 | 指标采集、日志等允许丢失场景 |
| 1 | Leader 写入本地日志即返回 | 中等,Leader 挂了可能丢 | 较高 | 大部分业务场景 |
| -1 / all | 等待 ISR 内所有副本都写入才返回 | 最高,几乎不丢 | 最低 | 金融、订单、对账等核心链路 |
选acks=all并不等于 100% 不丢,还得同时满足一个前提:ISR 里至少有一个副本处于同步状态。如果 ISR 只剩 Leader 自己,acks=all也就是 Leader 写完就返回了。为了堵住这个漏洞,Kafka 提供了min.insync.replicas参数。我建议核心业务至少设置为 2,这样当 ISR 中副本数量不足 2 时,生产请求会被拒绝,宁可让业务报错,也不能假装写成功了。
这里有一个很多新手容易踩的坑:把acks=all和min.insync.replicas=2一起设置后,如果集群只剩一个副本在线,生产者会持续报NotEnoughReplicasException。这不是配置错了,而是在用可用性换可靠性。你得在业务层面想清楚,这个场景下是让消息失败重试好,还是让消息悄悄丢失好。我的经验是:核心交易链路选前者,日志采集链路选后者。
3.3 幂等与事务:更高层次的保障
在acks基础上,Kafka 还提供了幂等生产者(enable.idempotence=true)。它的原理是给每条消息加一个序列号,Broker 端根据序列号去重。开启了幂等之后,生产者重试造成的重复消息可以在 Broker 端被过滤掉,不会出现“网络超时重发但实际第一条已经写入成功”导致的重复。3.0 之后的版本,幂等已经是默认开启的。
但幂等只能保证单个分区内不重复,跨分区的事务性写入需要 Kafka 事务 API。事务这个东西,能用上的人其实不多,但面试爱问。我的建议是:先搞清楚幂等解决什么问题、事务解决什么问题,别一上来就给业务套事务,事务会显著拉低吞吐,得不偿失。
4. 吞吐量背后的调度细节:批量、压缩与异步
4.1 生产者端:批量发送与缓冲池的设计
Kafka 生产者高性能的一个重要来源是批量发送。它不是来一条消息就发一条,而是先把消息攒在内存缓冲里,达到一定大小或一定时间后再批量发出。
这里有两个核心参数:
batch.size:一个批次的最大字节数,默认 16KB。这个值不是越大越好,太大了单批次的构建时间变长,反而增加延迟;太小了批量效果不明显。linger.ms:批次在内存里等待的时间,默认 0。很多人误解这个参数,以为设为 0 就完全不等待、来一条发一条。其实不是,设为 0 表示“只要有能发送的线程空出来就立即发送”,但如果发送线程正忙,消息照样会攒成批次。真正想提高吞吐,可以适当调到 5ms 到 20ms,用一点延迟换更高的吞吐。
生产者内存缓冲的总大小由buffer.memory控制,默认 32MB。如果生产速度过快、Broker 端 ack 太慢,缓冲会被填满,此时send()会阻塞(受max.block.ms控制,默认 60 秒)。你可以把buffer.memory调大,但根本上要排查是不是 Broker 端成了瓶颈,或者acks=all导致确认太慢。不要把生产者内存盲目调大,调得越大,OOM 时炸得越厉害。
4.2 压缩算法选型:存储、带宽与 CPU 的三方博弈
压缩是 Kafka 提升吞吐的又一利器,尤其在消息体比较大的场景。Kafka 支持gzip、snappy、lz4、zstd四种压缩算法。它们的关系大致是:zstd压缩率最高但 CPU 开销较大;snappy和lz4在压缩率和 CPU 开销之间比较均衡;gzip比较中庸。
我个人的选型策略是:
| 场景 | 推荐算法 | 理由 |
|---|---|---|
| 消息体小、要求低延迟 | 不压缩 | 节省 CPU,反正带宽够 |
| 日志类大数据量传输 | zstd | 压缩率最高,节省存储和带宽 |
| 通用业务消息 | lz4 或 snappy | 性能均衡,CPU 开销可控 |
需要强调的是,压缩并不只影响生产端,消费者端也必须解压,所以开启压缩会把一部分 CPU 开销从生产端挪到消费端。如果你发现消费者 CPU 很高,先看是不是消息压缩格式太耗 CPU。另外,Broker 端默认不会重新压缩消息,除非你配置了compression.type=producer以外的值,所以生产者和消费者用的压缩算法必须匹配。
4.3 消费者端:拉取模型与长轮询的优势
Kafka 选择的是拉取(Pull)模型,消费者主动去 Broker 拉数据,而不是 Broker 把数据推给消费者。这个选择非常关键:推模型的最大问题是不知道消费者的处理能力,容易把慢消费者压垮;拉模型让消费者自己掌握节奏,处理得快就多拉,处理得慢就少拉,天然具备背压能力。
消费者拉取的时候,fetch.min.bytes默认是 1 字节,也就是说只要有一点点数据就会返回;fetch.max.wait默认 500ms,表示如果数据不够,最多等 500ms 再返回。如果追求高吞吐,可以把fetch.min.bytes调到 1KB 或更大,让 Broker 攒一批数据再响应。如果追求低延迟,就把fetch.max.wait调小一点。这里又是一次典型的延迟与吞吐的权衡。
拉取模型还有一个隐藏优势:消费者重启后可以从任意 Offset 重新拉取历史消息,实现重放。这是推模型很难做到的。
5. 实战视角:从原理到线上问题排查
5.1 Kafka 消息延迟高的排查思路
“Kafka 消息延迟高”是网上出现频率非常高的问题。首先要定义清楚“延迟”发生在哪一段。我一般把链路分成三段来看:生产端到 Broker、Broker 内部、Broker 到消费端。
第一段:生产端是否阻塞。看生产者有没有buffer.memory打满的日志,看send()回调里有没有大量异常。如果linger.ms设置得很大,比如 100ms,那每批消息最长要等 100ms 才发出去,延迟自然高。我曾经接手过一个案例,同事为了吞吐把linger.ms调到了 300ms,结果业务方反馈端到端延迟 300ms+,这就是典型的用延迟换吞吐没换明白。
第二段:Broker 是否瓶颈。看系统指标:磁盘 IO 使用率、网络带宽、CPU 使用率。如果磁盘 IO 持续 100%,说明 PageCache 命中率太低,通常是消费者消费速度跟不上、数据被迫落盘导致读磁盘。 在高吞吐场景下,也要看看是不是单分区热点,所有流量都打在一个 Broker 上。
第三段:消费端是否落后。用命令查看消费组积压情况:
kafka-consumer-groups.sh --bootstrap-server localhost:9092 --describe --group your-group看LAG列,如果持续增长,说明消费速度小于生产速度。这时先看消费者单个消息处理耗时,再看消费者线程数。很多延迟不是 Kafka 的锅,而是消费者里的 SQL 慢查询、外部 HTTP 调用超时。
我遇到过最典型的一个“假 Kafka 延迟”案例:业务方把消费线程池从 10 个扩到 50 个,结果延迟不降反升。查了半天发现他们消费的消息里有一个字段需要调用远程服务补全,远程服务被 50 个线程打挂了,超时重试导致消费速度更慢。把远程调用改成批量异步之后,延迟立刻降下来了。消息中间件能解决的是“传输”问题,解决不了“消费逻辑本身太慢”的问题。
5.2 如何实现“延迟 30 分钟消费”
网上的热搜词里有“Kafka 如何延迟 30 分钟消费”,这是一个很典型的业务需求。先说结论:Kafka 天然不支持定时消息,它没有 RabbitMQ 那种延迟队列插件。要延迟消费,只能自己在业务层设计。
我常用的方案有三种:
方案一:消息带目标消费时间,消费者轮询判断。生产者发送消息时带上expectedConsumeTime,消费者拉取到消息后判断当前时间是否达到目标时间。如果没到,就把消息重新发回到一个“延迟缓冲主题”,或者用定时器稍后再处理。这个方案最简单,但要注意一个坑:不要用Thread.sleep在消费线程里等 30 分钟,这会阻塞消费者线程,导致心跳超时触发 Rebalance,整个消费组都乱了。正确做法是把未到期的消息转到一个内部重试主题,或者使用消费者pause/resume机制配合定时任务。
方案二:时间轮(TimingWheel)方案。思路是用一个时间轮把延迟消息先暂存起来,时间到了再投递到真实的业务 Topic。开源界有现成的实现,比如基于 Netty 的 HashedWheelTimer,Java 生态也可以用ScheduledExecutorService。这种方式吞吐高、延迟比较精准,但要自己维护组件,复杂度高一些。
方案三:引入外部延迟队列中间件。比如用 Redis 的 ZSet 按到期时间排序,每分钟轮询到期的消息再投递给 Kafka。这是很多团队在用的折中方案:Kafka 管主链路,Redis 管延迟调度。
我个人在项目里最常用的是方案一,因为它不引入额外组件,改造成本最低。只要把“未到期消息重新投递”和“到期判断”写清楚,就能在 Kafka 上实现按业务需求的延迟消费。延迟精度取决于轮询间隔,能做到秒级,对大部分业务足够了。
5.3 Kafka OOM 的常见原因与规避手段
“Kafka OOM”也是个高频词,但很多人没分清是 Broker OOM 还是客户端(生产者/消费者)OOM。这两者的处理方式完全不同。
Broker OOM其实很少见,因为数据主要在 PageCache 而不在堆内。如果 Broker 的 JVM 堆内存暴涨,通常不是消息数据造成的,而是某个消费者拉取参数太激进,fetch.max.bytes设得过大,或者开启了什么客户端缓存。遇到 Broker OOM,先jstack看线程堆栈,再jstat -gcutil看 GC 情况,基本能定位。
客户端 OOM才是重灾区。生产端最常见的是buffer.memory设置过大,同时max.block.ms设为了 -1(无限阻塞),导致生产者线程把堆内存耗光。消费端最常见的是单次poll拉取的消息太多,或者max.poll.records设得很大,而每条消息的处理逻辑又依赖大量内存。另外还有一个很隐蔽的坑:消费者消费 Kafka 时把消息对象存进了全局缓存或者让消息逃逸到长期存活的对象上,导致 JVM 堆内存不断膨胀。
规避手段归纳起来就三句话:分清楚堆内堆外;单次拉取量要与业务处理能力匹配;处理完的消息不要让引用滞留。如果还是炸,就用jmap拉堆转储文件分析大对象,别靠猜。
5.4 KRaft 模式与集群部署的演进
热词里频繁出现“Kafka 集群安装”“KRaft 模式”“Docker 部署”,这说明现在很多新同学开始上手时,面对的选择已经和几年前不一样了。
早期 Kafka 强依赖 ZooKeeper 管理元数据,部署一套集群要额外维护 ZK。从 2.8 版本开始引入 KRaft(Kafka Raft)模式,直接用 Raft 协议让 Kafka 自己管理元数据,逐步摆脱 ZooKeeper。3.3 版本之后 KRaft 已经可用于生产,新版 Kafka(4.x)则彻底移除了 ZooKeeper 依赖。所以新手上手时,我建议直接学 KRaft 模式,不用再碰 ZK 那套老架构。
KRaft 模式下,Broker 节点分成两类角色:Controller(控制器)和 Broker。Controller 负责管理元数据,通过 Raft 协议在多个 Controller 节点间达成共识。部署时你需要在配置里声明:
process.roles=broker,controller node.id=1 controller.quorum.voters=1@host1:9093,2@host2:9093,3@host3:9093如果节点只是 Broker,就把process.roles设为broker,并配置controller.quorum.voters指向 Controller 列表。用 Docker 部署时,还需要注意KAFKA_CFG_前缀的写法,每个环境变量的映射容易写错。网上的参考不少,但很多版本比较老,一定先确认你和用的 Kafka 镜像版本兼容。这里我想说的是:部署方式再怎么变,底层原理没变。你把分区、副本、ISR、页缓存的原理搞明白了,不管是用 Docker 还是二进制部署,都是“换汤不换药”。
6. 面试高频考点与实践参数速查
6.1 面试官爱问的几个经典问题
这里整理一些网上高频出现的 Kafka 面试题,我按自己的理解给一个回答思路,不展开写成八股文:
- Kafka 为什么这么快?回答要落在四个关键词上:分区并行、顺序写磁盘、PageCache、零拷贝。面试官如果追问,你能把
sendfile的数据路径画出来,基本就过关了。 - 如何保证消息不丢失?要从生产端、Broker、消费端三段分别回答:生产端
acks=all+ 重试 + 幂等;Broker 端副本因子 3 +min.insync.replicas=2;消费端手动提交 Offset,确保业务处理成功后再提交。 - 如何保证消息有序?一句话:Kafka 只保证单分区内有序。要全局有序,就只用一个分区,但吞吐会受限;业务上尽量按业务主键分发到同一分区。
- 为什么分区数不是越多越好?文件句柄多、元数据膨胀、Rebalance 时间长、副本同步压力大。这个题想听到的是“权衡”思维,不是让你背极限值。
- 消费者 Rebalance 是什么,怎么避免?Rebalance 是消费者组成员变化或订阅 Topic 变化时触发的分区重分配。频繁 Rebalance 的常见原因是消费者处理时间超过
max.poll.interval.ms,或者消费者频繁加入退出。解决办法是提高处理效率、适当调大这两个超时参数。
还有热词里提到“Pulsar 和 Kafka 哪个资料丰富一些”。客观说,Kafka 在国内的生态和资料量要丰富得多,无论是博客、面试题、开源组件还是搜索引擎结果的完整度,都要比 Pulsar 更成熟。Pulsar 在云原生架构上也有自己的优势,但多数团队在选型时还是会因为“好招人、好排查、好找资料”而选 Kafka。
6.2 核心参数速查表:不同场景的推荐值
我把这几个常用参数整理成速查表,方便你调优的时候直接翻:
| 参数 | 默认值 | 适用场景 | 备注 |
|---|---|---|---|
| acks | 1(新版本默认 all 场景视客户端版本而定) | 核心业务设 all,日志设 0 或 1 | 可靠性优先设 all |
| batch.size | 16384 | 吞吐优先可调大到 32768 | 不宜无限加大 |
| linger.ms | 0 | 延迟敏感保持 0,吞吐优先 5~20 | 延迟和吞吐的调节旋钮 |
| buffer.memory | 33554432 | 单机吞吐很大时调大 | 防止 OOM,别盲目调大 |
| compression.type | none | 大数据量推荐 zstd 或 lz4 | 消费端也会吃 CPU |
| fetch.min.bytes | 1 | 吞吐优先调大 | 配合 fetch.max.wait |
| fetch.max.wait | 500 | 低延迟调小 | 单位毫秒 |
| max.poll.records | 500 | 单条消息大时调小 | 避免一次拉取太多导致 OOM |
| min.insync.replicas | 1 | 核心业务设 2 | 配合 acks=all |
| replica.lag.time.max.ms | 30000 | 一般不用改 | 影响 ISR 判定 |
调参的时候记住一个原则:不要同时把所有参数都调到“最理想”的数值,调完一个参数要压测看整体效果。性能调优是环环相扣的,把batch.size调大了但不调linger.ms,实际效果可能微乎其微;把acks=all和min.insync.replicas=2配好了,但忘了在消费者端检查提交时机,数据照样可能重复或丢失。
最后说一点个人体会
如果你问我 Kafka 的高性能架构设计里最值得学习的是什么,我会说不是零拷贝,也不是 PageCache,而是它对“权衡”的理解。它用分区换来并行度,却承认了全局有序无法完美实现;它用副本和 ISR 换可靠性,却用acks参数把这个代价的旋钮交给了使用者;它用拉取模型换背压能力,却让消费者承担了更多调优责任。好架构不是把所有指标都做到最优,而是把每个选择的利弊亮出来,让使用者根据场景做决定。这也是我从一次次调优和排障中真正学到的:把原理吃透,参数就只是你手里的工具而已。