一、先搞懂什么是消费者组和分区分配
咱们先把两个最基础的概念掰明白,不然后面说分配策略都是白搭。首先得说消息队列(比如大家常用的Kafka),它会把一个主题(就是一类消息的集合,比如外卖订单主题)拆成好几个分区——你可以把分区理解成外卖站点的快递柜格子,每个格子只放一种类型的外卖订单,比如格子1放北京的订单,格子2放上海的,这样分是为了能同时处理更多消息,不然一个大柜子挤所有订单,处理速度肯定慢。
然后是消费者组,就是一群干活的人凑成的团队,专门处理某个主题的消息。比如一个订单主题有3个分区,消费者组里有2个消费者,那这俩消费者得商量着来,谁负责哪几个分区的消息,不能两个人抢同一个分区,也不能漏了某个分区没人管——这个“商量着来”的规则,就是咱们要讲的分配策略。
举个最直观的例子,比如你开了个奶茶店,每天要处理1000杯奶茶的订单,为了快,你把订单分成5个堆(分区):堆1放加珍珠的,堆2放加椰果的,堆3放加芋圆的,堆4放无糖的,堆5放热饮的。然后你找了3个店员(消费者)帮忙做奶茶,这3个店员怎么分这5堆订单?不同的分配策略,分法完全不一样。
二、第一种策略:Range分配(按范围分)
2.1 什么是Range分配
Range分配是最“朴素”的分配规则,简单说就是“按数字范围平均分”。比如有N个分区,M个消费者,先把分区按编号排好(比如刚才的5个分区编号是0、1、2、3、4),然后尽量让每个消费者拿到的分区数量差不超过1。
具体计算方法是:先算每个消费者最少能拿到多少个分区,公式是总分区数 ÷ 消费者数,余数就是多出来的分区数,把这些多的分区按顺序分给前几个消费者。
还是用奶茶店的例子:5个分区(N=5),3个消费者(M=3)。5÷3=1,余数是2。所以前2个消费者每人拿2个分区,最后1个拿1个。具体分法就是:
- 消费者0:拿编号0、1的分区(珍珠、椰果订单)
- 消费者1:拿编号2、3的分区(芋圆、无糖订单)
- 消费者2:拿编号4的分区(热饮订单)
2.2 Range分配的代码示例(技术栈:Java,基于Kafka客户端)
咱们用真实的Kafka客户端代码来写一个Range分配的配置,你可以直接复制到项目里用:
import org.apache.kafka.clients.consumer.ConsumerConfig;
import org.apache.kafka.clients.consumer.KafkaConsumer;
import java.util.Properties;
public class RangeConsumerExample {
public static void main(String[] args) {
Properties props = new Properties();
// 配置Kafka服务器地址
props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
// 配置消费者组ID,所有同组消费者的这个值必须一样
props.put(ConsumerConfig.GROUP_ID_CONFIG, "milk-tea-order-group");
// 配置序列化方式(Kafka要求消息必须是字节数组,这里用字符串转字节)
props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.StringDeserializer");
props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.StringDeserializer");
// 核心配置:指定Range分配策略
props.put(ConsumerConfig.PARTITION_ASSIGNMENT_STRATEGY_CONFIG, "org.apache.kafka.clients.consumer.RangeAssignor");
// 创建消费者实例
KafkaConsumer<String, String> consumer = new KafkaConsumer<>(props);
// 订阅要处理的主题
consumer.subscribe(java.util.Collections.singletonList("milk-tea-order-topic"));
// 持续拉取消息(这里简化写,实际项目要加异常处理)
while (true) {
consumer.poll(java.time.Duration.ofMillis(100)).forEach(record -> {
System.out.println("消费者拿到消息:" + record.value() + ",来自分区:" + record.partition());
});
}
}
}
2.3 Range分配的优缺点和适用场景
先讲优点:规则简单,实现成本低,Kafka默认的分配策略就是它,不需要额外配置,新手用起来最省心。
再讲缺点:有个很明显的坑——如果主题的分区数是固定的,消费者数变化时,容易出现分区分配不均的情况。比如刚才的例子,5个分区3个消费者,前两个消费者各拿2个,第三个拿1个,要是第三个消费者的分区消息特别多(比如热饮订单每天有300杯,其他每个分区只有100杯),那第三个消费者就会忙死,前两个闲得慌。还有一种情况,如果你有多个主题,每个主题的分区数不一样,Range分配会给每个主题单独算范围,容易导致某个消费者拿到的分区总数特别多,其他的特别少。
所以适用场景很明确:适合对分区分配均衡性要求不高的场景,比如测试环境、数据量小的生产环境,或者主题的分区数和消费者数刚好能整除的情况(比如6个分区,3个消费者,每个刚好拿2个,完全均衡)。
三、第二种策略:RoundRobin分配(轮询分)
3.1 什么是RoundRobin分配
RoundRobin的规则是“轮流拿”,比Range更公平一点。具体分两种情况: 第一种是“全局轮询”:把所有主题的所有分区按编号排好,然后按顺序轮流分给每个消费者。比如有两个主题,主题A有3个分区(0、1、2),主题B有2个分区(0、1),总共5个分区,3个消费者,那分法是:
- 消费者0:拿A0、B1
- 消费者1:拿A1、A2
- 消费者2:拿B0
第二种是“主题内轮询”:每个主题单独按顺序轮询分配。还是刚才的例子,主题A的3个分区轮询分给3个消费者,每个拿1个;主题B的2个分区轮询分给前2个消费者,最终分法是:
- 消费者0:A0、B0
- 消费者1:A1、B1
- 消费者2:A2
咱们用奶茶店的例子再讲一遍:5个分区(0、1、2、3、4),3个消费者。轮询的话就是按顺序分,第一个消费者拿0,第二个拿1,第三个拿2,然后再从第一个开始拿3,第二个拿4,最终分法是:
- 消费者0:0、3(珍珠、无糖订单)
- 消费者1:1、4(椰果、热饮订单)
- 消费者2:2(芋圆订单)
你看,这样分的话,每个消费者拿到的分区数还是差1,但分区的编号是分散的,不像Range那样前两个拿连续的分区,第三个拿最后一个。
3.2 RoundRobin分配的代码示例
还是用Java的Kafka客户端,只需要改分配策略的配置就行:
import org.apache.kafka.clients.consumer.ConsumerConfig;
import org.apache.kafka.clients.consumer.KafkaConsumer;
import java.util.Properties;
public class RoundRobinConsumerExample {
public static void main(String[] args) {
Properties props = new Properties();
props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
props.put(ConsumerConfig.GROUP_ID_CONFIG, "milk-tea-order-group");
props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.StringDeserializer");
props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.StringDeserializer");
// 核心配置:指定RoundRobin分配策略
props.put(ConsumerConfig.PARTITION_ASSIGNMENT_STRATEGY_CONFIG, "org.apache.kafka.clients.consumer.RoundRobinAssignor");
KafkaConsumer<String, String> consumer = new KafkaConsumer<>(props);
consumer.subscribe(java.util.Collections.singletonList("milk-tea-order-topic"));
while (true) {
consumer.poll(java.time.Duration.ofMillis(100)).forEach(record -> {
System.out.println("消费者拿到消息:" + record.value() + ",来自分区:" + record.partition());
});
}
}
}
3.3 RoundRobin分配的优缺点和适用场景
优点:比Range分配更公平,尤其是多个主题的情况下,能尽量让每个消费者拿到的分区数更均衡,不会出现某个消费者拿一堆连续分区的情况。
缺点:也有个坑——如果某个主题的分区数比消费者数少,那会有消费者拿不到这个主题的分区;另外,它的分配是完全随机按顺序的,不会考虑每个分区的消息量,比如某个分区消息特别多,刚好分给了一个能力弱的消费者,那这个消费者还是会忙死。
适用场景:适合对分区分配均衡性要求中等的场景,比如生产环境中主题数量多、分区数变化不频繁的情况,比如电商平台的商品主题、用户主题,这类主题的分区消息量比较平均,轮询分配不会有太大问题。
四、第三种策略:Sticky分配(粘滞分)
4.1 什么是Sticky分配
Sticky分配是目前最智能的分配策略,解决了Range和RoundRobin的很多坑。它的核心规则有两个: 第一个是“尽量保持分配不变”:比如消费者A之前拿了分区0和1,下次重新分配的时候,只要消费者A还在组里,就尽量让它继续拿这两个分区,不要随便换。 第二个是“尽量让分配更均衡”:如果消费者组里的消费者数量变了(比如加了一个新的消费者,或者某个消费者挂了),Sticky会重新调整分配,尽量让每个消费者拿到的分区数差不超过1,同时尽量少动原来的分配。
咱们用奶茶店的例子讲:第一次分配,5个分区3个消费者,Sticky会尽量让每个消费者拿的分区数差不超过1,比如分法是:
- 消费者0:0、1(珍珠、椰果)
- 消费者1:2、3(芋圆、无糖)
- 消费者2:4(热饮)
过了一会儿,消费者3加入了组,现在有4个消费者,Sticky会重新分配,尽量少动原来的分配,同时让每个消费者拿的分区数差不超过1。5个分区4个消费者,每个最少拿1个,余数1,所以有一个消费者拿2个,其他拿1个。Sticky会尽量不把原来的分区换给别人,比如调整后:
- 消费者0:0、1(保持原来的)
- 消费者1:2(原来拿2、3,现在把3给消费者3)
- 消费者2:4(保持原来的)
- 消费者3:3(拿原来消费者1的3)
这样的调整只动了一个分区,其他都保持不变,比Range和RoundRobin重新全部分配要稳定很多。
4.2 Sticky分配的代码示例
还是用Java的Kafka客户端,改分配策略的配置:
import org.apache.kafka.clients.consumer.ConsumerConfig;
import org.apache.kafka.clients.consumer.KafkaConsumer;
import java.util.Properties;
public class StickyConsumerExample {
public static void main(String[] args) {
Properties props = new Properties();
props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
props.put(ConsumerConfig.GROUP_ID_CONFIG, "milk-tea-order-group");
props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.StringDeserializer");
props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.StringDeserializer");
// 核心配置:指定Sticky分配策略
props.put(ConsumerConfig.PARTITION_ASSIGNMENT_STRATEGY_CONFIG, "org.apache.kafka.clients.consumer.StickyAssignor");
KafkaConsumer<String, String> consumer = new KafkaConsumer<>(props);
consumer.subscribe(java.util.Collections.singletonList("milk-tea-order-topic"));
while (true) {
consumer.poll(java.time.Duration.ofMillis(100)).forEach(record -> {
System.out.println("消费者拿到消息:" + record.value() + ",来自分区:" + record.partition());
});
}
}
}
4.3 Sticky分配的优缺点和适用场景
优点:解决了Range和RoundRobin的两个大问题——一是分配更均衡,不会出现某个消费者拿一堆分区的情况;二是粘滞性,尽量保持分配不变,减少了重新分配带来的开销(比如重新分配时,消费者要重新拉取消息,可能会重复消费或者漏消费)。
缺点:实现逻辑比前两个复杂,对Kafka客户端的版本有要求(Kafka 0.11.0.0及以上版本才支持),另外如果消费者组变化频繁,它的调整逻辑会稍微复杂一点,但总体还是比前两个好。
适用场景:适合对分区分配均衡性和稳定性要求高的生产环境,比如实时数据处理、订单处理、日志收集这类核心业务场景,比如外卖平台的订单主题、支付主题,这类场景对分配的稳定性和均衡性要求高,Sticky是最优选择。
五、三种策略的对比和选型总结
咱们把三种策略放在一起对比,你就能更清楚怎么选了:
- 如果你是新手,在测试环境玩,或者主题的分区数和消费者数刚好能整除,选Range,不用想,配置简单。
- 如果你是生产环境,主题数量多、分区数变化不频繁,对均衡性要求中等,选RoundRobin,比Range公平。
- 如果你是生产环境的核心业务,对均衡性和稳定性要求高,选Sticky,虽然配置稍微多一点,但最省心。
最后再给你一个选型的小口诀:测试用Range,普通生产用RoundRobin,核心生产用Sticky。
Comments