ARTICLE · INTELLIGENCE

战地情报 · 详情页

来自尧图项目组的一线实战观察与深度解析

Kafka核心概念与实战:从Topic、Partition到高吞吐原理

Kafka核心概念与实战:从Topic、Partition到高吞吐原理 先说明一下很多初学者第一次接触 Kafka 时最容易卡住的地方不是代码而是脑子里对“消息队列到底在干嘛”没有画面感。我见过不少朋友把官方文档翻了好几遍Producer、Consumer、Broker 这些词都认识但一说到“为什么需要 Topic 和 Partition”“消费者组解决什么问题”就开始含糊。这篇内容就是冲着解决这个问题来的从核心概念讲到环境搭建从命令行实操写到 Java API再补上重复消费、消息延迟、连接报错这些高频坑最后给出一份可以直接背的面试速查。不管你是刚准备入行的新手还是被迫接手 Kafka 项目的后端开发照着这篇文章走一遍至少能独立搭出一个可用的环境并且能说清楚它到底是怎么工作的。1. 先搞懂 Kafka 到底是什么1.1 消息队列在解决什么问题要理解 Kafka先理解它出现之前的困境。假设你有一套电商系统用户下单后需要触发扣库存、发短信、送积分、更新搜索索引四个动作。最简单的方式就是在下单接口里同步调这四个服务但这种做法有两个问题一是接口响应时间被最慢的那个服务拖死二是某个下游服务挂了整个下单流程直接失败。消息队列的解法是“中间加一层”。下单服务只把订单事件写入消息队列然后立刻返回成功扣库存、发短信这些服务各自去队列里取事件处理。这样下单接口不再关心下游服务的状态下游服务即使短暂不可用消息也会待在队列里等着恢复后继续消费。这就是消息队列最核心的价值异步解耦和削峰填谷。Kafka 能在这类系统中胜出靠的是比 RabbitMQ 高一个量级的吞吐能力以及消息从生产到消费的完整持久化机制。1.2 核心概念逐个拆解Kafka 里的概念不算多但每个都对应一个明确的角色。我按数据流的方向从前到后讲。Broker一台运行 Kafka 服务的机器就是一个 Broker多个 Broker 组成集群。数据不是存在某个中心节点而是分散在各 Broker 上。Topic消息的分类名类似数据库里的表。比如订单系统可以建一个order-topic用户行为系统建一个user-log。Partition一个 Topic 可以切分成多个分区分区才是真正存储数据的单元。消息写入分区时是追加写每个消息会在分区内获得唯一的偏移量 Offset。Replica分区的副本用于故障容灾。Leader 副本负责读写请求Follower 副本负责同步数据Leader 挂了就从 Follower 里重新选主。Producer消息生产者负责把消息发送到指定 Topic 的某个分区。Consumer 与 Consumer Group消费者从分区拉取数据。多个消费者可以组成一个消费组组内每个消费者负责一个或多个分区保证每条消息只被组内一个消费者处理。Offset消费者在分区内的消费位置相当于书签。Kafka 通过 Offset 记录消费者读到了哪一条重启后能接着读。这些概念之间的关系可以理解为“Topic 是一条大河Partition 是大河分出的河道消息是河道里的船消费者组是岸上的装卸队每条船装了什么由生产者说了算”。1.3 Kafka 的“快”是有原因的很多人第一次测 Kafka 性能时会被它的吞吐惊到单机轻松跑到每秒几十万条消息。这背后不是玄学而是三个关键设计。第一是顺序写盘。普通消息队列写入磁盘是随机 IO磁头到处找位置慢到没法看。Kafka 每个分区都是一个追加日志文件新消息永远写到文件尾部顺序写的速度接近内存操作。第二是零拷贝技术。消费者读取消息时传统做法是磁盘到内核缓冲区、到用户程序、再拷贝回内核、最后通过网卡发出绕了四趟。Kafka 利用操作系统的 sendfile 系统调用数据从磁盘直接进入网卡发送中间跳过用户态拷贝省下大量 CPU 开销。第三是批量发送与压缩。Producer 端并不会每来一条消息就发一次网络请求而是攒到一定数量或达到时间阈值后批量发送配合 gzip、snappy 或 lz4 压缩网络传输效率大幅提升。理解了这三个机制很多问题就自然有解了。比如“Kafka 读写最大值与硬件有什么关系”这种问题其实核心就是看磁盘的顺序读写性能和网卡带宽。单块 SATA 固态的顺序写大概在 500 MB/s 左右NVMe 固态能到 3000 MB/s 以上而 Kafka 的吞吐瓶颈绝大多数场景都在网络而不是磁盘。2. 环境准备与安装部署2.1 前置条件与版本选择安装 Kafka 之前先理清楚两件事JDK 版本和 Kafka 版本。Kafka 是 Java 写的运行环境必须有 JDK。目前主流 Kafka 版本是 3.x而 3.x 不同小版本对 JDK 的要求不一样比如 Kafka 3.4 之前用 JDK 8/11 都没问题3.5 之后推荐 JDK 11 以上。我个人的建议是直接装 JDK 17长期维护版本兼容性也最好实测跑 Kafka 3.6 和 3.7 都很稳。版本选择上不要盲目追新。生产环境建议选择某个 3.x 大版本的较新小版本因为 Apache Kafka 的小版本迭代很快新特性往往伴随新问题。比如早期 3.0 刚引入 KIP-500 的 KRaft 模式时很多功能还不完善如果直接上生产会踩不少坑。稳妥的选择是 3.4 或 3.6 这种相对成熟的版本。另外必须提一下历史包袱老的 Kafka 版本强依赖 ZooKeeper 来保存元数据和做集群协调这也是很多新手配置时会困惑的地方——为什么装 Kafka 还要装 ZK从 Kafka 2.8 开始官方引入了 KRaft 模式可以在没有 ZooKeeper 的情况下运行到 3.5 之后 KRaft 在社区被标记为生产可用3.7 之后更是逐步向“完全移除 ZK”演进。不过大量现存生产环境仍然是 ZooKeeper 模式两种模式你最好都了解。注意KRaft 模式下集群节点分为 Controller 和 Broker 两种角色Controller 负责元数据管理Broker 负责数据读写。如果只是本地测试可以配置一个节点同时承担两种角色。2.2 单机安装实操记录这里我给出一份可以直接照抄的单机安装流程以 Kafka 3.6.2 为例环境是 CentOS 7 或者 Ubuntu 20.04 都适用。# 1. 下载并解压建议用官方镜像站 wget https://archive.apache.org/dist/kafka/3.6.2/kafka_2.13-3.6.2.tgz tar -xzf kafka_2.13-3.6.2.tgz cd kafka_2.13-3.6.2解压后目录结构里有两个关键目录bin存放所有命令行脚本config存放配置文件。先改配置文件config/server.properties中的三个核心项# 每个 Broker 的唯一标识集群中不能重复 broker.id0 # 监听地址默认只监听 localhost生产环境改内网 IP listenersPLAINTEXT://localhost:9092 # 日志存储目录默认在 /tmp/kafka-logs建议改成独立数据盘 log.dirs/data/kafka-logs改完之后启动因为这个版本我用的 KRaft 模式所以不需要单独启 ZooKeeper但首次启动需要执行一个格式化步骤# 生成集群 ID 并格式化存储目录 KAFKA_CLUSTER_ID$(bin/kafka-storage.sh random-uuid) bin/kafka-storage.sh format -t $KAFKA_CLUSTER_ID -c config/server.properties # 前台启动 Kafka方便看日志确认没问题后再用 systemd 托管 bin/kafka-server-start.sh config/server.properties看到日志中输出started (kafka.server.KafkaRaftServer)就代表启动成功了。首次跑通后的第一件事我建议立刻执行下面这条命令验证一下bin/kafka-topics.sh --bootstrap-server localhost:9092 --create --topic test-topic --partitions 3 --replication-factor 1 bin/kafka-topics.sh --bootstrap-server localhost:9092 --describe --topic test-topic能正常创建并描述 Topic说明整个链路已经通了。Windows 用户下载解压后执行的是bin\windows\目录下的.bat脚本配环境变量 JAVA_HOME 时要注意安装 JDK 时别选 JRE。2.3 3 节点集群部署要点真实生产环境不可能只跑单机因为单机的宕机意味着整个通道瘫痪。搭建 3 节点集群时三台机器上装好 Kafka 后需要改动的地方只有server.properties逐台配置如下。节点一broker.id1 listenersPLAINTEXT://192.168.1.11:9092 log.dirs/data/kafka-logs # 集群节点互相通信需要以下两项 controller.quorum.voters1192.168.1.11:9093,2192.168.1.12:9093,3192.168.1.13:9093 advertised.listenersPLAINTEXT://192.168.1.11:9092节点二和节点三结构相同只改 broker.id、listeners 和 advertised.listeners 中的 IP。关键点有两个一是controller.quorum.voters这项三台机器配的内容必须完全一致否则节点间无法建立共识二是每台机器防火墙要放通 9092 端口如果 KRaft 模式的话还要放通 9093 端口。三台都启动后执行下面的命令确认集群健康状态bin/kafka-metadata.sh --bootstrap-server 192.168.1.11:9092,192.168.1.12:9092,192.168.1.13:9092 --describe能看到三个 Broker 都处于 Active 状态就是成功了。创建 Topic 时设置--replication-factor 3这样每个分区会有 3 个副本允许最多挂掉两台节点而不丢失消息。副本数不是越大越好因为副本同步会占用额外磁盘和网络资源生产场景默认 3 就够用了。3. 从命令行到代码的实操上手3.1 命令行快速跑通生产消费Kafka 自带的命令行工具是最直观的学习途径也是排查问题时的利器。先开一个终端启动生产者bin/kafka-console-producer.sh --bootstrap-server localhost:9092 --topic test-topic启动后光标停在输入状态输入一行字回车就完成了一条消息的发送。再开另一个终端启动消费者bin/kafka-console-consumer.sh --bootstrap-server localhost:9092 --topic test-topic --from-beginning这条命令加了--from-beginning意思是消费者从头开始读这个 Topic 里已有的所有消息。不加的话消费者只会从启动之后新产生的消息开始读这个细节在排查问题时特别重要。我建议你按下面这套动作做一遍实验先后台启动消费者再发 10 条消息记录消费者收到了哪些重启消费者再加--from-beginning观察它是否把历史消息再读一遍。这个过程能帮你直观理解 Offset 的作用——消费者的位置是持久化的但从头消费会忽略已有位置强行重置。3.2 用 Java 写一个生产者命令行只是辅助工具实际业务一定得写代码。以 Spring Boot 项目为例先引入依赖dependency groupIdorg.springframework.kafka/groupId artifactIdspring-kafka/artifactId version3.1.2/version /dependency然后在application.yml里配置连接信息和序列化器。一个长期运行平滑不报错的配置是这样spring: kafka: bootstrap-servers: localhost:9092 producer: key-serializer: org.apache.kafka.common.serialization.StringSerializer value-serializer: org.apache.kafka.common.serialization.StringSerializer acks: all retries: 3 batch-size: 16384 linger-ms: 5 compression-type: snappy逐项解释一下这些参数背后的含义。acksall表示消息写入 Leader 且所有 ISR 副本都确认后才返回成功这是最可靠但也是延迟相对最高的一种方式retries3表示发送失败后自动重试配合enable.idempotencetrue可以避免重试导致的消息重复linger-ms5和batch-size16384的意思是生产者攒够 16 KB 数据或者等够 5 毫秒才发一批这两个参数直接决定吞吐调大延迟上升调小吞吐下降。发送消息的代码很简单Service public class OrderProducer { Autowired private KafkaTemplateString, String kafkaTemplate; public void sendOrder(String orderId, String payload) { kafkaTemplate.send(order-topic, orderId, payload); } }这个例子里的第二个参数orderId是消息的 Key。Kafka 分区分配规则是“相同 Key 的消息进入同一分区”所以如果你希望某个订单的创建、支付、完成事件严格有序地消费就应该用订单 ID 做 Key。3.3 消费者与 Consumer Group 的核心机制消费者的代码同样分成配置和业务两部分。基础配置如下spring: kafka: consumer: group-id: order-group key-deserializer: org.apache.kafka.common.serialization.StringDeserializer value-deserializer: org.apache.kafka.common.serialization.StringDeserializer auto-offset-reset: earliest enable-auto-commit: falsegroup-id是消费者所属组的名字auto-offset-resetearliest表示当消费者没有历史 Offset 时从头消费enable-auto-commitfalse表示手动提交 Offset这是生产环境推荐的配置因为自动提交可能造成消息已经处理完但 Offset 没来得及提交重复消费就发生了。监听消费的代码如下Component public class OrderConsumer { KafkaListener(topics order-topic, groupId order-group) public void onOrder(String payload) { // 业务处理逻辑 System.out.println(收到消息: payload); // 处理完成后手动提交 Offset ack.acknowledge(); } }这里最核心的概念是消费者组与分区的对应关系。Kafka 的规则是一个分区只会被同一个消费组内的一个消费者实例消费而一个消费者实例可以消费多个分区。如果订单 Topic 有 6 个分区你启动 3 个消费者实例每个实例分到 2 个分区如果启动 8 个实例会有 2 个实例空转。所以消费者实例数不是越多越好最好与分区数匹配。另外还要理解消费者的“分区分组均衡”机制当一个实例加入或退出消费组时Kafka 会自动触发 Rebalance把分区重新分配。Rebalance 期间该消费者组的所有分区消费都会暂停这也是生产环境频繁重启消费者导致消费延迟变高的原因之一。3.4 可视化工具AKHQ 与 Kafka UI只看命令行日志排查问题效率太低了可视化工具几乎是必需品。现在社区里用得比较多的是 AKHQ 和 Kafka UI。AKHQ 是功能全面型的代表界面能直接查看 Topic 列表、分区数量和每个分区的 Produce 速率、Consume 速率、Offset 差距还能创建和删除 Topic甚至管理 Kafka Connect 的 connector 任务。如果你想查看某个 Connector 任务的状态进入 AKHQ 的 Connect 菜单就能看到每个 Connector 的 Tasks 列表以及state字段是 RUNNING 还是 FAILED点击任务还能看详细的错误日志。这个功能在排查 Flink 或 CDC 同步链路时特别实用。Kafka UI 是 Provectus 开源的项目界面比 AKHQ 更现代消息预览、消费者组管理、查看每个分区的最新 Offset 都做得比较顺手。选哪个取决于习惯我的建议是本地学习用 Kafka UI因为配置简单运行轻量生产环境排查 Kafka Connect 时用 AKHQ因为它对 Connect 的支持更深入。用 Docker 一分钟就能拉起一个 Kafka UIdocker run -d \ -p 8080:8080 \ -e KAFKA_CLUSTERS_0_NAMElocal \ -e KAFKA_CLUSTERS_0_BOOTSTRAPSERVERShost.docker.internal:9092 \ provectuslabs/kafka-ui:latest打开http://localhost:8080就能在界面上看到集群的信息。可视化工具最大的价值是让你快速发现消费积压——如果消费者的消费速率低于生产速率Offset 差值会持续拉大界面上每个 Topic 旁边的 Lag 数值会不断上涨。出现这种情况基本就是消费者处理太慢或者出现了 Rebalance下一步就该顺着我刚才讲的方向排查了。4. 核心机制深度解析与故障排查4.1 消息可靠性不丢不重到底怎么做到可靠性是消息队列落地时最绕不开的话题面试也几乎必问。先分清两个问题消息丢失和消息重复它们的成因和解法完全不同。消息丢失可能发生在三个环节。生产阶段丢失Ack 机制没设好生产者发出消息后没确认就离开了消息没到达 Broker。解决方法是acksall配合重试机制。Broker 阶段丢失Leader 各自写入后没来得及同步副本就宕机这条消息就丢了。解决方法是副本数至少 2 或 3并且min.insync.replicas设置为 2表示至少有 2 个副本写入成功才算成功。消费阶段丢失消费者拉到消息后自动提交了 Offset但业务逻辑还没来得及执行就宕机消息就丢了。解决方法是enable-auto-commitfalse处理完业务逻辑再手动提交 Offset。消息重复则恰恰相反它往往发生在“消息其实已经处理成功了但 Offset 没来得及提交”或者“生产者重试发送”的场景。比如消费者从数据库取出一条订单记录插入成功正准备提交 Offset 时进程崩溃重启后又从原来的位置消费那条订单就被插入两次。想彻底避免重复从消息队列本身是做不到的业界通用方案是“消费幂等”——即同一个操作执行多少次结果都相同。最实用的做法是建立一张消费记录表以消息的唯一 ID 作为主键消费前先查表存在就跳过不存在就执行业务并写入记录。或者利用数据库的唯一索引兜底插入重复时直接捕获冲突异常忽略即可。4.2 消息延迟高的排查思路“消息延迟高”是运维场景里最常收到的告警。所谓延迟指消息从生产到被消费的间隔时间远超正常水位通常用“消费滞后量”量化也就是每分区当前写入的最大 Offset 和消费者已提交 Offset 的差值。排查时我按下面这套顺序来第一看生产端是否积压。kafka-console-consumer.sh或者可视化工具里如果 Topic 的总消息量在短时间内暴涨说明生产速率远高于正常水平此时消费 Lag 自然变大。这种情况通常是业务上游异常可以先确认是否有人改了生产者配置或者存在定时任务集中执行。第二看消费端是否出现了 Rebalance。消费者实例频繁重启、心跳超时、处理时间过长导致自己掉出消费组都会触发 Rebalance而 Rebalance 期间消费是暂停的。日志中如果看到RebalanceInProgress和频繁的Assignment记录问题就在这里。解决办法是调大session.timeout.ms和heartbeat.interval.ms并检查消费者的实际处理耗时。第三看单个消息处理耗时。这个最好定位给消费逻辑加一个耗时统计如果单条消息处理普遍超过 500 毫秒那瓶颈就在业务代码。需要注意别忽略慢外部调用比如消费时同步查数据库或者调第三方接口一个慢接口就能拖垮整个分区。第四排查网络和磁盘问题。消费者拉取数据依赖网络延迟和带宽如果消费者与 Broker 不在同一机房内网带宽被占满时消费速率也会骤降。磁盘 IO 有问题时Broker 端读取已提交日志的速度下降一样会造成 Lag。用iostat看磁盘等待时间用网卡监控看带宽占用率这两项很容易被忽略。4.3 常见报错与处理实录这里整理几个我实际遇到过的报错都是新手很容易撞上的。第一个是org.apache.kafka.common.network.InvalidReceiveException: Invalid receive。这个报错通常发生在客户端版本和服务端版本不兼容时或者有人用不规范的客户端工具连接了你的 Broker。处理办法是先确认客户端 Kafka 版本是否与服务器端接近再看访问地址是否正确生产环境最容易犯的错是把外网地址发给了客户端导致客户端反复重连也连不通。第二个是Connection to node -1 could not be established. Broker may not be available。这基本是advertised.listeners配置错误。Kafka 有个容易踩的坑你连接时用的是localhost:9092但 Broker 注册到集群的地址可能是内网 IP客户端拿到内网 IP 后连不通。解决办法是保证advertised.listeners用的是客户端能够访问到的地址。第三个是消费者日志出现大量的Offset commit failed。最常见的原因是消费者组的 Offset 提交太频繁或者事务冲突。处理办法是调大max.poll.interval.ms并检查是否有消费者实例在 Rebalance 期间尝试提交 Offset。旧版本会有auto.commit.interval.ms过短导致频繁提交的问题这个参数直接改为5000到10000毫秒比较合理。第四个是Failed to update metadata after 60000 ms。这个报错通常与网络抖动或认证配置有关但最典型的原因是生产者在启动时无法连接到任何一个 Broker。检查 bootstrap-servers 配置的 IP 和端口是否被防火墙拦截然后用telnet命令直接测端口通不通能避免很多无谓的排查。4.4 Kafka、RabbitMQ、RocketMQ 选型对比选型问题不算故障但不选对的代价比故障更大。网上对这三款消息队列的对比文章很多我直接用一张表总结核心差异。维度KafkaRabbitMQRocketMQ吞吐量极高百万级消息/秒级别较低万级消息/秒级别高十万到百万级消息延迟毫秒级但吞吐越高延迟会上升微秒级延迟极低毫秒级消息可靠性高靠副本和 ACK 保证较高支持事务消息高支持事务消息消息顺序性分区内有序全局有序配置麻烦分区队列内有序功能丰富度核心功能简单需要扩展路由灵活、TTL、死信队列事务消息、延迟消息、重试机制完整运维复杂度中高集群需要打理低单机或简单集群中配套组件较多典型场景日志采集、大数据流处理、用户行为追踪低频业务通知、任务调度、内部系统解耦电商订单、交易流水、金融级可靠性场景选择的核心逻辑是场景驱动。如果业务量不大比如每天只有几万条消息用 RabbitMQ 最舒服轻量、可靠、社区资料多。如果是金融、订单这类业务消息不能丢又需要事务和延迟消息能力RocketMQ 是最合适的。而一旦涉及海量日志、行为埋点、流式计算Kafka 基本是唯一解。顺带说一句很多人问“Kafka 和 RabbitMQ 哪个好用”这种问题其实没有标准答案——一台测试机一个 Topic 的负载下RabbitMQ 的性能也很能打但同样规模放到百万级消息场景RabbitMQ 和 Kafka 的差别就体现出来了。5. 面试高频问题速查这部分整理几个我在面试中高频遇到的问题每道题给出精练的回答思路可以直接作为复习提纲。5.1 Kafka 为什么吞吐量那么高从四个机制回答顺序磁盘写、零拷贝、批量发送与压缩、分区并行。重点讲清楚零拷贝的流程磁盘数据经过 sendfile 直接由内核发送到网卡不需要经过用户态省掉了大量上下文切换和数据复制。追加回答时可以补充“顺序写使得磁盘磁头不需要来回移动”这是性能的物理基础。5.2 Kafka 会不会重复消费会。从两个角度解释生产者重试导致的重复写入以及消费者处理成功但 Offset 未提交导致的重复消费。解决方案是让消费端做幂等处理例如数据库唯一键、Redis 去重或者记录消费表。注意提问人真正想听的是你有没有实际处理过重复消费问题的经验。5.3 为什么需要分区分区带来三个价值扩展性、并行度、负载均衡。单个 Broker 的存储和吞吐是有限的多分区能横向扩展数据容量消费者组内多实例并行消费依赖分区分配分区的副本机制提供数据冗余。再补充一句分区增加了数据的重排成本这也是为什么需要谨慎设置分区数。5.4 消费者组内 Rebalance 是如何触发的触发条件有三类消费者实例加入或退出、Topic 元数据发生变化、消费者心跳超时被判定为死亡。Rebalance 期间消费暂停可能导致消息延迟上升。解决过度 Rebalance 的思路包括稳定消费者实例、合理处理耗时、设置合适的会话超时时间。5.5 如何保证消息顺序Kafka 只提供分区内的消息顺序性不提供跨分区全局顺序。保证顺序的办法是让需要有序的消息带相同 Key相同 Key 的消息一定进入同一分区。在消费端单分区由单消费者实例消费顺序自然保持。如果业务需要全局顺序那就只能设置成单分区但同时会牺牲吞吐。在这里我还想补充一个观点Kafka 的面试题翻来覆去就是可靠性、顺序性、消费组、性能这些点真正拉开差距的不是背得多熟而是能不能结合自己实际踩过的坑说清楚“当时为什么这么配”。我在面试别人的时候几乎不问概念直接问“你线上遇到过消息积压三千万条怎么办”能把这个问题讲清楚的人前面那些概念早就内化了。写在最后的一点体会接触 Kafka 这几年我最大的感受是入门难的不是 API而是理解“分区模型”和“副本机制”带来的底层思维转变。很多新手一上来就写生产者消费者跑通了就觉得自己会了等到生产环境抛出来重复消费、Rebalance 风暴才发现自己对 Kafka 的理解停留在表面。如果你按照这篇文章把环境搭了一遍建议再花点时间做一组实验把一个消费组从 1 个实例扩到 6 个观察分区分配变化然后把某个消费者直接 kill 掉观察 Rebalance 后的 Lag 波动。这些实验比单纯读文档有用得多也是我目前带新人时必定安排的任务。
RELATED READING

延伸阅读

更多一线实战笔记与深度复盘,助您持续精进