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

    • 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 题

Partition

定义与作用

Partition(分区)是 Kafka 的并行度单位和顺序保证边界。每个 Topic 被划分为 1 个或多个 Partition,每个 Partition 是一个有序的、不可变的日志(ordered, immutable log)。消息在 Partition 内被追加到尾部,以 Offset 标识位置。

Partition 解决了分布式系统中的两个核心矛盾:

  1. 顺序 vs 并行:单 Partition 保证消息严格有序,多 Partition 允许并行读写
  2. 单机容量 vs 数据量:不同 Partition 可分布在不同 Broker 上,突破单机磁盘 / 吞吐上限

核心原理

Partition 的分布模型

图中绿色标记的为 Leader 副本——所有读写都发生在 Leader 上。Follower 通过拉取 Leader 的数据保持同步。

消息路由:如何决定消息去哪个分区

旧版本(< 2.4)默认使用 Round-Robin,每次发送都切换分区,导致批处理效率低。Kafka 2.4+ 引入 Sticky Partitioner(粘性分区),在无 Key 时将同一批次的消息发往同一分区,批次满后再切换,显著提升了小消息场景的吞吐。

分区内的顺序保证

关键约束:Kafka 只保证单个 Partition 内的消息严格有序(按写入顺序),不保证跨分区的全局顺序。

如果要保证全局有序:将 Topic 的分区数设为 1。代价是失去并行能力。

完整示例

示例一:多分区并行写入验证

场景:创建一个 3 分区的 Topic,发送有 Key 的消息,验证相同 Key 的消息进入同一分区。

操作前环境:Kafka 单节点,Topic partition-demo 不存在。

步骤:

# 1. 创建 3 分区 Topic
bin/kafka-topics.sh --create \
  --topic partition-demo \
  --partitions 3 \
  --bootstrap-server localhost:9092

# 2. 发送带 Key 的消息
bin/kafka-console-producer.sh \
  --topic partition-demo \
  --bootstrap-server localhost:9092 \
  --property parse.key=true \
  --property key.separator=:
> user-1:message A from user 1
> user-2:message B from user 2
> user-1:message C from user 1
> user-3:message D from user 3
> user-2:message E from user 2

# 3. 查看各分区的消息数
bin/kafka-run-class.sh kafka.tools.GetOffsetShell \
  --topic partition-demo --time -1 --bootstrap-server localhost:9092
# partition-demo:0:2   ← user-1 的两条消息
# partition-demo:1:2   ← user-2 的两条消息
# partition-demo:2:1   ← user-3 的一条消息

操作前后对比:

消息Keyhash(key) % 3进入分区分区内 Offset
message Auser-10Partition 00
message Buser-21Partition 10
message Cuser-10Partition 01
message Duser-32Partition 20
message Euser-21Partition 11

核心验证:user-1 的消息全在 Partition 0,user-2 全在 Partition 1。这保证了"同一用户的消息有序"。

示例二:分区数对吞吐量的影响

场景:在同一 Topic、不同分区数下,使用 kafka-producer-perf-test.sh 测试吞吐量。

步骤:

# 测试 1 分区
bin/kafka-topics.sh --create --topic perf-1p --partitions 1 --bootstrap-server localhost:9092
bin/kafka-producer-perf-test.sh --topic perf-1p --num-records 1000000 \
  --record-size 100 --throughput -1 --producer-props linger.ms=5 \
  --bootstrap-server localhost:9092
# 1000000 records sent, 45623.7 records/sec (4.35 MB/sec)

# 测试 6 分区
bin/kafka-topics.sh --create --topic perf-6p --partitions 6 --bootstrap-server localhost:9092
bin/kafka-producer-perf-test.sh --topic perf-6p --num-records 1000000 \
  --record-size 100 --throughput -1 --producer-props linger.ms=5 \
  --bootstrap-server localhost:9092
# 1000000 records sent, 198342.1 records/sec (18.92 MB/sec)

操作后对比:

分区数吞吐量说明
145,623 records/sec单分区,单线程写入
6198,342 records/sec6 分区并行写入,~4.3x 提升

分析:吞吐量提升比例(4.3x)低于分区倍数(6x),因为共享磁盘 I/O 和网络带宽。实际提升取决于硬件。

易错场景

易错 1:以为多分区就能保证消息"完全并行"

场景:创建 10 分区 Topic,但只有一个 Producer 线程顺序发送,期望吞吐 10x 提升。

事实:并行度来自多个 Producer 实例或线程同时写入不同分区。单线程顺序发送时,消息仍然逐条发送(虽然可以轮换分区),但不会获得并行写入的吞吐红利。

// 错误:单线程,消息依然串行发送
for (int i = 0; i < 10000; i++) {
    producer.send(new ProducerRecord<>("topic", "msg-" + i));
}

// 正确:多线程并发写入不同分区
ExecutorService executor = Executors.newFixedThreadPool(4);
for (int i = 0; i < 10000; i++) {
    executor.submit(() -> producer.send(new ProducerRecord<>("topic", "msg-" + i)));
}

易错 2:分区数设置过大

场景:为"以后可能需要"设置了 200 个分区,但实际只有 3 个消费者。

后果:

  1. 每个 Broker 上打开大量文件句柄
  2. Controller 选举和元数据同步变慢
  3. 无分区 Key 的消费可能遇到大量空分区

面试高频考点

Q:Kafka 如何保证消息顺序?如果要全局有序怎么办?

A:

  • 分区内有序:默认保证——一个 Partition 内消息按写入 Offset 严格有序
  • 按 Key 有序:使用消息 Key 进行分区路由,相同 Key 的消息总是进入同一分区,保证该 Key 下的消息有序
  • 全局有序:分区数设为 1。缺点:丧失水平扩展能力,单分区吞吐是瓶颈。这是 Kafka 的架构性 trade-off——它优先选择了吞吐和扩展性
上一页
Topic
下一页
Offset