一次 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 默认的分区策略是这样的:
List<PartitionInfo> partitions = cluster.partitionsForTopic(topic);
return Math.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 调到 64KB,linger.ms 调到 10ms,吞吐量直接翻了 3 倍。代价是延迟会增加那么一点点,对异步业务来说完全可以接受。
消息丢失和重复消费:坑过才知道痛
运维最怕的就是消息丢失和重复消费,这两个问题我都被坑过。
重复消费的常见原因有这么几个:
- Rebalance 的时候。比如一个消费者正在处理一条消息,还没处理完,这时候组里加了新消费者,触发 Rebalance,那条消息就被新消费者拿到又处理了一遍。
- 先消费后提交 offset。如果消费完消息,offset 还没提交,服务挂了,重启后这条消息会再被消费一次。
- 生产者重试。生产者发送消息没收到 ACK,会重试,broker 端如果没做幂等控制,就会收到重复消息。
消息丢失更可怕,常见原因有:
- acks 没设置为 all。如果 acks=1,leader 副本写入后就返回 ACK,这时候 follower 副本还没同步,leader 挂了,消息就丢了。
- 消费者先提交 offset 再消费。offset 提交了但消息还没处理完,这时候消费者挂了,消息就丢了。
- 生产者是 fire-and-forget 模式。只管发不管结果,失败了就丢了。
针对消息丢失,Kafka 引入了幂等性机制,每个生产者有个 PID(Producer ID),每条消息有个序列号,broker 端会校验,重复的就丢弃。如果还不够用,就上事务,能保证"消费-生产-提交 offset"这三步是原子的。
我那次故障之后,就把核心业务的生产者 acks 全改成了 all,retries 调到很大,再加上幂等性,基本上就不怕丢了。
ISR、HW、LEO:副本同步原理
Kafka 副本同步这块有几个关键概念,新手很容易搞混。
每个 Partition 可以配置多个副本(replication),分布在不同 broker 上。副本里有个特殊的角色叫 leader,处理所有读写请求,其他副本叫 follower,只负责从 leader 同步数据。
AR(Assigned Replicas) 是分区分配的所有副本,ISR(In-Sync Replicas) 是和 leader 保持同步的副本集合,OSR 是 Out-of-Sync Replicas,就是掉队的那些。
leader 副本负责维护 ISR 集合,有个关键参数 replica.lag.time.max.ms,默认 10 秒。如果 follower 副本超过 10 秒没追上 leader,就会被踢出 ISR。
HW(High Watermark) 是高水位,消费者只能消费 HW 之前的消息。LEO(Log End Offset) 是下一条要写入的消息位置。
打个比方,HW 就是水库的水位线,消费者只能喝到水位线以下的水,LEO 就是当前水库的入水口位置。ISR 集合里所有副本里最小的 LEO,就是整个分区的 HW。
这里有个坑,如果 ISR 里的副本都掉线了,Kafka 默认不允许从 OSR 里选 leader(unclean.leader.election.enable=false),这时候分区就不可用了。但如果你把这个参数改成 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:单个日志分段大小,默认 1GB
compact:日志压缩,相同 key 只保留最新 value
这个适合那种"只关心最新状态"的场景,比如用户配置信息。每个 key 对应的旧版本消息会被清理掉,只保留最新的那条。
我当时那个故障,根本原因是没设置 log.retention.bytes,导致分区可以无限增长。设个上限,比如 log.retention.bytes=107374182400(100GB),就保险多了。
消费组协调那点事
Kafka 的消费者是以消费组(Consumer Group)为单位管理的。一个分区只能被同一个消费组内的一个消费者消费,但可以被多个不同的消费组消费。
如果消费者数超过分区数,多出来的消费者就分配不到任何分区,闲在那儿。
协调消费者和分区分配的是两个组件:
- GroupCoordinator:运行在 broker 上,管理消费组
- ConsumerCoordinator:运行在客户端,和 GroupCoordinator 通信
Rebalance 的过程有四个阶段:
- FIND_COORDINATOR:找到消费组对应的 GroupCoordinator
- JOIN_GROUP:所有消费者加入组,选举 leader
- SYNC_GROUP:leader 把分区分配方案同步给所有人
- HEARTBEAT:消费者定期发心跳,保持成员关系
Rebalance 期间所有消费者都会停止消费,所以这个过程要尽量快。别随便加消费者,也别让消费者频繁挂掉。
我之前有个项目,消费者用的是短链接方式(处理完一条就断开重连),结果 Rebalance 频繁触发,整个消费组基本处于瘫痪状态。后来改成常驻进程,稳稳当当。
多线程消费的正确姿势
KafkaConsumer 是非线程安全的,一个实例不能被多个线程同时用。
有两种多线程消费方案:
方案一:每个线程一个 KafkaConsumer 实例
最简单也最稳,每个线程自己消费自己的分区,彼此隔离。适合消费逻辑比较重的场景。
方案二:单线程拉取 + 线程池处理
一个线程负责从 Kafka 拉消息,放到内存队列里,线程池负责处理消息。这样消费逻辑可以并行,但拉取还是单线程的。
我一般用方案一,简单粗暴不容易出问题。方案二性能更好,但要处理消息的顺序性、内存队列溢出等问题,复杂度高不少。
写在最后
Kafka 这东西,上手容易,玩精通难。底层涉及到的东西太多,光是副本同步、存储结构、消费者协调这几块就够啃上一阵。
这次磁盘撑爆的事故给我最大的教训是:任何中间件上线前,都要把容量、清理策略、监控这些想清楚。不然半夜被叫起来处理故障,那种感觉太酸爽了。
下期我打算把 Kafka 的监控指标体系讲一下,怎么用 JMX 抓数据,怎么配告警规则,这些实战经验全是坑里趟出来的,感兴趣的可以留意。
如果你也是踩过 Kafka 的坑,或者对某个点有疑问,欢迎在评论区交流。
我是@运维躬行录,专注分享运维、监控、自动化踩坑经验。
公众号:耕云躬行录(回复 kafka 领这份《Kafka 运维实战手册》)
个人博客:躬行笔记
关注我,下期讲《Kafka 监控告警体系搭建:从 JMX 到 Grafana 实战》