Skip to content

Kafka 基础

Apache Kafka 是一个分布式的流处理平台,最初由 LinkedIn 开发,现已成为 Apache 顶级项目。Kafka 以高吞吐量、低延迟、持久化存储和水平扩展能力著称,广泛应用于日志收集、流数据处理、事件驱动架构和消息队列场景。本节将讲解 Kafka 的核心概念(Topic、Partition、Offset、Consumer Group)以及架构设计和安装配置。

前置知识

阅读本节前,建议先了解:RabbitMQ 基础分布式系统基础概念

核心概念

什么是 Kafka

Kafka 是一个分布式、分区、可复制、基于日志提交的流处理平台。它的核心特性包括:

  • 高吞吐量:单集群每秒可处理百万级消息
  • 持久化存储:消息以日志形式持久化到磁盘,支持长期保存
  • 水平扩展:通过增加 Broker 和 Partition 实现线性扩展
  • 多消费者:同一消息可被多个 Consumer Group 独立消费
  • 容错性:基于副本机制实现数据冗余
  • 顺序保证:同一 Partition 内的消息保证顺序

Kafka 与 RabbitMQ 对比

特性KafkaRabbitMQ
设计目标高吞吐日志流灵活消息路由
消息模型发布-订阅(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.msheartbeat.interval.ms

架构与组件

Kafka 架构图

+------------+     +-------------+     +------------+
|  Producer  |────→│  Kafka      │────→│  Consumer  |
|  (PHP App) │     │  Cluster    │     |  (PHP App) │
+------------+     +-------------+     +------------+
                    │             │
                    │ Broker 1    │
                    │ Broker 2    │
                    │ Broker 3    │
                    │             │
                    │ ZooKeeper   │ (KRaft 模式下可不需要)
                    │ / KRaft     │
                    +-------------+

核心组件

组件说明
BrokerKafka 服务节点,每个节点称为一个 Broker
Controller集群中的主控节点(从 Broker 中选举),管理 Partition 和副本
ZooKeeper旧版 Kafka 使用 ZooKeeper 做集群管理(Kafka 2.8+ 可使用 KRaft 替代)
KRaftKafka 原生共识协议,替代 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.0

Docker 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

注意事项

  1. Partition 数量规划:创建时根据业务量预估,创建后只能增加不能减少
  2. 副本因子:生产环境建议 replication-factor >= 3
  3. 消息大小:单条消息建议不超过 1MB(默认 max.message.bytes=1048576
  4. Offset 管理:Consumer 需要定期提交 Offset,避免 Rebalance 后重复消费
  5. 版本选择:新项目使用 Kafka 3.x + KRaft 模式,避免 ZooKeeper 运维开销

最佳实践

  1. Topic 规划:按业务领域划分 Topic,命名格式 {domain}.{entity}.{event-type}
  2. Partition 数量:通常等于或略大于消费者数量,留出扩展空间
  3. 消息格式:使用 JSON/Avro/Protobuf,包含消息版本号和唯一 ID
  4. 清理策略:根据业务需求配置日志保留策略(时间/大小/压缩)
  5. 监控:监控消息速率、Partition 分布、Consumer Lag、Broker 资源

下一节

继续学习:Kafka PHP 客户端操作

参考链接