上个月某天凌晨告警群里突然炸了Kafka 集群某台 broker 磁盘使用率 95%马上要撑爆。我赶紧爬起来处理登上去一看日志目录里堆了一堆 .log 文件大的有几十 G小的也有几个 G。当时心里挺懵的Kafka 不是号称高吞吐、零丢失吗怎么磁盘会爆后来花了半宿才把问题搞定顺带把 Kafka 底层存储这块东西翻了个底朝天。今天把这次踩坑的收获还有我整理的一些核心原理全部分享出来能让你少走不少弯路。Kafka 到底是咋工作的先把 Kafka 的运行机制理一遍后面讲的所有东西都跟这个有关。Kafka 是个分布式消息系统由几类角色组成Producer生产者发消息Consumer消费者收消息Broker 是 Kafka 服务本身集群就是多台机器起多个 broker 进程。消息按 Topic主题分类一个 Topic 可以分成多个 Partition分区分区分布在不同 broker 上。Consumer 启动后加入 Consumer Group消费组组内每个 Partition 只被一个 Consumer 消费。早期版本依赖 ZooKeeper 存元数据、协调集群新版本2.8开始支持 KRaft 模式可以脱离 ZooKeeper 跑。Kafka 既能当消息队列解耦、削峰、异步也能当存储系统消息持久化、多副本还能当流处理平台的数据源。不过生产上最常见的用法还是第一种。明白了这些再看下面的内容就不费劲了。目录结构分区和分段存储Kafka 的存储是分层的几个核心概念得搞明白。主题Topic是逻辑上的概念下面分为多个分区Partition分区是物理存储的基本单位。每个分区对应磁盘上一个文件夹文件夹名字一般是主题名-分区编号比如order-topic-0。分区下面不是一个大文件而是分成了多个日志分段LogSegment每个 LogSegment 包含三类文件.log文件实际的消息数据.index文件偏移量索引.timeindex文件时间戳索引比如你看到这样的目录结构order-topic-0/ ├── 00000000000000000000.log ├── 00000000000000000000.index ├── 00000000000000000000.timeindex ├── 00000000000000123456.log ├── 00000000000000123456.index └── 00000000000000123456.timeindex每个 LogSegment 默认是 1G由log.segment.bytes控制当一个 .log 文件写满了就会滚动生成新的 LogSegment文件名就是第一条消息的 offset。这就解释了为什么我那台机器磁盘会爆分区里 .log 文件一个接一个副本数又是 1没法通过副本分散压力。消息顺序性的保证按 Key 路由很多人问 Kafka 怎么保证消息顺序严格来说Kafka 只能保证单个分区内的消息是有序的。生产者发消息的时候如果没有指定分区会走一个分区器Partitioner。Kafka 默认的分区策略是这样的ListPartitionInfopartitionscluster.partitionsForTopic(topic);returnMath.abs(key.hashCode())%partitions.size();也就是说同一个 Key 的消息一定会进同一个分区同一个分区内的消息是顺序写、顺序读的。我之前遇到过一个生产事故某个业务把消息的 Key 随机生成用的 UUID结果消息被分散到所有分区消费者并行消费的时候消息顺序就乱了。后来把 Key 改成业务主键 ID问题才解决。所以如果你对消息顺序有要求一定要保证 Key 是固定的业务标识比如订单 ID、用户 ID。生产者客户端的两个线程Kafka 的生产者客户端跑得这么溜靠的是两个线程的协作主线程和Sender 线程也叫发送线程。主线程负责把消息经过拦截器、序列化器、分区器之后放到消息累加器RecordAccumulator里。注意这个累加器不是简单地把消息扔进去它内部是按分区分组的每个分区对应一个双端队列队列里是一批一批的消息Batch。Sender 线程在后台不断轮询从累加器里把攒够的 Batch 发送到 broker。为什么要攒批减少网络 IO 次数啊一批发几十条肯定比一条一条发快得多。这里面有几个关键参数batch.size攒批的大小默认 16KBlinger.ms最长等待时间默认 0不等待compression.type压缩算法默认 none我之前调优过一个项目把batch.size调到 64KBlinger.ms调到 10ms吞吐量直接翻了 3 倍。代价是延迟会增加那么一点点对异步业务来说完全可以接受。消息丢失和重复消费坑过才知道痛运维最怕的就是消息丢失和重复消费这两个问题我都被坑过。重复消费的常见原因有这么几个Rebalance 的时候。比如一个消费者正在处理一条消息还没处理完这时候组里加了新消费者触发 Rebalance那条消息就被新消费者拿到又处理了一遍。先消费后提交 offset。如果消费完消息offset 还没提交服务挂了重启后这条消息会再被消费一次。生产者重试。生产者发送消息没收到 ACK会重试broker 端如果没做幂等控制就会收到重复消息。消息丢失更可怕常见原因有acks 没设置为 all。如果 acks1leader 副本写入后就返回 ACK这时候 follower 副本还没同步leader 挂了消息就丢了。消费者先提交 offset 再消费。offset 提交了但消息还没处理完这时候消费者挂了消息就丢了。生产者是 fire-and-forget 模式。只管发不管结果失败了就丢了。针对消息丢失Kafka 引入了幂等性机制每个生产者有个 PIDProducer ID每条消息有个序列号broker 端会校验重复的就丢弃。如果还不够用就上事务能保证消费-生产-提交 offset这三步是原子的。我那次故障之后就把核心业务的生产者 acks 全改成了 allretries 调到很大再加上幂等性基本上就不怕丢了。ISR、HW、LEO副本同步原理Kafka 副本同步这块有几个关键概念新手很容易搞混。每个 Partition 可以配置多个副本replication分布在不同 broker 上。副本里有个特殊的角色叫 leader处理所有读写请求其他副本叫 follower只负责从 leader 同步数据。ARAssigned Replicas是分区分配的所有副本ISRIn-Sync Replicas是和 leader 保持同步的副本集合OSR 是 Out-of-Sync Replicas就是掉队的那些。leader 副本负责维护 ISR 集合有个关键参数replica.lag.time.max.ms默认 10 秒。如果 follower 副本超过 10 秒没追上 leader就会被踢出 ISR。HWHigh Watermark是高水位消费者只能消费 HW 之前的消息。LEOLog End Offset是下一条要写入的消息位置。打个比方HW 就是水库的水位线消费者只能喝到水位线以下的水LEO 就是当前水库的入水口位置。ISR 集合里所有副本里最小的 LEO就是整个分区的 HW。这里有个坑如果 ISR 里的副本都掉线了Kafka 默认不允许从 OSR 里选 leaderunclean.leader.election.enablefalse这时候分区就不可用了。但如果你把这个参数改成 true虽然能恢复可用性但可能会丢数据。高性能的秘密这几个设计缺一不可Kafka 能扛这么高的吞吐量靠的是几个关键设计顺序写盘很多人以为 Kafka 是内存存储其实它是磁盘存储。但 Kafka 用的是顺序写盘机械磁盘顺序写盘的速度比随机写内存还快。.log 文件的写入是 append-only 模式永远在文件末尾追加。零拷贝传统的数据传输要走磁盘 → 内核缓冲区 → 用户缓冲区 → Socket 缓冲区 → 网卡。Kafka 用sendfile()系统调用数据直接从磁盘到网卡少了用户态和内核态的来回拷贝速度快很多。页缓存操作系统会把磁盘上的数据缓存在内存里这就是页缓存Kafka 充分利用了这一点。它不维护进程内的缓存而是依赖操作系统的页缓存这样重启后缓存还在。批量发送前面说过的攒批机制一批消息一次网络 IO 搞定。端到端压缩生产者压缩一批消息broker 存的是压缩后的数据消费者拿到的也是压缩数据整个链路不解压省 CPU 也省带宽。我之前做的性能测试单 broker 用普通机械盘顺序写也能达到几百 MB/s 的吞吐量。换成 SSD 更是直接起飞。日志清理策略磁盘满了不慌回到最开始那个磁盘撑爆的问题其实和日志清理策略有很大关系。Kafka 有两种日志清理策略由log.cleanup.policy控制delete默认按时间或大小删除过期日志可以配置这几个参数log.retention.hours保留时间默认 168 小时7 天log.retention.bytes分区最大大小默认 -1不限制log.segment.bytes单个日志分段大小默认 1GBcompact日志压缩相同 key 只保留最新 value这个适合那种只关心最新状态的场景比如用户配置信息。每个 key 对应的旧版本消息会被清理掉只保留最新的那条。我当时那个故障根本原因是没设置log.retention.bytes导致分区可以无限增长。设个上限比如log.retention.bytes107374182400100GB就保险多了。消费组协调那点事Kafka 的消费者是以消费组Consumer Group为单位管理的。一个分区只能被同一个消费组内的一个消费者消费但可以被多个不同的消费组消费。如果消费者数超过分区数多出来的消费者就分配不到任何分区闲在那儿。协调消费者和分区分配的是两个组件GroupCoordinator运行在 broker 上管理消费组ConsumerCoordinator运行在客户端和 GroupCoordinator 通信Rebalance 的过程有四个阶段FIND_COORDINATOR找到消费组对应的 GroupCoordinatorJOIN_GROUP所有消费者加入组选举 leaderSYNC_GROUPleader 把分区分配方案同步给所有人HEARTBEAT消费者定期发心跳保持成员关系Rebalance 期间所有消费者都会停止消费所以这个过程要尽量快。别随便加消费者也别让消费者频繁挂掉。我之前有个项目消费者用的是短链接方式处理完一条就断开重连结果 Rebalance 频繁触发整个消费组基本处于瘫痪状态。后来改成常驻进程稳稳当当。多线程消费的正确姿势KafkaConsumer 是非线程安全的一个实例不能被多个线程同时用。有两种多线程消费方案方案一每个线程一个 KafkaConsumer 实例最简单也最稳每个线程自己消费自己的分区彼此隔离。适合消费逻辑比较重的场景。方案二单线程拉取 线程池处理一个线程负责从 Kafka 拉消息放到内存队列里线程池负责处理消息。这样消费逻辑可以并行但拉取还是单线程的。我一般用方案一简单粗暴不容易出问题。方案二性能更好但要处理消息的顺序性、内存队列溢出等问题复杂度高不少。写在最后Kafka 这东西上手容易玩精通难。底层涉及到的东西太多光是副本同步、存储结构、消费者协调这几块就够啃上一阵。这次磁盘撑爆的事故给我最大的教训是任何中间件上线前都要把容量、清理策略、监控这些想清楚。不然半夜被叫起来处理故障那种感觉太酸爽了。下期我打算把 Kafka 的监控指标体系讲一下怎么用 JMX 抓数据怎么配告警规则这些实战经验全是坑里趟出来的感兴趣的可以留意。如果你也是踩过 Kafka 的坑或者对某个点有疑问欢迎在评论区交流。
一次 Kafka 集群磁盘写满后,我把这些底层原理全捋了一遍
上个月某天凌晨告警群里突然炸了Kafka 集群某台 broker 磁盘使用率 95%马上要撑爆。我赶紧爬起来处理登上去一看日志目录里堆了一堆 .log 文件大的有几十 G小的也有几个 G。当时心里挺懵的Kafka 不是号称高吞吐、零丢失吗怎么磁盘会爆后来花了半宿才把问题搞定顺带把 Kafka 底层存储这块东西翻了个底朝天。今天把这次踩坑的收获还有我整理的一些核心原理全部分享出来能让你少走不少弯路。Kafka 到底是咋工作的先把 Kafka 的运行机制理一遍后面讲的所有东西都跟这个有关。Kafka 是个分布式消息系统由几类角色组成Producer生产者发消息Consumer消费者收消息Broker 是 Kafka 服务本身集群就是多台机器起多个 broker 进程。消息按 Topic主题分类一个 Topic 可以分成多个 Partition分区分区分布在不同 broker 上。Consumer 启动后加入 Consumer Group消费组组内每个 Partition 只被一个 Consumer 消费。早期版本依赖 ZooKeeper 存元数据、协调集群新版本2.8开始支持 KRaft 模式可以脱离 ZooKeeper 跑。Kafka 既能当消息队列解耦、削峰、异步也能当存储系统消息持久化、多副本还能当流处理平台的数据源。不过生产上最常见的用法还是第一种。明白了这些再看下面的内容就不费劲了。目录结构分区和分段存储Kafka 的存储是分层的几个核心概念得搞明白。主题Topic是逻辑上的概念下面分为多个分区Partition分区是物理存储的基本单位。每个分区对应磁盘上一个文件夹文件夹名字一般是主题名-分区编号比如order-topic-0。分区下面不是一个大文件而是分成了多个日志分段LogSegment每个 LogSegment 包含三类文件.log文件实际的消息数据.index文件偏移量索引.timeindex文件时间戳索引比如你看到这样的目录结构order-topic-0/ ├── 00000000000000000000.log ├── 00000000000000000000.index ├── 00000000000000000000.timeindex ├── 00000000000000123456.log ├── 00000000000000123456.index └── 00000000000000123456.timeindex每个 LogSegment 默认是 1G由log.segment.bytes控制当一个 .log 文件写满了就会滚动生成新的 LogSegment文件名就是第一条消息的 offset。这就解释了为什么我那台机器磁盘会爆分区里 .log 文件一个接一个副本数又是 1没法通过副本分散压力。消息顺序性的保证按 Key 路由很多人问 Kafka 怎么保证消息顺序严格来说Kafka 只能保证单个分区内的消息是有序的。生产者发消息的时候如果没有指定分区会走一个分区器Partitioner。Kafka 默认的分区策略是这样的ListPartitionInfopartitionscluster.partitionsForTopic(topic);returnMath.abs(key.hashCode())%partitions.size();也就是说同一个 Key 的消息一定会进同一个分区同一个分区内的消息是顺序写、顺序读的。我之前遇到过一个生产事故某个业务把消息的 Key 随机生成用的 UUID结果消息被分散到所有分区消费者并行消费的时候消息顺序就乱了。后来把 Key 改成业务主键 ID问题才解决。所以如果你对消息顺序有要求一定要保证 Key 是固定的业务标识比如订单 ID、用户 ID。生产者客户端的两个线程Kafka 的生产者客户端跑得这么溜靠的是两个线程的协作主线程和Sender 线程也叫发送线程。主线程负责把消息经过拦截器、序列化器、分区器之后放到消息累加器RecordAccumulator里。注意这个累加器不是简单地把消息扔进去它内部是按分区分组的每个分区对应一个双端队列队列里是一批一批的消息Batch。Sender 线程在后台不断轮询从累加器里把攒够的 Batch 发送到 broker。为什么要攒批减少网络 IO 次数啊一批发几十条肯定比一条一条发快得多。这里面有几个关键参数batch.size攒批的大小默认 16KBlinger.ms最长等待时间默认 0不等待compression.type压缩算法默认 none我之前调优过一个项目把batch.size调到 64KBlinger.ms调到 10ms吞吐量直接翻了 3 倍。代价是延迟会增加那么一点点对异步业务来说完全可以接受。消息丢失和重复消费坑过才知道痛运维最怕的就是消息丢失和重复消费这两个问题我都被坑过。重复消费的常见原因有这么几个Rebalance 的时候。比如一个消费者正在处理一条消息还没处理完这时候组里加了新消费者触发 Rebalance那条消息就被新消费者拿到又处理了一遍。先消费后提交 offset。如果消费完消息offset 还没提交服务挂了重启后这条消息会再被消费一次。生产者重试。生产者发送消息没收到 ACK会重试broker 端如果没做幂等控制就会收到重复消息。消息丢失更可怕常见原因有acks 没设置为 all。如果 acks1leader 副本写入后就返回 ACK这时候 follower 副本还没同步leader 挂了消息就丢了。消费者先提交 offset 再消费。offset 提交了但消息还没处理完这时候消费者挂了消息就丢了。生产者是 fire-and-forget 模式。只管发不管结果失败了就丢了。针对消息丢失Kafka 引入了幂等性机制每个生产者有个 PIDProducer ID每条消息有个序列号broker 端会校验重复的就丢弃。如果还不够用就上事务能保证消费-生产-提交 offset这三步是原子的。我那次故障之后就把核心业务的生产者 acks 全改成了 allretries 调到很大再加上幂等性基本上就不怕丢了。ISR、HW、LEO副本同步原理Kafka 副本同步这块有几个关键概念新手很容易搞混。每个 Partition 可以配置多个副本replication分布在不同 broker 上。副本里有个特殊的角色叫 leader处理所有读写请求其他副本叫 follower只负责从 leader 同步数据。ARAssigned Replicas是分区分配的所有副本ISRIn-Sync Replicas是和 leader 保持同步的副本集合OSR 是 Out-of-Sync Replicas就是掉队的那些。leader 副本负责维护 ISR 集合有个关键参数replica.lag.time.max.ms默认 10 秒。如果 follower 副本超过 10 秒没追上 leader就会被踢出 ISR。HWHigh Watermark是高水位消费者只能消费 HW 之前的消息。LEOLog End Offset是下一条要写入的消息位置。打个比方HW 就是水库的水位线消费者只能喝到水位线以下的水LEO 就是当前水库的入水口位置。ISR 集合里所有副本里最小的 LEO就是整个分区的 HW。这里有个坑如果 ISR 里的副本都掉线了Kafka 默认不允许从 OSR 里选 leaderunclean.leader.election.enablefalse这时候分区就不可用了。但如果你把这个参数改成 true虽然能恢复可用性但可能会丢数据。高性能的秘密这几个设计缺一不可Kafka 能扛这么高的吞吐量靠的是几个关键设计顺序写盘很多人以为 Kafka 是内存存储其实它是磁盘存储。但 Kafka 用的是顺序写盘机械磁盘顺序写盘的速度比随机写内存还快。.log 文件的写入是 append-only 模式永远在文件末尾追加。零拷贝传统的数据传输要走磁盘 → 内核缓冲区 → 用户缓冲区 → Socket 缓冲区 → 网卡。Kafka 用sendfile()系统调用数据直接从磁盘到网卡少了用户态和内核态的来回拷贝速度快很多。页缓存操作系统会把磁盘上的数据缓存在内存里这就是页缓存Kafka 充分利用了这一点。它不维护进程内的缓存而是依赖操作系统的页缓存这样重启后缓存还在。批量发送前面说过的攒批机制一批消息一次网络 IO 搞定。端到端压缩生产者压缩一批消息broker 存的是压缩后的数据消费者拿到的也是压缩数据整个链路不解压省 CPU 也省带宽。我之前做的性能测试单 broker 用普通机械盘顺序写也能达到几百 MB/s 的吞吐量。换成 SSD 更是直接起飞。日志清理策略磁盘满了不慌回到最开始那个磁盘撑爆的问题其实和日志清理策略有很大关系。Kafka 有两种日志清理策略由log.cleanup.policy控制delete默认按时间或大小删除过期日志可以配置这几个参数log.retention.hours保留时间默认 168 小时7 天log.retention.bytes分区最大大小默认 -1不限制log.segment.bytes单个日志分段大小默认 1GBcompact日志压缩相同 key 只保留最新 value这个适合那种只关心最新状态的场景比如用户配置信息。每个 key 对应的旧版本消息会被清理掉只保留最新的那条。我当时那个故障根本原因是没设置log.retention.bytes导致分区可以无限增长。设个上限比如log.retention.bytes107374182400100GB就保险多了。消费组协调那点事Kafka 的消费者是以消费组Consumer Group为单位管理的。一个分区只能被同一个消费组内的一个消费者消费但可以被多个不同的消费组消费。如果消费者数超过分区数多出来的消费者就分配不到任何分区闲在那儿。协调消费者和分区分配的是两个组件GroupCoordinator运行在 broker 上管理消费组ConsumerCoordinator运行在客户端和 GroupCoordinator 通信Rebalance 的过程有四个阶段FIND_COORDINATOR找到消费组对应的 GroupCoordinatorJOIN_GROUP所有消费者加入组选举 leaderSYNC_GROUPleader 把分区分配方案同步给所有人HEARTBEAT消费者定期发心跳保持成员关系Rebalance 期间所有消费者都会停止消费所以这个过程要尽量快。别随便加消费者也别让消费者频繁挂掉。我之前有个项目消费者用的是短链接方式处理完一条就断开重连结果 Rebalance 频繁触发整个消费组基本处于瘫痪状态。后来改成常驻进程稳稳当当。多线程消费的正确姿势KafkaConsumer 是非线程安全的一个实例不能被多个线程同时用。有两种多线程消费方案方案一每个线程一个 KafkaConsumer 实例最简单也最稳每个线程自己消费自己的分区彼此隔离。适合消费逻辑比较重的场景。方案二单线程拉取 线程池处理一个线程负责从 Kafka 拉消息放到内存队列里线程池负责处理消息。这样消费逻辑可以并行但拉取还是单线程的。我一般用方案一简单粗暴不容易出问题。方案二性能更好但要处理消息的顺序性、内存队列溢出等问题复杂度高不少。写在最后Kafka 这东西上手容易玩精通难。底层涉及到的东西太多光是副本同步、存储结构、消费者协调这几块就够啃上一阵。这次磁盘撑爆的事故给我最大的教训是任何中间件上线前都要把容量、清理策略、监控这些想清楚。不然半夜被叫起来处理故障那种感觉太酸爽了。下期我打算把 Kafka 的监控指标体系讲一下怎么用 JMX 抓数据怎么配告警规则这些实战经验全是坑里趟出来的感兴趣的可以留意。如果你也是踩过 Kafka 的坑或者对某个点有疑问欢迎在评论区交流。