一、先搞懂什么是消费者组和分区分配

咱们先把两个最基础的概念掰明白,不然后面说分配策略都是白搭。首先得说消息队列(比如大家常用的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。