先说个我在社区窜贴时经常看到的事儿:有人搜“kafuka”,一搜全是Apache Kafka的用法,还会顺手把名字拼成“卡夫卡”。其实Kafka这个名字本身就爱被拼错,但这并不影响它是目前后端技术栈里最值得花时间搞明白的中间件之一。我见过太多人把Kafka当成普通消息队列来学,上来就对着“Producer/Consumer”两个词死磕,结果连“为什么消费后消息还在”都解释不了。这篇内容不是官方文档翻译,而是我从零把Kafka用起来之后,整理的一套真正对实战有帮助的使用闭环,覆盖本地部署、客户端代码、高频排错和生产注意事项,适合刚接触消息中间件、或者已经在用但总踩坑的开发者。
1. 为什么我建议你从Kafka入手消息队列
1.1 Kafka到底解决什么问题
Kafka最常见的身份是“分布式消息队列”,但如果你只把它当队列,很多设计会显得很奇怪。举个最简单的例子:普通队列如RabbitMQ,消息被一个消费者取走后,这条消息基本就从队列里消失了;但Kafka里消费者只是移动了一下“读指针”,消息本身还留在磁盘上,过多久删除由保留策略决定。
这个特性让它实际上变成一个分布式日志系统。想象一本实时更新的账本,每个客户都可以翻到自己想看的位置,可以从头读,也可以从任意位置读,甚至把之前读过的账目再读一遍。Kafka就是这样一个“账本”,生产者往末尾追加记录,消费者各自记录自己看到哪一页,互不干扰。这种模型带来的三个核心能力正好是业务系统最常见的需求:
- 解耦:上游不用关心下游有多少个服务在等数据,下游也不用关心上游什么时候把数据送完,中间只用Topic这个逻辑通道衔接。
- 削峰填谷:突发流量来了,生产者只管往Kafka里写,消费者按自己的节奏慢慢处理,不会因为一瞬间的高流量把下游数据库打爆。
- 多路广播:同一个Topic可以被多个消费者组独立消费,每个组都能拿到全量数据,天然实现数据分发给不同处理链路,比如一个组做实时计算,另一个组做离线归档。
理解了这三个能力,你就明白为什么很多公司哪怕内部只有两三个服务,也坚持中间放一个Kafka:它不是为了让系统更复杂,而是为了让业务之间的数据通道变得可控。
1.2 别用“队列”的惯性思维去理解Kafka
初学者最容易卡住的一点,是拿着RabbitMQ的经验来套Kafka。在RabbitMQ里,一条消息被某个消费者签收后,基本不会给另一个消费者;但在Kafka里,同样一条消息完全可以被不同消费者组的消费者各自消费一遍。你甚至可以让同一个组里的多个消费者并行消费同一个分区,也能让另一个组重新从头消费所有历史数据。
这正是Kafka和“消息队列”在定位上最根本的分歧。Kafka把消息当成一段持久化的数据流,而消费者组不过是这段流上的一个“游标”。所以很多在队列里需要特殊机制才能解决的问题,在Kafka里天然就是顺理成章的:晚到的消费者想补数据,直接重置偏移量从头拉就行;业务要回放几天前的数据排查问题,把消费者组改成另一个保留着旧offset的实例即可。
当然,这种设计也有代价。消息不是“被消费完就删除”的,所以磁盘占用必须靠保留策略来管理;消息在分区间不是全局有序的,所以“严格有序”的需求得单独做设计。这些代价会在后面的章节里反复出现,你只要先记住一句话:Kafka不是作为一个消息中转站设计的,它是作为一个可回放、可多路消费的持久化日志设计的。
1.3 Kafka适合谁,不适合谁
按我的经验,Kafka最适合三类场景:一是大数据生态里的数据采集与流处理,它是Hadoop生态之外的绝对主角;二是业务系统需要对用户行为埋点做实时分析,数据量可能一天几亿条几十亿条,丢点数据可以接受但吞吐必须顶住;三是系统里有多条数据链路需要对同一份数据进行不同维度的处理,用Kafka做分发比各种回调硬编码要干净得多。
如果你只是想在一个很小的单体应用里给两个模块递个消息,每天就几万条,也不需要回放历史数据,那Kafka确实有点重。你完全可以用更轻量的Redis Stream或者RabbitMQ解决问题。但即便如此,我依然建议你花半天时间把Kafka用起来,因为它的设计思路会深刻影响你怎么看待“数据如何流动”这件事,这种收益会一直带到后续任何技术栈里。
2. 本机跑一个Kafka:新版KRaft模式真的省心
2.1 下载和版本选择
Kafka本机部署在过去是一件让新人头疼的事,因为旧教程都会让你先装一个ZooKeeper,Kafka元数据全存在ZK里,部署一个Kafka等于要维护两个分布式系统。好在这几年Kafka引入了KRaft模式,从3.3版本开始可以作为实验特性使用,3.5之后逐步稳定,到我现在用的3.7版本,单机实验已经完全不需要ZooKeeper了,一条命令起controller,一条命令起broker,整个人都舒坦不少。
下载时你会看到类似这样的文件:kafka_2.13-3.7.0.tgz。前面的2.13是Kafka服务端编译时用的Scala版本,换句话说它是一个构建标识,不是你业务代码需要关心的。真正要认准的是后面的3.7.0,这是Kafka的版本号,我建议你从官方下载页找一个3.6以上的版本。官网下载后,解压到你习惯的工作目录就行,比如我一般放在/opt/kafka下,存成kafka目录。
2.2 用三道命令完成启动
解压完之后,进入目录看一眼config/kraft/server.properties,这里是KRaft模式的主配置。默认情况下,里面会配置process.roles=broker,controller,意思是一个进程同时充当Kafka的数据节点和元数据管理节点,单机实验保证最省资源。监听端口默认是PLAINTEXT://:9092和CONTROLLER://:9093,想改再改,刚上手时直接用默认就好。
接着按顺序执行下面三条命令:
KAFKA_CLUSTER_ID="$(bin/kafka-storage.sh random-uuid)" bin/kafka-storage.sh format -t $KAFKA_CLUSTER_ID -c config/kraft/server.properties bin/kafka-server-start.sh config/kraft/server.properties第一条命令生成一个集群ID,一个Kafka集群在KRaft模式下靠这个ID来标识身份。第二条命令是格式化存储目录,本质上是在本地创建一个“空的日志空间”,这个操作只要做一次,重复执行会报错说存储目录已存在,不用慌,直接说明你之前格式化过了。第三条命令把服务拉起来,看到类似INFO Kafka Server started的日志,就说明成功了一半。
我踩过的一个坑是:格式化成功后,我又手滑执行了第二次,结果启动时一直报storage already exists,当时以为是配置坏了,折腾了十分钟才反应过来是格式化重复了。如果你真要重新初始化,先删掉配置文件里指定的log.dirs目录再执行格式化,别直接在原目录上来第二遍。
起来之后,可以用jps看看进程,也可以直接在另一个终端执行bin/kafka-broker-api-versions.sh --bootstrap-server localhost:9092,能正常返回一堆api版本就说明端口没问题。
2.3 用控制台脚本直观验证生产与消费
服务跑起来后,先创建主题:
bin/kafka-topics.sh --create --topic quickstart-events --partitions 3 --replication-factor 1 --bootstrap-server localhost:9092这里--partitions 3表示把Topic分成3个分区,--replication-factor 1表示每个分区只存1份副本。单机环境副本只能是1,不然没有多余节点可复制。然后开两个终端,一个运行生产者控制台,一个运行消费者控制台:
# 终端A bin/kafka-console-producer.sh --topic quickstart-events --bootstrap-server localhost:9092 # 终端B bin/kafka-console-consumer.sh --topic quickstart-events --from-beginning --bootstrap-server localhost:9092在终端A里打一行字回车,终端B马上就能看到。这虽然很基础,但请你务必亲手敲一遍,因为后续所有代码都是在替这两个控制台脚本干活。
你可以在消费者命令后面加一个--group demo-group,再启动第二个消费者实例,就会看到同一个组的多个消费者会把分区分配出去,一人负责一部分分区。Kafka控制台脚本的输出会把分配情况直接打出来,那种“分区在组内被平分”的直觉,比看十篇博客都管用。
2.4 必须建立的分区与副本概念
到这一步,你会看到Kafka里最核心的词汇:Topic、Partition、Replica、Offset、ConsumerGroup。我尽量用大白话给你拆清楚。
一个Topic就是一类消息的逻辑名称,比如orders就是订单消息。Topic底下被切成多个Partition,Partition才是真正存储数据的最小单位,也是并行处理的基本单位。为什么要有Partition?打个比方,一个图书馆如果只有一个柜台,所有还书的人都排一队,那再快也会排队;如果把柜台分成多个窗口,不同窗口处理不同区域的书,整体速度就上来了。Kafka也是这么干的:一个Topic只有单分区时,写和读都集中在一个文件上,吞吐自然有限;分成3个分区后,生产者和消费者就可以各自并行读取多个文件。
每条消息写入Partition时,Kafka会给它一个从0递增的编号,叫Offset。消费者读分区时记着“我读到Offset 5了”,下次接着从6读。所以同一个消费组里的消费者在同时消费多个分区时,其实谁读到哪儿是独立维护的。这里替大家提前划一个重点:Kafka只在分区内部保证消息顺序,跨分区不保证。所以“严格全局有序”在Kafka里是很奢侈的事情,业务设计时要尽量做到“同一个业务实体的消息尽可能落在同一个分区”,靠什么保证?后面写生产者代码时,我会说Key的作用。
Replica简单理解就是分区的“备份”。生产环境里副本数通常是3,这样Broker挂了一个也不会丢数据。副本有主从之分,所有读写都走主副本,从副本只负责同步数据。这些概念你都能在kafka-topics.sh --describe命令的输出里看到,分区编号、Leader副本、ISR列表清清楚楚。建议你现在就跑一下:
bin/kafka-topics.sh --describe --topic quickstart-events --bootstrap-server localhost:9092看到输出里Partition 0的Leader、Replicas、Isr之后,你对Kafka的分布式骨架就有了实感,后面再谈高可用就不会觉得空洞。
3. 手写第一个Java客户端:跨过框架遮蔽的坎
3.1 依赖只用kafka-clients就行
很多人一写Kafka客户端,就习惯性地去配Spring Boot的spring-kafka,然后把KafkaTemplate注入进来,几条注解一写就完事。但我的建议是,第一遍一定要绕过框架,直接用官方kafka-clients写。原因很简单:Spring Kafka帮你封装了太多东西,你以为自己在控制offset,其实框架在背后悄悄处理了;等出问题时,日志里的线程名、异常栈全被框架翻译了一遍,你根本不知道原生的Kafka在抱怨什么。
Maven项目里只需要加一个依赖:
<dependency> <groupId>org.apache.kafka</groupId> <artifactId>kafka-clients</artifactId> <version>3.7.0</version> </dependency>一定要记得mvn dependency:tree看一眼,别因为项目里别的依赖带了旧版本的kafka-clients,硬生生把客户端版本拉低,导致和服务端Protocol版本对不上。我见过不止一次,Spring Boot依赖传递进来的是2.x的kafka-clients,然后连Kafka 3.x,报一堆UnsupportedVersionException,其实都是版本冲突闹的。
3.2 生产者代码与核心参数解读
先看一段能直接跑起来的生产者代码:
import org.apache.kafka.clients.producer.KafkaProducer; import org.apache.kafka.clients.producer.ProducerRecord; import org.apache.kafka.clients.producer.ProducerConfig; import org.apache.kafka.common.serialization.StringSerializer; import java.util.Properties; public class DemoProducer { public static void main(String[] args) { Properties props = new Properties(); props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092"); props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName()); props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName()); props.put(ProducerConfig.ACKS_CONFIG, "all"); props.put(ProducerConfig.RETRIES_CONFIG, 3); props.put(ProducerConfig.LINGER_MS_CONFIG, 5); props.put(ProducerConfig.BATCH_SIZE_CONFIG, 16384); try (KafkaProducer<String, String> producer = new KafkaProducer<>(props)) { for (int i = 0; i < 100; i++) { producer.send(new ProducerRecord<>("quickstart-events", String.valueOf(i), "value-" + i)); } } } }这里最值得展开的是三个参数。
第一个是acks,它决定“消息要等到多少个副本确认才算发送成功”。acks=all表示要等所有ISR副本都写入成功才返回成功,这是最可靠的配置,代价是延迟会高一些;acks=0表示发出去就不管了,吞吐最高但可能丢消息;acks=1表示只要Leader副本写成功就算成功,丢数据的窗口比all大很多。单机实验时3个值都跑不出区别,一旦上生产,你就得认真考虑业务能不能接受丢失。
第二个是retries,它表示发送失败后的重试次数。Kafka的发送失败原因很多:网络闪断、Leader副本切换、元数据刷新超时等。设成3到5是比较务实的做法,不要设成0,也不要设成Integer.MAX_VALUE然后不设max.in.flight.requests.per.connection,否则重试可能导致本就可能的消息乱序更明显。
第三个是linger.ms和batch.size,这俩共同决定“攒多少消息再发出去”。Kafka生产者在内存里会先攒一批消息,攒够batch.size字节,或者超过linger.ms时间,再一次性发给服务端。这不是为了省事,而是为了大幅减少网络包数量,提升吞吐。如果你每条消息都想着马上发出去,把linger.ms设成0,性能会肉眼可见地下降。但设得太大又会抬高消息延迟,一般来说linger.ms=5到20是常见区间,具体值看你对延迟的容忍度。
还有个细节是ProducerRecord里的Key。如果你传了Key,Kafka会按Key的哈希值算分区,同一个Key基本永远落在同一个分区;不传Key,则按轮询方式均匀分布。所以讲“相同业务实体的消息要保证有序”时,你只需要在send时把订单ID作为Key传进去。代码里new ProducerRecord<>(topic, key, value)的用法就是这个目的。
3.3 消费者代码与offset提交的真相
消费者代码长这样:
import org.apache.kafka.clients.consumer.ConsumerConfig; import org.apache.kafka.clients.consumer.ConsumerRecord; import org.apache.kafka.clients.consumer.ConsumerRecords; import org.apache.kafka.clients.consumer.KafkaConsumer; import org.apache.kafka.common.serialization.StringDeserializer; import java.time.Duration; import java.util.Collections; import java.util.Properties; public class DemoConsumer { public static void main(String[] args) { Properties props = new Properties(); props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092"); props.put(ConsumerConfig.GROUP_ID_CONFIG, "demo-group"); props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName()); props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName()); props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, "false"); props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest"); try (KafkaConsumer<String, String> consumer = new KafkaConsumer<>(props)) { consumer.subscribe(Collections.singletonList("quickstart-events")); while (true) { ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(1000)); for (ConsumerRecord<String, String> record : records) { System.out.printf("offset=%d, key=%s, value=%s%n", record.offset(), record.key(), record.value()); } consumer.commitSync(); } } } }subscribe的意思是让Kafka帮你管理分区分配。同一个消费者组里,如果起了两个实例,Kafka会自动把主题下的分区分配给这两个实例,尽量公平。比如主题有3个分区,两个消费者就是一人拿1/2个,而不是每条消息都被两个消费者重复处理。
消费者最绕的地方是offset怎么提交。代码里enable.auto.commit=false,然后每处理完一批消息后手动commitSync(),这个操作告诉服务端“这批消息我已经处理完了,下次从这里继续读”。如果改成enable.auto.commit=true,Kafka会在后台每隔一段时间自动提交当前已拉取消息的offset——注意是“已拉取”,不是“已处理”。一旦消息拉到本地后业务处理还没结束,却发生了宕机,重平衡后就会有一部分消息被重复消费。所以我在实战里几乎永远把自动提交关掉,手动在业务处理完成后再提交,虽然代码多两行,但至少你明确知道重复消费的边界在哪里。
另外一个auto.offset.reset参数也容易懵。它只在当前消费组没有已提交offset时生效——比如这个组第一次消费这个Topic,或者offset数据过期被清了。earliest表示从最早的数据开始读,latest表示从最新位置开始读。如果你想验证“历史消息回放”,想从头消费,就可以用一个全新的消费者组ID加earliest。
这里再提醒你一个很多新人忽略的问题:poll(Duration.ofMillis(1000))的意思是“至少等1秒,拿不到数据就返回空集合”,不是“每1秒只poll一次”。Kafka的消费者需要在循环里快速调用poll,因为心跳和分区分配这些操作都依赖这个循环的持续运转。如果你在poll和下次poll之间执行了一段很长很长的业务逻辑,超过max.poll.interval.ms,Kafka会认为这个消费者出问题了,把它踢出消费组,然后触发rebalance,这简直是重复消费问题的最高频来源,后面专门讲。
3.4 跑通后你该做的一次小实验
代码跑通之后,别急着收工,我建议你做一次10分钟的小实验:保持生产者发送100条消息,先启动一个消费者把数据拉完,然后关闭这个消费者,再换一个新的消费者组ID,配上earliest启动。你会看到新消费者把这100条消息从头又消费了一遍。这一刻你对Kafka“消息还在磁盘上”这个本质的体会,会比任何文档都深刻。
反过来,如果一直用同一个消费者组ID,重启消费者后它会从上次提交的offset继续消费,不会把老数据再发一遍。这就是“消费位置由消费者组管理”的具体表现。把这两个现象对比着看,你就明白为什么Kafka被称为“分布式日志系统”,而不是传统意义上的消息队列。
4. 三个高频翻车现场:从日志到解决方案的完整排查链
4.1 定位问题的通用起手式
Kafka用多了之后,你会发现绝大多数使用问题不是Kafka本身坏了,而是你根本没按照它的用法去配参数。我在排错时的通用顺序是:先看服务端日志,再看主题元数据,最后看客户端日志。
服务端日志默认在logs/server.log,里面有所有错误信息,比如磁盘空间不足、分区Leader切换、网络连接关闭等。主题元数据用我之前说的kafka-topics.sh --describe看,重点看分区数和ISR列表是不是健康。客户端日志就要注意,有时候错误被框架包装了,最有效的办法是先把Spring Kafka摘掉,用原生客户端最小复现一次。
下面三个场景,是我自己在真实环境里碰到过,也在公司带新人时几乎每个月都看得到的高频问题。
4.2 场景一:消息发送一直超时
现象很典型:生产者代码启动后,控制台抛出的异常是org.apache.kafka.common.errors.TimeoutException: Topic quickstart-events not present in metadata after 60000 ms或者类似的“元数据拉取超时”。
很多人第一反应是“Broker挂了,Kafka要炸”,但排查下来十有八九是网络或配置问题。这个异常的意思是:生产者启动后需要先向bootstrap.servers里的地址发一个元数据请求,拿到你要发送的Topic在哪里、Leader在哪。如果这个请求在一个超时时间内一直没成功,就会报这个错。
排查链路是:
- 确认
bootstrap.servers里的地址能和Kafka地址对上。本机写localhost:9092没问题,但如果Kafka跑在Docker容器里,容器和宿主机之间直接用localhost就很容易踩坑,要换成宿主机内网IP。 - 用
telnet localhost 9092或者nc -vz localhost 9092确认端口通不通。如果不通,排查防火墙和Kafka的listeners配置。 - 如果端口通但还是超时,去确认
advertised.listeners。这个参数是Kafka对外广播的地址,客户端拿到后会用这个地址连接Broker。很多云主机上,这个值默认是localhost,客户端的bootstrap地址写的是公网IP,但Broker广播的是localhost,于是客户端在拿到元数据后又去连localhost,自然失败。
我最开始栽在第三步上。本地VM里Kafka跑得好好的,一放到云服务器上就超时,最后就是advertised.listeners没改为云服务器内网IP。看到服务端日志里的Advertised listeners一行,会瞬间明白问题所在。
4.3 场景二:消费端无限重复消费
还有个更揪心的场景:消息明明处理成功了,日志也打出去了,但重启服务后,同一批消息又来了一遍,而且每次重启都在老地方重来。这是典型的offset没提交成功,或者提交了但没生效。
排查链路第一步:检查enable.auto.commit。如果设成true,那就看auto.commit.interval.ms默认值。这里有个很隐蔽的坑:enable.auto.commit=true时,Kafka是每隔一段时间自动提交“当前已poll的消息offset”,不是提交“已处理完的消息offset”。假设你一次poll拉回100条,刚处理到第10条时自动提交线程就把这批消息的offset都提交了,结果第10条之后的处理逻辑抛出异常,程序退出。下次重启时会从已提交的offset继续消费,那90条消息就永远处理不了了。这属于丢数据,不是重复消费;如果自动提交间隔比较长,处理完了但offset没来及提交,又宕机了,就是重复消费。
排查链路第二步:如果你已经用了手动提交commitSync,但仍然重复,那多半是处理逻辑耗时太长,触发了rebalance。每个消费者在两次poll之间都在执行业务处理,Kafka会检测“这个消费者还活着吗”,活着的标准就包括它有没有在max.poll.interval.ms内继续poll。一旦超时,Kafka会把消费者踢出组,区分配给其他成员;重新开启后从最近一次提交的offset开始读,这一批正在处理但没有提交的,自然就又读一遍。
解决办法很直接:要么把max.poll.interval.ms调大,比如从默认的300000改成600000,让单次处理有更多时间;要么把max.poll.records调小,比如从默认的500改成100,让每次poll带回的消息量变小,既是控制内存也是减少单次处理时长。最稳妥的还是手动提交,并且确保“业务处理成功之后再commit”,顺序别搞反。
4.4 场景三:有一台Broker磁盘突然打满
单机实验没有这个问题,但一旦你起了多节点集群,Kafka某台磁盘打满会直接导致生产客户端大量报Disk usage low或NotEnoughReplicasException。本质原因是Kafka消息全部落盘,而落盘所在的分区有一个清洗和过期机制:默认log.retention.hours=168,也就是保留7天,过期后由日志清理线程删除。如果这7天里你狂灌数据,磁盘就被吃光了。
排查思路要分成两步。先看是不是真的“写太多了”。用JMX指标里的kafka.server:type=BrokerTopicMetrics,name=BytesInPerSec看写入速率,估算每天写入量,再对比磁盘容量。如果确实超了,优先调整Topic级别的retention.ms——比如业务只允许保留3天的日志,就不要留着默认7天。命令像这样:
bin/kafka-configs.sh --bootstrap-server localhost:9092 --entity-type topics --entity-name quickstart-events --alter --add-config retention.ms=259200000再退一步,如果Topic本身确实需要保留很久,但对少量丢失不敏感,可以考虑开启Kafka的压缩(log.cleanup.policy=compact),用“按Key保留最新值”的策略替代“按时间删除所有旧值”。很多大数据工程师会把长期保留的Topic设置成compact,只保留每个Key的最新状态,这样磁盘占用就能压下来。
从实战角度看,磁盘打满背后还有个让人容易忽略的连锁反应:当磁盘剩余空间不够时,Kafka会拒绝新写入,但不会自动把分区迁移到别的Broker。你必须要自己通过kafka-reassign-partitions.sh扩容或迁移,比如给新Broker加磁盘,再把高占用分区迁移过去。这个动作建议在集群搭建时就预先设计好,别等磁盘满了再慌忙操作。
5. 上生产前必须说的几句大实话
5.1 默认配置真的不能直接用
Kafka的很多默认值都偏保守或偏开发场景,直接上生产会撞出你想象不到的问题。举个我常见的例子:log.retention.hours默认168小时,也就是7天,如果你业务数据只允许保留24小时,那默认值就会白白占用大量磁盘;如果把数据误当成必须永久保留的也没好到哪去,磁盘满了服务直接罢工。
还有num.partitions默认是1,意味着新建Topic如果不指定分区数就只有一个分区,单分区的写并发和读并发都会受限。我建议在生产环境给Broker配一个Topic级别的默认值,比如num.partitions=3,避免业务部门创建Topic时忘记分区数,长期低性能运行。同理,default.replication.factor默认是1,生产环境至少设成3,不然某台Broker挂掉时,这台机器上的分区在只剩ISR不到min.insync.replicas阈值时会停止读写。
这些配置不是越多越好,但生产环境里确实有几个必须过一遍:log.retention.hours、num.partitions、default.replication.factor、log.segment.bytes。log.segment.bytes我一般会从默认的1GB调小到512MB,因为段文件越小,清理线程做日志删除时越平缓,不至于让磁盘IO出现尖峰;当然太小又会增加文件数量,增加打开句柄,建议不要低于256MB。
5.2 数据可靠性和性能的取舍
Kafka使用里最核心的权衡,就是可靠性到底做到什么程度。如果你只想要吞吐,acks=0加linger.ms调高,每秒可以刷出惊人的消息量,但机器一断电,缓存里的消息就可能没了。如果一定要不丢消息,acks=all加min.insync.replicas=2,再配合生产者的重试机制,才是比较稳妥的组合。
这里说下min.insync.replicas:它表示“至少要有多少个ISR副本确认写入,才认为这条消息发送成功”。如果副本数为3,min.insync.replicas=2意味着至少要等两个副本都同步完成,Leader才给生产者返回成功;如果ISR成员数目低于2,写入会被直接拒绝。很多人只把acks调到all就不管了,但你没设min.insync.replicas,默认是1,那就等于只等Leader确认,和数据要求是矛盾的。
终端用户可能觉得这些参数很绕,但我用一个真实案例告诉你为什么重要:之前有个团队部署三副本Kafka,acks=all、min.insync.replicas忘设,某天一台Broker磁盘损坏,Leader主动把两个同步副本踢出ISR,然后继续用孤零零的Leader给客户端返回成功。等那台坏机器上的数据彻底丢了,整个Topic的副本还是1个,但客户端毫无感知,业务已经往里面写了一天的数据。后来他们把所有Topic的ISR下限调整为2,再配合告警,才算真正把“不丢消息”落实下来。
5.3 分区数不是越多越好
我见过一些工程师一听分区能提高并行度,就恨不得把每个Topic都设成几十个分区。分区数提升吞吐的前提是,生产者和消费者确实能达到足够并行。生产者默认每个分区都会维护一个队列,分区多、发送侧为了按分区排队,管理成本也高;消费者线程数如果小于分区数,多余分区也只能闲置。更麻烦的是,分区一旦增多,发生rebalance时协调时间会变长,因为这些分区要重新分配。
我的经验法则是:先估算单分区能扛的吞吐量,再去定分区数。常见配置下单分区QPS能做到几千到上万条,如果你的业务峰值是5万QPS,3到8个分区通常够了。再结合消费者的实例数,比如你有5个消费者实例,分区数最好不要是质数之类的偏数,尽量能被实例个数整除或者为实例数的整数倍,这样分配会更均匀。刚上线的主题建议分区数控制在6到12之间,不要因为看着别人50个分区就觉得越多越好。等吞吐需求明显增长后,再通过Kafka 3.x对分区扩容的方式动态调整,但那是后话。
5.4 这些需求不建议用Kafka硬扛
最后再泼几盆冷水。Kafka并不是万能的,有两个需求我强烈不建议直接拿Kafka实现。
一个是严格的全局有序。刚才说了,Kafka只保证分区内有序,你想让所有消息都按到达顺序被处理,宁可放弃分区并行度,用单分区主题,或者把数据按业务ID路由到固定分区,在业务侧再做排序聚合。否则就会出现“A消息先发,B消息后发,A和B在不同分区,B反而被先消费”的乱序问题。
另一个是延迟消息/定时消息。Kafka没有内置的延迟队列能力,你要做订单超时关闭、定时任务调度,与其自己造轮子在Kafka上穷折腾,不如用Redis的过期监听,或者直接用现成的消息调度平台。Kafka擅长的是“高吞吐、可回放的数据流”,不是“精确到秒级的定时投递”。工具选型最怕的就是硬把业务场景往一个中间件上生搬硬套,这种教训比配置参数踩坑还要贵。
我在实际项目里最后养成了一个习惯:任何团队用Kafka之前,我都会先帮他们画一张简单的表格,包括Topic名称、分区数、副本数、保留时间、允许丢失多少、有序性要求。这张表填完,90%的后期灾难都能规避。新人学Kafka特别容易陷在“API怎么调”里,但真正的分界线往往是“业务给Kafka提出什么要求,Kafka用它的机制怎么响应”。上手之后多想想这一层,你就不会只在“怎么用”上打转,而是真正摸到这套数据流的脾气。