详解 Kafka 核心架构(核心流程):分区分布、消息路由与消费分配原理(三)

📅 2026/7/27 23:21:41
详解 Kafka 核心架构(核心流程):分区分布、消息路由与消费分配原理(三)
最近在总结Kafka的核心流程对于生产者如何选择分区、分区在集群中如何摆放、消费者组又如何分配分区等问题容易混淆。这篇文章把Kafka的Topic、分区、生产者路由、集群分布以及消费者组机制彻底讲清楚。一、Topic只是逻辑上的“频道”在Kafka中Topic主题是一个逻辑概念你可以把它理解成一个消息的分类标签或频道。生产者向某个Topic发消息消费者订阅某个Topic收消息但在磁盘上你找不到一个叫“Topic”的文件。真正存储数据、支撑并行和顺序性的是——Partition分区。每个Topic可以划分为多个Partition。Partition是物理上的有序、不可变的消息队列底层是一个个的日志分段LogSegment。一个Topic的多个分区可以分布在不同的Broker上从而实现水平扩展。因此当我们讨论消息的存储、顺序、负载时本质上都是在讨论分区。二、生产者如何把消息放进分区生产者发送消息时必须决定这条消息落到Topic的哪一个分区。Kafka提供了三种基本策略外加一种自定义扩展1. 指定分区手动指定在ProducerRecord中直接指明partition字段消息就被强制写入该分区。这适用于需要严格顺序控制的业务场景比如同一个用户的所有操作必须有序。ProducerRecordString, String record new ProducerRecord(my-topic, 0, key, value); // 指定分区02. 基于Key的哈希Key Hash如果指定了消息的Key但没有指定分区Kafka会默认用murmur2算法对Key取哈希然后模以分区总数得到目标分区。这样相同Key的消息总是会落到同一个分区从而保证该Key的消息顺序性。ProducerRecordString, String record new ProducerRecord(my-topic, user123, message); // 按key哈希3. 轮询Round-Robin既没指定分区也没指定Key时Kafka生产者会采用轮询的方式将消息依次放入各个可用分区。在早期的版本2.4中轮询是批量的每个batch轮换一次新版本默认改为黏性分区sticky即把一批消息发给同一个分区等该批次满了再切换到下一个分区以减少网络请求。但宏观上依然是均匀分布。4. 自定义分区器如果以上三种都不满足需求你还可以实现org.apache.kafka.clients.producer.Partitioner接口打造属于自己的分区逻辑。重要结论无论哪种策略单个分区内的消息天然是有序的写入顺序即偏移量顺序但整个Topic跨分区则是无序的。这就是Kafka“分区有序主题无序”的由来。三、分区在集群Broker上的分布艺术你可能会问一个Topic有3个分区Kafka集群有3个Broker这些分区到底怎么“摆放”这正好对应你最初描述里的那段话“先随机选一个放0接着按broker顺序往后放1、2循环一遍后可以重复broker接着放”这其实是在描述Kafka分区副本分配算法的简化版。我们来具体解释。分区副本的默认分配策略假设3个BrokerBroker-0, Broker-1, Broker-21个Topic3个分区副本因子1无冗余Kafka控制器在创建Topic时会随机选取一个起始Broker索引例如随机到Broker-1将Partition-0放在Broker-1上按Broker顺序向后轮询Partition-1放到Broker-2Partition-2放到Broker-0因为到了末尾又回头如果副本因子3每个分区一主二从那么每个分区的副本也会用类似的轮询方式依次放到不同的Broker上同时引入“机架感知”来尽量避免所有副本堆在同一个机架。算法伪逻辑如下将Broker列表随机打乱然后以轮询方式为每个分区选取主副本位置其余副本依次放到后续Broker上最后保证同一个分区的不同副本不会在同一个Broker上且尽量跨机架。这就是你所说的“先随机选一个放0接着往后放1、2循环复用Broker”的底层含义。这种分配方式能够让分区Leader均匀地分散到各个Broker避免热点。分区分布与消费顺序的关系即使分区散布在不同的机器上由于每个分区在物理上就是一个顺序写的日志同一个分区内的消息顺序是绝对保证的无论Leader在哪个Broker。消费者只要按照偏移量顺序拉取就能重现写入时的顺序。四、生产者如何找到分区的Leader —— 路由发现三重奏你在最后括号里提到的“生产者找分区有三种方法一通过代理二通过重定向三客户端先查询路由表根据路由表找到broker”其实是在描述生产者定位分区Leader Broker的过程。下面我把它转化为更准确的Kafka实践1. 客户端主动拉取路由表主流方式现代Kafka客户端Java、librdkafka等启动时会通过配置的bootstrap.servers连接任意一个Broker然后立即发送一个元数据请求Metadata Request获取集群中所有Topic的分区信息包括每个分区的主副本Leader在哪个Broker各个副本的分布情况每个Broker的IP和端口客户端将这些信息缓存在内存中形成一张“路由表”。之后生产者发送消息时直接从路由表中查出目标分区Leader所在的Broker建立TCP连接将消息发送过去。2. 重定向转发与异常处理如果生产者将消息发到了非Leader的Broker比如某分区Leader已经发生转移而客户端元数据还未刷新该Broker会返回一个NotLeaderForPartition异常并在响应中携带当前Leader的信息。客户端接收到这个异常后会立即更新自己的元数据缓存并重试发送到正确的Leader。这个过程就是“重定向”。3. 通过代理历史特殊场景严格来说现代Kafka架构并不推荐也不依赖于中心代理。早期0.8版本之前生产者默认连接到某个Broker该Broker可能作为代理将消息转发到真正的Leader但这会带来单点瓶颈。现在的“代理”更多指的是外部组件如Kafka REST Proxy、网关等它们作为中间层接收消息再转发。这不是原生产者的标准工作方式所以我们可以把它看作一种架构变体而非分区寻找的核心手段。用一张图总结现在的正常流程生产者 --(1)发送元数据请求-- Bootstrap Broker--(2)返回路由表--------生产者 --(3)直连Leader Broker 发送消息-- Leader所在Broker若返回NotLeader错误则回到(1)刷新路由重试。五、消费者组与分区的“分配契约”消费者这一侧同样绕不开分区。同一个Consumer Group内的消费者共同订阅一个或多个Topic它们之间会根据规则划分所负责的分区让每个分区在同一组内只被一个消费者消费。分区分配策略常见的策略有可在partition.assignment.strategy配置Range范围将分区连续段分配给消费者。比如3个分区p0,p1,p22个消费者可能会把p0,p1分给C1p2分给C2。容易导致分配不均。RoundRobin轮询把所有分区按顺序逐个轮询分配给消费者。上述例子中C1分到p0,p2C2分到p1相对均匀。Sticky粘性尽量保持现有分配不变再均衡时仅移动最少的分区减少开销。Cooperative Sticky协作式粘性也是粘性策略但再均衡是分步增量进行避免“Stop-the-world”式的全局暂停。再均衡Rebalance当消费者加入、离开或Topic分区数变化时消费者组会触发再均衡重新分配分区。这个过程中整个组会暂时停止消费直到新的分配方案达成。正是由于每个分区只能被同组的一个消费者消费且分区内部有序我们才能设计出严格保序的业务逻辑而多消费者并行消费不同分区又实现了高吞吐。六、一图总结核心数据流Topic (逻辑)/ | \Partition0 Partition1 Partition2 (物理, 分散在不同Broker)| | |[Leader Broker0] [Leader Broker1] [Leader Broker2]^ ^| |生产者 (轮询/key哈希/指定) 消费者组 (按策略分配分区)生产者通过元数据缓存定位各分区Leader直连发送。分区是顺序存储的最小单元同一分区保序整体通过多分区并行。消费者组内分区独占既保证顺序又实现水平扩展。七、写在最后Kafka以分区为核心的架构设计巧妙地在顺序性和并行度之间取得了平衡。理解Topic的逻辑抽象、生产者的三种分区选择方式、分区在Broker间的轮询分布算法以及消费者的分区分配与再均衡是驾驭Kafka的基础。