乐途乐途
主页
  • 计算机基础

    • TCP/IP
    • Linux
    • HTTP
  • 数据库

    • SQL
    • MySQL 5.7
  • 编程语言

    • C
    • C++
    • Java SE
    • Python2
    • Python3
  • 数据格式

    • JSON
    • XML
  • 认证与安全

    • JWT
  • 工具

    • Markdown
  • Git

    • GitFlow
  • Quartz

    • Quartz
  • Java

    • Maven 入门
    • Maven 进阶
    • MyBatis
    • Spring
    • Spring MVC
  • Java

    • Spring Boot
    • Spring Cloud
    • Spring Cloud Alibaba
    • Spring Security
    • Spring AI
    • Spring Batch
    • Kafka
    • Java 设计模式
  • 缓存

    • Redis
  • 搜索引擎

    • Elasticsearch
  • 分布式协调

    • ZooKeeper
联系
阿里云
主页
  • 计算机基础

    • TCP/IP
    • Linux
    • HTTP
  • 数据库

    • SQL
    • MySQL 5.7
  • 编程语言

    • C
    • C++
    • Java SE
    • Python2
    • Python3
  • 数据格式

    • JSON
    • XML
  • 认证与安全

    • JWT
  • 工具

    • Markdown
  • Git

    • GitFlow
  • Quartz

    • Quartz
  • Java

    • Maven 入门
    • Maven 进阶
    • MyBatis
    • Spring
    • Spring MVC
  • Java

    • Spring Boot
    • Spring Cloud
    • Spring Cloud Alibaba
    • Spring Security
    • Spring AI
    • Spring Batch
    • Kafka
    • Java 设计模式
  • 缓存

    • Redis
  • 搜索引擎

    • Elasticsearch
  • 分布式协调

    • ZooKeeper
联系
阿里云
  • 学习路径
  • 第1章 消息队列与 Kafka 概述

    • 章节导读
    • 消息队列基础
    • Kafka 概述
    • Kafka 为什么快
  • 第2章 快速上手:单机环境搭建

    • 章节导读
    • 环境准备与安装
    • 快速启动
    • Topic 管理
  • 第3章 核心概念:主题、分区与日志

    • 章节导读
    • Topic
    • Partition
    • Offset
    • Segment 与存储结构
  • 第4章 生产者详解

    • 章节导读
    • Producer 概述
    • Producer 发送机制
    • Producer 分区策略
    • Producer 幂等与事务
    • Producer 配置
  • 第5章 消费者与消费组

    • 章节导读
    • Consumer 概述
    • Consumer Group 消费组
    • Consumer 分区分配策略
    • Consumer Offset 提交
    • Consumer 多线程
    • Consumer 配置
  • 第6章 Broker 与控制器

    • 章节导读
    • Broker 概述
    • Broker 配置
    • Controller
    • KRaft 共识协议
  • 第7章 副本与数据可靠性

    • 章节导读
    • 副本机制
    • ISR 与副本同步
    • Leader 选举
    • ACK 与一致性保证
    • 高水位与 Leader Epoch
  • 第8章 存储与性能优化

    • 章节导读
    • 存储架构
    • 日志清理与压缩
    • Page Cache 与零拷贝
    • Producer 性能优化
    • Consumer 性能优化
    • Broker 性能优化
  • 第9章 生产环境运维与监控

    • 章节导读
    • Topic 管理
    • Kafka 运维工具
    • 监控
    • 常见故障排查
  • 第10章 Kafka生态与面试考点

    • 章节导读
    • Kafka 生态全景
    • 面试高频 30 题

Consumer 分区分配策略

定义与作用

分区分配策略决定了消费组内各 Consumer 如何瓜分 Partition。不同的分配策略影响负载均衡度和Rebalance 效率。Kafka 提供了四种内置策略,从经典的 Range/RoundRobin 到新一代的 Cooperative Sticky。

核心原理

四种策略对比

Range 策略(默认)详解

Range 策略按 Topic 独立分配,每个 Topic 的分区排序后均分给消费者。问题:如果一个消费组订阅了多个分区数不均匀的 Topic,数据倾斜会很严重。

Sticky 策略

Sticky 试图保持旧分配不变,仅移动必要的分区:

Rebalance 前: C1→[P0,P1,P3],  C2→[P2]
C3 加入后(Range): C1→[P0,P1], C2→[P2,P3], C3→[]  (大量迁移)
C3 加入后(Sticky): C1→[P0,P1,P3], C2→[P2], C3→[]   (暂无迁移,等下次)

Cooperative Sticky(增量重平衡)

传统 Rebalance 是"Stop-The-World"的:全部撤销再重新分配。Cooperative Sticky 允许分批执行:

特性传统(Eager)Cooperative Sticky
撤销范围全部撤销只撤销需迁移的分区
消费暂停全组暂停仅被撤销分区的 Consumer 短暂暂停
适用版本所有Kafka 2.4+
配置partition.assignment.strategy=<Range>partition.assignment.strategy=org.apache.kafka.clients.consumer.CooperativeStickyAssignor

完整示例

示例一:Range 策略导致的数据倾斜

场景:消费组订阅 Topic A(3 分区)和 Topic B(1 分区),2 个 Consumer。

操作前环境:两个 Topic 已创建,Consumer 使用默认 Range 策略。

# 创建 Topic
bin/kafka-topics.sh --create --topic topic-a --partitions 3 --bootstrap-server localhost:9092
bin/kafka-topics.sh --create --topic topic-b --partitions 1 --bootstrap-server localhost:9092

分配结果(Range):

ConsumerTopic ATopic B总分区数
Consumer 1P0, P1P03
Consumer 2P2—1

问题:Consumer 1 负载是 Consumer 2 的 3 倍。

改用 RoundRobin:

props.put("partition.assignment.strategy",
    "org.apache.kafka.clients.consumer.RoundRobinAssignor");
Consumer总分区(RoundRobin)
Consumer 1topic-a-P0, topic-a-P2, topic-b-P0 → 3 分区
Consumer 2topic-a-P1 → 1 分区

(注意:RoundRobin 在这个场景改善有限,因为总数 4 分区 / 2 消费者是均匀的,但 Range 的分配方式不均)

示例二:Cooperative Sticky 减少 Rebalance 影响

场景:秒杀活动中 Consumer 3 因 OOM 重启,需要最小化 Rebalance 对消费的影响。

Properties props = new Properties();
// ... 基础配置 ...

// 关键配置
props.put("partition.assignment.strategy",
    "org.apache.kafka.clients.consumer.CooperativeStickyAssignor");
props.put("group.instance.id", "consumer-3");  // 静态成员,重启后不触发 Rebalance
props.put("session.timeout.ms", "30000");      // 合理设置心跳超时

操作前后对比:

事件Eager Rebalance 影响Cooperative Sticky 影响
C3 OOM 重启全组暂停消费 10-30s仅 C3 的分区暂停,C1/C2 继续消费
C3 恢复再次全组暂停 10-30sC3 重新加入,接收回原分区

易错场景

易错 1:多 Topic 订阅时默认 Range 策略的倾斜陷阱

场景:一个消费组订阅了 10 个 Topic,其中 8 个是 1 分区的低流量 Topic,2 个是 100 分区的高流量 Topic。

后果:Range 策略在每个 Topic 内独立分配,低流量和高流量混合分配后,数据倾斜非常严重。部分 Consumer 闲置,部分 Consumer 过载。

最佳实践:订阅多 Topic 时使用 RoundRobin 或 Sticky。

易错 2:忘记 group.instance.id 导致静态成员失效

场景:设置了 group.instance.id 但没有设置足够大的 session.timeout.ms。

后果:Consumer 重启时间超过 session.timeout.ms 时,静态成员仍然会触发 Rebalance。

正确做法:

props.put("group.instance.id", "consumer-1");
props.put("session.timeout.ms", "60000");  // 给重启留足时间

面试高频考点

Q:什么时候用 Range?什么时候用 Cooperative Sticky?

A:

Range(默认)适合:

  • 订阅单个 Topic 且分区数能被消费者数整除的场景
  • 简单场景,不想引入额外复杂性

Cooperative Sticky(推荐)适合:

  • 大规模消费组(10+ Consumer),Rebalance 代价大
  • Consumer 频繁扩缩容的场景(K8s 弹性伸缩)
  • 对消费延迟敏感的系统

RoundRobin 适合:

  • 订阅多 Topic 且希望负载均匀
  • 不介意 Rebalance 时全量重分配

Sticky(非 Cooperative) 是过渡方案,新项目直接用 Cooperative Sticky。

上一页
Consumer Group 消费组
下一页
Consumer Offset 提交