Skip to content

Kafka 分区分配策略:Range/Sticky/CooperativeSticky 对比

提出问题

Kafka 消费组在重平衡时,需要将 Topic 的多个分区分配给组内的各个消费者实例。分配策略直接决定了三个关键结果:负载是否均匀重平衡时多少分区需要移动消费组是否需要全量暂停。很多线上问题——热点消费、消费倾斜、频繁重平衡导致的抖动——追根溯源都跟分配策略选错了有关。

Kafka 2.x 时代默认的是 RangeAssignor,但它在多 Topic 订阅场景下存在严重的不均匀问题。Kafka 3.0 之后默认切换到了 CooperativeStickyAssignor,但仍有大量存量集群还在用旧策略。面试官问这个问题,本质是考察你是否理解分配策略的演进背景生产环境选型决策

分析问题

RangeAssignor:按 Topic 独立分配,多 Topic 时严重倾斜

RangeAssignor 的策略是按 Topic 逐个分配:对每个 Topic,将它的分区按字典序排序,然后平均分配给消费者。

假设有 2 个消费者(C1、C2)和 2 个 Topic(T1、T2,各 3 个分区):

T1 分区: [0, 1, 2] → C1 拿 T1-0, T1-1,C2 拿 T1-2
T2 分区: [0, 1, 2] → C1 拿 T2-0, T2-1,C2 拿 T2-2

结果:C1 承担 4 个分区,C2 承担 2 个分区

问题:每个 Topic 独立计算后,消费者数量不能整除分区数时,编号靠前的消费者总是多拿一个分区。当 Topic 数量增多时,这种倾斜会被放大:

java
// 假设 20 个 Topic,每个 10 分区,5 个消费者
// 每个 Topic 分配:C1~C2 各 2 分区,C3~C5 各 2 分区
// 累积:C1~C2 各 40 分区,C3~C5 各 20 分区
// 差了一倍!

核心源码逻辑(Kafka 2.x RangeAssignor):

java
public Map<String, List<TopicPartition>> assign(Map<String, Integer> partitionsPerTopic,
                                                  Map<String, Subscription> subscriptions) {
    Map<String, List<String>> consumersPerTopic = consumersPerTopic(subscriptions);
    Map<String, List<TopicPartition>> assignment = new HashMap<>();
    
    for (Map.Entry<String, List<String>> topicEntry : consumersPerTopic.entrySet()) {
        String topic = topicEntry.getKey();
        List<String> consumersForTopic = topicEntry.getValue();
        int numPartitions = partitionsPerTopic.get(topic);
        int numConsumers = consumersForTopic.size();
        
        // 每个 Topic 独立计算,不均匀
        int partitionsPerConsumer = numPartitions / numConsumers;
        int consumersWithExtraPartition = numPartitions % numConsumers;
        
        // ...
    }
    return assignment;
}

线上踩坑实例:我经手过的一个业务线,订阅了 12 个 Topic,每个 8 分区,部署了 6 个消费者实例。Range 模式下,C1 消费者承担了 24 个分区,C6 只承担了 12 个分区。C1 的机器 CPU 长期在 75% 运行,C6 只有 30% 左右,且 C1 的消息处理延迟最高达到 3 秒,C6 不到 1 秒。排查时 kafka-consumer-groups --describe --group <group> 一看,分区分配明显倾斜。切换到 CooperativeSticky 后,每个消费者稳定在 14~18 个分区,CPU 分布均匀在 40%~55% 之间。

StickyAssignor:更均匀 + 最小移动

StickyAssignor 在 Kafka 2.3 引入,改进点有两个:

  1. 跨 Topic 统一分配:不是按 Topic 独立算,而是对所有消费者订阅的所有分区做全局优化,保证各消费者拿到的分区数量尽可能接近。
  2. 最小化移动:重平衡时,尽可能保留已有的分区分配,只移动必要的分区来达到均匀。

回到上面的例子,StickyAssignor 会尽量让 C1 和 C2 各拿 3 个分区(T1 的 1.5 个 + T2 的 1.5 个),而不是 C1 拿 4 个、C2 拿 2 个。

但 StickyAssignor 仍然是 Eager 模式——重平衡时所有消费者停止消费(Stop The World),释放所有分区,再重新分配。对于高吞吐场景,每次重平衡会造成几秒到几十秒的消费中断。

一个典型的重平衡风暴案例:某日志采集系统,Consumer Group 有 30 个实例,消费 15 个 Topic 共 120 个分区。使用 StickyAssignor。某次凌晨发布滚动重启,每个实例重启时触发一次 Eager 重平衡,30 个实例依次重启,导致 30 次全量重平衡,每次持续 8~15 秒。整个发布窗口持续 8 分钟,其中约 4 分钟消费者处于不可用状态。日志堆积从 0 飙升到 300 万条。换成 CooperativeSticky 后,同场景下每次重平衡只影响 2~3 个分区的消费者,单次停顿不超过 200ms,堆积分分钟消化掉。

CooperativeStickyAssignor:增量合作式,告别全量暂停

CooperativeStickyAssignor 在 Kafka 3.0 引入,现在是默认策略。它的核心改进是 增量合作式重平衡

  1. 重平衡时,Coordinator 不要求所有消费者释放全部分区,而是只通知需要调整的消费者
  2. 消费者逐步释放(Revoke)少部分分区,其他消费者继续消费。
  3. 多次通信协商后,最终达到均匀分配。

对比三种策略的重平衡过程

mermaid
sequenceDiagram
    participant C1 as Consumer 1
    participant C2 as Consumer 2
    participant C3 as Consumer 3 (新加入)
    
    Note over C1,C3: Range/Sticky (Eager 模式)
    Coordinator->>C1: Revoke 全部
    Coordinator->>C2: Revoke 全部
    C1->>Coordinator: 已释放
    C2->>Coordinator: 已释放
    Coordinator->>C1: Assign 新分区
    Coordinator->>C2: Assign 新分区
    Coordinator->>C3: Assign 新分区
    
    Note over C1,C3: CooperativeSticky (增量模式)
    Coordinator->>C1: Revoke 分区 A
    Coordinator->>C2: Revoke 分区 B
    C1->>Coordinator: 已释放
    C2->>Coordinator: 已释放
    Coordinator->>C3: Assign 分区 A, B

配置方式

properties
# Kafka 3.0+ 默认已是 CooperativeStickyAssignor
# 显式配置
partition.assignment.strategy=org.apache.kafka.clients.consumer.CooperativeStickyAssignor

# Spring Boot 配置
spring.kafka.consumer.properties.partition.assignment.strategy=org.apache.kafka.clients.consumer.CooperativeStickyAssignor

为什么不停全量也能完成分配? 核心在于新版的 JoinGroup 协议引入了 member.epoch 字段。每次 JoinGroup 响应中,Coordinator 会递增 epoch,消费者根据 epoch 变化判断是否增量还是全量重平衡。CooperativeSticky 在第一次 JoinGroup 时,只 Revoke 那些需要迁移的分区,保留已有的分区继续消费,并带上 CurrentAssignment 元数据。Coordinator 在全局视图中,发现所有消费者都提交了 CurrentAssignment 后,才下发最终的 Assign。这个过程分 2~3 轮 JoinGroup 完成,每轮只有几百毫秒。

自定义分配策略:什么时候需要?

当集群中消费者实例的硬件配置不一致时(比如 8C16G 的机器和 4C8G 的机器混部),默认策略的均匀分布反而有问题。需要自定义分配策略,让高性能实例多拿分区。

java
public class WeightedAssignor extends AbstractPartitionAssignor {
    
    @Override
    public Map<String, List<TopicPartition>> assign(
            Map<String, Integer> partitionsPerTopic,
            Map<String, Subscription> subscriptions) {
        
        // 从订阅信息中提取权重(通过 userData 传入)
        Map<String, Integer> weights = new HashMap<>();
        for (Map.Entry<String, Subscription> entry : subscriptions.entrySet()) {
            String consumerId = entry.getKey();
            ByteBuffer userData = entry.getValue().userData();
            // 假设 userData 前 4 字节是 int 权重
            weights.put(consumerId, userData.getInt());
        }
        
        int totalWeight = weights.values().stream().mapToInt(Integer::intValue).sum();
        int totalPartitions = partitionsPerTopic.values().stream().mapToInt(Integer::intValue).sum();
        
        // 按权重分配
        Map<String, List<TopicPartition>> assignment = new HashMap<>();
        int assigned = 0;
        for (Map.Entry<String, Integer> w : weights.entrySet()) {
            int expected = (int) Math.round((double) w.getValue() / totalWeight * totalPartitions);
            // 分配 expected 个分区给该消费者
            // ...
        }
        return assignment;
    }
    
    @Override
    public String name() {
        return "weighted";
    }
}

消费者端传入权重:

java
Properties props = new Properties();
props.put(ConsumerConfig.PARTITION_ASSIGNMENT_STRATEGY_CONFIG, 
          WeightedAssignor.class.getName());
// 通过 userData 传入权重
props.put(ConsumerConfig.INTERNAL_LEAVE_GROUP_ON_CLOSE_CONFIG, false);

不过实际生产中,我更推荐的做法是均匀分配 + 在消费者侧用自适应限流来控制消费速率,而不是在分配策略中做差异化。因为权重策略一旦某个实例挂了,需要 Coordinator 重新计算权重分布,复杂度高且容易出错。除非你运维的集群规模在 100+ 实例且机器配置差异明显,否则不值得自定义。

升级与迁移陷阱

从 Range 切换到 CooperativeSticky 的兼容性问题

这是面试中高频踩坑点。Kafka 2.x 集群直接升级到 3.x 后,如果旧消费者还在用 RangeAssignor,不会自动切换。要改配置。

但更坑的是混合版本消费者:同一个 Group 内,部分消费者用 RangeAssignor、部分用 CooperativeStickyAssignor,会导致 Coordinator 抛出 InconsistentGroupProtocolException,消费者不断触发重平衡,陷入死循环。

正确的迁移步骤

  1. 先停掉所有消费者,统一升级到 Kafka 3.x 客户端
  2. 在一个 Group 中一次性切换所有消费者实例的 partition.assignment.strategy
  3. 先灰度一个 Group,观察 15 分钟,确认无异常重平衡后再逐步灰度其他 Group

旧版 Coordinator 的 UnsupportedVersionException

如果使用了 CooperativeStickyAssignor,但 Group Coordinator 所在的 Broker 版本低于 2.3,JoinGroup 请求会返回 UnsupportedVersionException,消费者会降级为 Eager 模式。表现为配置了增量策略但实际效果还是全量暂停,非常隐蔽。

排查方法:查看消费者日志中的 JoinGroup response 行,看是否有 UnsupportedVersionException 的降级日志。

总结

三选一:生产环境直接用 CooperativeStickyAssignor

策略负载均匀重平衡停顿版本适用场景
RangeAssignor❌ 多 Topic 时严重倾斜全量暂停(Eager)0.8+仅单 Topic 订阅
StickyAssignor✅ 均匀全量暂停(Eager)2.3+低重平衡频率场景
CooperativeStickyAssignor✅ 均匀增量暂停(Cooperative)3.0+所有场景,默认推荐

关键要点

  • Range 倾斜在消费者订阅多 Topic 时是致命问题,小号消费者可能比大号消费者多扛一倍分区,导致热点消费。
  • CooperativeSticky 不仅解决了负载均匀,还解决了重平衡时的全量暂停问题,建议所有 ≥ 3.0 的集群直接使用。
  • 自定义分配策略只在大规模异构集群(不同 Consumer 实例性能差异大)时才需要,普通场景用默认的即可。
  • 升级注意:存量集群从 Range 切换到 CooperativeSticky 时,建议先灰度验证,因为增量模式的协商流程和旧版 Eager 模式不同,可能触发短暂的 UnsupportedVersionException混合版本消费者共存会触发重平衡死循环,务必一次性全部切换
  • 面试追问:如果面试官问"为什么 CooperativeSticky 不是一次完成的",回答:因为增量分配需要多轮 JoinGroup 协商,每轮 Coordinator 只分配部分分区,消费者需要确认接收后才能继续下一轮。这是为了确保所有消费者对新分配达成共识,类似分布式一致性协议的"两阶段提交"思想。
  • 排查命令kafka-consumer-groups --bootstrap-server <broker> --group <group> --describe 查看分区分配;kafka-consumer-groups --bootstrap-server <broker> --group <group> --describe --members 查看每个成员的分区详情。

参考:Apache Kafka 官方文档 — Consumer Rebalance Protocol;KIP-429: Incremental Rebalance Protocol;KIP-54: Sticky Partition Assignment

手撕 → 框架 → 生产化,一步步把 AI Agent 工程化搞透。