Kafka 基础
Apache Kafka 是一个分布式的流处理平台,最初由 LinkedIn 开发,现已成为 Apache 顶级项目。Kafka 以高吞吐量、低延迟、持久化存储和水平扩展能力著称,广泛应用于日志收集、流数据处理、事件驱动架构和消息队列场景。本节将讲解 Kafka 的核心概念(Topic、Partition、Offset、Consumer Group)以及架构设计和安装配置。
前置知识
阅读本节前,建议先了解:RabbitMQ 基础 和 分布式系统基础概念
核心概念
什么是 Kafka
Kafka 是一个分布式、分区、可复制、基于日志提交的流处理平台。它的核心特性包括:
- 高吞吐量:单集群每秒可处理百万级消息
- 持久化存储:消息以日志形式持久化到磁盘,支持长期保存
- 水平扩展:通过增加 Broker 和 Partition 实现线性扩展
- 多消费者:同一消息可被多个 Consumer Group 独立消费
- 容错性:基于副本机制实现数据冗余
- 顺序保证:同一 Partition 内的消息保证顺序
Kafka 与 RabbitMQ 对比
| 特性 | Kafka | RabbitMQ |
|---|---|---|
| 设计目标 | 高吞吐日志流 | 灵活消息路由 |
| 消息模型 | 发布-订阅(Log) | 多种(Direct/Fanout/Topic) |
| 消息保留 | 基于时间/大小(可回溯) | 消费后删除 |
| 吞吐量 | 极高(MB/s 级) | 中等(万级/s) |
| 延迟 | 中等(ms 级) | 低(us 级) |
| 消息顺序 | Partition 内有序 | Queue 内有序 |
| 消费模型 | 拉取(Pull) | 推送/拉取(Push/Pull) |
| 消息回溯 | 支持(Offset 可重置) | 不支持 |
| 流量控制 | 分区 + 消费者组 | prefetch_count |
| 适用场景 | 大数据、日志、事件流 | 业务消息、任务队列 |
Topic 与 Partition
Topic(主题)
Topic 是消息的逻辑分类,类似数据库中的"表"。Producer 将消息发送到 Topic,Consumer 从 Topic 订阅消息。
Topic: order-events
├── Partition 0: [msg0, msg1, msg2, msg3, ...]
├── Partition 1: [msg0, msg1, msg2, msg3, ...]
├── Partition 2: [msg0, msg1, msg2, msg3, ...]
└── Partition 3: [msg0, msg1, msg2, msg3, ...]Partition(分区)
Partition 是 Topic 的物理分片,是 Kafka 并行处理的基本单位。
Partition 内部结构(追加日志):
+------+----------+--------+--------+--------+--------+
| Offset| Message | Message | Message | Message | ... |
+------+----------+--------+--------+--------+--------+
| 0 | msg_0 | msg_1 | msg_2 | msg_3 | |
+------+----------+--------+--------+--------+--------+
- 消息以追加(Append-only)方式写入
- 每条消息有唯一的 Offset(逻辑序号,从 0 开始递增)
- Offset 是 Partition 级别的(不同 Partition 的 Offset 独立)
- Consumer 通过 Offset 追踪消费位置关键概念
- Offset 只在 Partition 内唯一,不是全局唯一的
- 消费者可以重置 Offset 来重新消费历史消息
- Partition 数量决定了 Topic 的并行度(最大消费者数 = Partition 数)
- Partition 数量创建后可以增加,但不能减少
分区策略
Producer 发送消息到哪个 Partition,由分区策略决定:
| 策略 | 说明 | 适用场景 |
|---|---|---|
| 轮询(Round-robin) | 默认策略,消息均匀分配到各 Partition | 无 Key 的消息 |
| 按 Key 哈希 | hash(key) % numPartitions,相同 Key 到同一 Partition | 需要顺序保证的场景 |
| 手动指定 | 直接指定 Partition 编号 | 特殊路由需求 |
| 粘性分区(Sticky) | 尽量将同一批次发送到同一 Partition,减少请求 | 批量发送优化 |
php
// 分区策略示例
// 无 Key → 轮询分配
$producer->send(new TopicMessage('order-events', null, $payload));
// 有 Key → Hash 分区(同一 Key 的消息到同一 Partition)
$producer->send(new TopicMessage('order-events', 'order_id_1001', $payload));
// 手动指定分区
$producer->send(new TopicMessage('order-events', null, $payload, 2));Offset 与 Consumer Group
Offset
Offset 是消费者在 Partition 中的消费位置标记:
Partition 0 的消息和 Offset:
Consumer 已消费 Consumer 未消费
┌──────────┐┌──────────────────┐
│ msg_0 ││ msg_4 │
│ msg_1 ││ msg_5 │
│ msg_2 ││ msg_6 │
│ msg_3 ││ ... │
└──────────┘└──────────────────┘
↑ 已提交 Offset ↑ Log End Offset
(committed) (LEO)
- committed offset: 消费者确认已处理的消息位置
- position: 消费者当前正在处理的消息位置
- LEO (Log End Offset): 下一条待写入的消息位置Consumer Group(消费者组)
Consumer Group 是 Kafka 实现消息广播和负载均衡的核心机制。
Consumer Group 规则:
1. 同一 Consumer Group 内,每条消息只被一个 Consumer 消费(负载均衡)
2. 不同 Consumer Group 之间,消息独立消费(广播)
Topic: order-events (3 Partitions)
Consumer Group A:
Consumer A1 → Partition 0
Consumer A2 → Partition 1
Consumer A3 → Partition 2
Consumer Group B:
Consumer B1 → Partition 0, 1, 2(全部消费)
结果:
- Group A 的消息被 3 个 Consumer 分担
- Group B 的每个 Consumer 收到所有消息消费者组与分区分配
- Consumer 数量 > Partition 数:多余的 Consumer 闲置(不处理消息)
- Consumer 数量 = Partition 数:最优分配,每个 Consumer 处理一个 Partition
- Consumer 数量 < Partition 数:部分 Consumer 处理多个 Partition
- 理想情况:Consumer 数量 = Partition 数量
Rebalance(重平衡)
当 Consumer Group 中的成员发生变化时(新增/移除 Consumer),Kafka 会触发 Rebalance,重新分配 Partition 到 Consumer。
Rebalance 触发条件:
1. 新 Consumer 加入 Consumer Group
2. Consumer 离开 Consumer Group(崩溃/主动退出)
3. Topic 的 Partition 数量变化
4. 订阅的 Topic 变化
Rebalance 过程中:
- 所有 Consumer 暂停消费
- 重新分配 Partition
- Consumer 重新获取分配到的 Partition
- 恢复消费Rebalance 问题
Rebalance 期间消费者无法处理消息,可能导致:
- 消费延迟增加
- 如果 Rebalance 频繁发生(消费者频繁加入/离开),称为 "Rebalance Storm"
- 解决方案:合理设置
session.timeout.ms和heartbeat.interval.ms
架构与组件
Kafka 架构图
+------------+ +-------------+ +------------+
| Producer |────→│ Kafka │────→│ Consumer |
| (PHP App) │ │ Cluster │ | (PHP App) │
+------------+ +-------------+ +------------+
│ │
│ Broker 1 │
│ Broker 2 │
│ Broker 3 │
│ │
│ ZooKeeper │ (KRaft 模式下可不需要)
│ / KRaft │
+-------------+核心组件
| 组件 | 说明 |
|---|---|
| Broker | Kafka 服务节点,每个节点称为一个 Broker |
| Controller | 集群中的主控节点(从 Broker 中选举),管理 Partition 和副本 |
| ZooKeeper | 旧版 Kafka 使用 ZooKeeper 做集群管理(Kafka 2.8+ 可使用 KRaft 替代) |
| KRaft | Kafka 原生共识协议,替代 ZooKeeper(Kafka 3.3+ 默认使用) |
| Producer | 消息生产者 |
| Consumer | 消息消费者 |
| Admin Client | 管理客户端(创建 Topic、修改配置等) |
副本机制
Partition 的副本分布
Broker 1 Broker 2 Broker 3
┌──────────────┐ ┌──────────────┐ ┌──────────────┐
│ Partition 0 │ │ Partition 0 │ │ Partition 1 │
│ (Leader) │ │ (Follower) │ │ (Leader) │
├──────────────┤ ├──────────────┤ ├──────────────┤
│ Partition 1 │ │ Partition 2 │ │ Partition 2 │
│ (Follower) │ │ (Leader) │ │ (Follower) │
└──────────────┘ └──────────────┘ └──────────────┘
- Leader Replica: 处理读写请求
- Follower Replica: 从 Leader 同步数据(不直接服务客户端)
- ISR (In-Sync Replicas): 与 Leader 保持同步的副本集合
- Leader 故障时,从 ISR 中选举新的 Leader安装与配置
Docker 安装(KRaft 模式)
bash
# Kafka 3.4+ 支持纯 KRaft 模式(不需要 ZooKeeper)
# 使用官方 docker-compose 快速启动
# 单节点 Kafka(开发环境)
docker run -d \
--name kafka \
-p 9092:9092 \
-e KAFKA_NODE_ID=1 \
-e KAFKA_PROCESS_ROLES=broker,controller \
-e KAFKA_LISTENERS=PLAINTEXT://0.0.0.0:9092,CONTROLLER://0.0.0.0:9093 \
-e KAFKA_ADVERTISED_LISTENERS=PLAINTEXT://localhost:9092 \
-e KAFKA_CONTROLLER_QUORUM_VOTERS=1@kafka:9093 \
-e KAFKA_CONTROLLER_LISTENER_NAMES=CONTROLLER \
-e CLUSTER_ID=MkU3OEVBNTcwNTJENDM2Q0 \
apache/kafka:3.7.0Docker Compose 3 节点 KRaft 集群
yaml
# docker-compose.yml
version: '3.8'
services:
kafka1:
image: apache/kafka:3.7.0
container_name: kafka1
ports:
- "9092:9092"
environment:
KAFKA_NODE_ID: 1
KAFKA_PROCESS_ROLES: broker,controller
KAFKA_LISTENERS: PLAINTEXT://0.0.0.0:9092,CONTROLLER://0.0.0.0:9093
KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://kafka1:9092
KAFKA_CONTROLLER_QUORUM_VOTERS: 1@kafka1:9093,2@kafka2:9093,3@kafka3:9093
KAFKA_CONTROLLER_LISTENER_NAMES: CONTROLLER
CLUSTER_ID: MkU3OEVBNTcwNTJENDM2Q0
KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 3
KAFKA_TRANSACTION_STATE_LOG_REPLICATION_FACTOR: 3
KAFKA_TRANSACTION_STATE_LOG_MIN_ISR: 2
volumes:
- kafka1_data:/var/lib/kafka/data
kafka2:
image: apache/kafka:3.7.0
container_name: kafka2
ports:
- "9093:9092"
environment:
KAFKA_NODE_ID: 2
KAFKA_PROCESS_ROLES: broker,controller
KAFKA_LISTENERS: PLAINTEXT://0.0.0.0:9092,CONTROLLER://0.0.0.0:9093
KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://kafka2:9092
KAFKA_CONTROLLER_QUORUM_VOTERS: 1@kafka1:9093,2@kafka2:9093,3@kafka3:9093
KAFKA_CONTROLLER_LISTENER_NAMES: CONTROLLER
CLUSTER_ID: MkU3OEVBNTcwNTJENDM2Q0
KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 3
KAFKA_TRANSACTION_STATE_LOG_REPLICATION_FACTOR: 3
KAFKA_TRANSACTION_STATE_LOG_MIN_ISR: 2
volumes:
- kafka2_data:/var/lib/kafka/data
kafka3:
image: apache/kafka:3.7.0
container_name: kafka3
ports:
- "9094:9092"
environment:
KAFKA_NODE_ID: 3
KAFKA_PROCESS_ROLES: broker,controller
KAFKA_LISTENERS: PLAINTEXT://0.0.0.0:9092,CONTROLLER://0.0.0.0:9093
KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://kafka3:9092
KAFKA_CONTROLLER_QUORUM_VOTERS: 1@kafka1:9093,2@kafka2:9093,3@kafka3:9093
KAFKA_CONTROLLER_LISTENER_NAMES: CONTROLLER
CLUSTER_ID: MkU3OEVBNTcwNTJENDM2Q0
KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 3
KAFKA_TRANSACTION_STATE_LOG_REPLICATION_FACTOR: 3
KAFKA_TRANSACTION_STATE_LOG_MIN_ISR: 2
volumes:
- kafka3_data:/var/lib/kafka/data
volumes:
kafka1_data:
kafka2_data:
kafka3_data:常用管理命令
bash
# 创建 Topic
kafka-topics.sh --create \
--topic order-events \
--partitions 6 \
--replication-factor 3 \
--bootstrap-server localhost:9092
# 查看 Topic 列表
kafka-topics.sh --list --bootstrap-server localhost:9092
# 查看 Topic 详情
kafka-topics.sh --describe --topic order-events --bootstrap-server localhost:9092
# 修改 Partition 数量(只能增加)
kafka-topics.sh --alter \
--topic order-events \
--partitions 12 \
--bootstrap-server localhost:9092
# 删除 Topic
kafka-topics.sh --delete \
--topic order-events \
--bootstrap-server localhost:9092
# 控制台生产者
kafka-console-producer.sh \
--topic order-events \
--bootstrap-server localhost:9092
# 控制台消费者
kafka-console-consumer.sh \
--topic order-events \
--from-beginning \
--bootstrap-server localhost:9092
# 查看消费者组
kafka-consumer-groups.sh --list --bootstrap-server localhost:9092
# 查看消费者组详情
kafka-consumer-groups.sh --describe \
--group order-consumer-group \
--bootstrap-server localhost:9092注意事项
- Partition 数量规划:创建时根据业务量预估,创建后只能增加不能减少
- 副本因子:生产环境建议 replication-factor >= 3
- 消息大小:单条消息建议不超过 1MB(默认
max.message.bytes=1048576) - Offset 管理:Consumer 需要定期提交 Offset,避免 Rebalance 后重复消费
- 版本选择:新项目使用 Kafka 3.x + KRaft 模式,避免 ZooKeeper 运维开销
最佳实践
- Topic 规划:按业务领域划分 Topic,命名格式
{domain}.{entity}.{event-type} - Partition 数量:通常等于或略大于消费者数量,留出扩展空间
- 消息格式:使用 JSON/Avro/Protobuf,包含消息版本号和唯一 ID
- 清理策略:根据业务需求配置日志保留策略(时间/大小/压缩)
- 监控:监控消息速率、Partition 分布、Consumer Lag、Broker 资源
下一节
继续学习:Kafka PHP 客户端操作