Kafka 学习文档
Kafka 学习文档:PHP 开发者版
一份围绕 Kafka 核心原理、PHP 接入实践、生产环境问题、面试高频考点 的系统化学习指南。
目录
学习目标
学完本文档后,你应该能够:
- 理解核心概念:Producer、Broker、Topic、Partition、Consumer、Consumer Group、Offset、Replica、Leader/Follower。
- 理解存储与吞吐原理:顺序写磁盘、Batch、Page Cache、零拷贝、Partition 并行。
- 掌握可靠性设计:ACK 机制、ISR、手动 Commit、At Least Once、幂等、Retry、DLQ。
- 掌握 PHP 实战:使用
php-rdkafka实现 Producer 和 Consumer,封装服务层,处理异常与重试。 - 掌握生产级方案:Supervisor 进程管理、优雅退出、Outbox 模式、消费积压排查。
- 能够设计业务场景:订单事件流、退款异步处理、日志采集、用户行为分析。
- 能够应对面试:清晰回答 Kafka 为什么快、如何保证顺序/不丢消息、Partition 与 Consumer 关系等高频问题。
第一部分:Kafka 核心原理
第1章 Kafka 是什么,为什么需要它
1.1 Kafka 定义
Apache Kafka 是一个分布式事件流平台(Distributed Event Streaming Platform),日常开发中常被用作:
- 消息队列
- 事件总线
- 日志系统
- 数据同步通道
典型电商场景:
用户下单
│
▼
订单服务
│
▼
Kafka
│
├───▶ 库存服务
├───▶ 优惠券服务
├───▶ 积分服务
├───▶ 消息通知服务
└───▶ 数据统计服务订单服务只需发送一个 OrderCreated 事件,下游服务自行消费,无需逐个调用。
1.2 没有消息队列的问题
以用户下单后的操作为例:扣库存、扣优惠券、加积分、发短信、发 Push、更新统计、写操作日志。
同步调用的问题:
| 问题 | 说明 |
|---|---|
| 响应时间长 | 库存(50ms) + 优惠券(70ms) + 积分(40ms) + 短信(100ms) + ... ≈ 400ms+ |
| 系统耦合严重 | 订单服务必须知道所有下游服务,新增「推荐系统」可能需要修改订单代码 |
| 流量冲击 | 秒杀时 10万 QPS 直接打穿数据库(数据库只能处理 5000 QPS) |
使用 Kafka 后:
创建订单 → 写数据库 → 发送 Kafka → 立即返回(< 50ms)
│
└── 后续任务全部异步执行1.3 Kafka 解决的三大核心问题
解耦(Decoupling)
- 生产者不依赖消费者,消费者也不依赖生产者。
- 新增下游服务无需修改上游代码。
异步(Asynchronous)
- 将非核心业务(短信、积分、统计)从主流程剥离,降低接口延迟。
削峰(Peak Shaving)
- Kafka 作为流量缓冲池,先接住突发流量,Consumer 再以稳定速率处理。
- 秒杀场景:
10万请求 → Kafka → 5000/s 写数据库。
第2章 核心架构与概念
2.1 核心组件一览
| 组件 | 中文 | 职责 |
|---|---|---|
| Producer | 生产者 | 向 Kafka 发送消息 |
| Broker | 代理节点 | Kafka 服务器节点,多个 Broker 组成集群 |
| Topic | 主题 | 消息分类,类似数据库的 Table |
| Partition | 分区 | Topic 的物理分片,实现并行处理 |
| Consumer | 消费者 | 读取并处理 Kafka 消息 |
| Consumer Group | 消费者组 | 一组协同消费的 Consumer,共同消费一个 Topic |
| Offset | 偏移量 | 消息在 Partition 中的唯一位置标识 |
| Replica | 副本 | Partition 的数据拷贝,保证高可用 |
| Leader | 领导者 | 负责读写的主副本 |
| Follower | 跟随者 | 从 Leader 同步数据,Leader 宕机时可晋升 |
2.2 Producer(生产者)
职责:向指定 Topic 发送消息。
$producer->produce(
'order_created', // Topic 名
RD_KAFKA_PARTITION_UA, // 自动选择 Partition
json_encode(['order_id' => 10001]),
'10001' // Key(用于分区路由)
);2.3 Broker(服务器节点)
一台 Kafka Server 就是一个 Broker。生产环境通常部署为集群:
Kafka Cluster
├── Broker-1
├── Broker-2
└── Broker-32.4 Topic(主题)
用于消息分类,例如:
order_created
order_paid
order_cancelled
refund_created
user_registered2.5 Partition(分区)
一个 Topic 可拆分为多个 Partition,是 Kafka 水平扩展的核心。
Topic: order_created
├── Partition 0
├── Partition 1
└── Partition 2为什么需要 Partition?
- 提高吞吐:多 Partition 支持并行写入和并行消费。
- 水平扩展:增加 Partition 数量即可提升整体处理能力。
2.6 Consumer(消费者)
while (true) {
$message = $consumer->consume(1000);
if ($message->err === RD_KAFKA_RESP_ERR_NO_ERROR) {
$data = json_decode($message->payload, true);
handleOrder($data);
}
}2.7 Consumer Group(消费者组)
核心规则:同一个 Consumer Group 内,一个 Partition 同一时刻只能被一个 Consumer 消费。
Topic: order_created (P0, P1, P2)
Consumer Group: order-worker-group
├── Consumer A → 消费 P0
├── Consumer B → 消费 P1
└── Consumer C → 消费 P2第3章 Topic、Partition 与 Offset
3.1 Kafka 消息不会消费即删除
与传统队列不同,Kafka 中消息被消费后不会立即删除,而是根据 Retention Policy(保留策略)(如保存 7 天或达到指定大小)来决定清理时机。
传统队列:Producer → Queue → Consumer → 消息删除
Kafka: Producer → Kafka Log → Consumer 读取 → 消息依然存在3.2 Partition 本质是 Append-Only Log
Partition 内部是一个只追加日志(Append-Only Log):
Partition 0
├── offset 0 → message A
├── offset 1 → message B
├── offset 2 → message C
└── offset 3 → message DKafka 不断向 Partition 末尾追加消息,而非随机修改历史数据。
3.3 Offset(偏移量)
Offset 表示消息在 Partition 中的位置。Consumer 消费到 Offset 3 时,Kafka 会记录该 Consumer Group 的消费进度。
重要:Offset 是每个 Partition 独立维护的,而非 Topic 全局统一。
Topic: order_created
├── Partition 0 (Offset: 0, 1, 2, 3...)
├── Partition 1 (Offset: 0, 1, 2...)
└── Partition 2 (Offset: 0, 1, 2, 3, 4...)第4章 消息发送与高吞吐原理
4.1 Producer 发送流程
Producer
│
▼
序列化
│
▼
选择 Partition(根据 Key 或轮询)
│
▼
进入 Batch 缓冲区
│
▼
批量发送给 Broker
│
▼
写入 Partition(顺序追加)4.2 Kafka 为什么性能高
| 技术点 | 说明 |
|---|---|
| 顺序写磁盘 | 追加写入,避免随机 IO,磁盘顺序写速度接近内存 |
| Batch 批量发送 | 多条消息打包一次网络请求,减少 RTT |
| Page Cache | 利用操作系统缓存,减少磁盘读取 |
| 零拷贝(Zero Copy) | 数据直接从 Page Cache 发送到网卡,减少 CPU 拷贝 |
| Partition 并行 | 多 Partition 支持多线程并行读写 |
| 消息压缩 | 支持 gzip/snappy/lz4/zstd,减少网络和磁盘开销 |
第5章 消息顺序保证
5.1 Kafka 的顺序保证范围
Kafka 只能保证单个 Partition 内的消息有序,无法保证 Topic 全局有序。
P0: A → C P1: B → D
Kafka 无法保证全局顺序一定是 A → B → C → D5.2 业务级顺序保证方案
如果需要保证同一个订单的状态变更有序(创建 → 支付 → 发货 → 完成),可将 order_id 作为消息 Key。
Kafka 根据 Key 的 Hash 值选择 Partition:
$messageKey = (string) $orderId;
// 同一 order_id 永远进入同一 Partitionorder 10001 create ──┐
order 10001 pay ──┼──▶ Partition 2(按 hash(order_id) % partition数 路由)
order 10001 ship ──┘这样同一订单的所有状态变更都在同一个 Partition 内,自然保证了顺序性。
第6章 Consumer Group 与 Rebalance
6.1 Partition 与 Consumer 的分配关系
| 场景 | 分配结果 |
|---|---|
| Consumer < Partition | 部分 Consumer 消费多个 Partition |
| Consumer = Partition | 理想状态,1:1 分配,吞吐最优 |
| Consumer > Partition | 多余 Consumer 空闲,无法提升并发 |
核心结论:Consumer Group 的最大并发能力 受 Partition 数量限制。
6.2 Rebalance(重平衡)
当 Consumer Group 发生以下变化时,Kafka 会重新分配 Partition:
- 新 Consumer 加入
- Consumer 退出或宕机
- Partition 数量变化
示例:
变更前: 变更后(Consumer B 宕机):
A → P0, P1 A → P0, P1, P2, P3
B → P2, P3Rebalance 的影响:
- 期间可能暂停消费
- 频繁 Rebalance 会导致吞吐下降、消费延迟、消息积压
生产建议:
- 避免频繁扩缩容 Consumer
- 设置合理的
session.timeout.ms和heartbeat.interval.ms - 消费逻辑尽量轻量,避免处理超时导致被踢出 Group
第7章 副本机制与高可用
7.1 Replica(副本)
为保证数据不丢失,每个 Partition 可配置多个副本(replication.factor)。
Partition 0 (副本数=3)
├── Broker 1: P0 Leader
├── Broker 2: P0 Follower
└── Broker 3: P0 Follower7.2 Leader 与 Follower
- Leader:唯一接受 Producer 写入和 Consumer 读取的副本。
- Follower:从 Leader 拉取数据同步。Leader 宕机时,从 ISR 中选举新 Leader。
7.3 ISR(In-Sync Replicas)
ISR 是与 Leader 保持同步的副本集合。
Leader + Follower1 + Follower2 同步正常 → ISR = {Leader, F1, F2}
Leader + Follower1(正常) + Follower2(滞后太多) → ISR = {Leader, F1}滞后的 Follower 会被移出 ISR,直到追上进度后重新加入。
第8章 ACK 机制与消息可靠性
Producer 发送消息时通过 acks 参数控制可靠性级别:
8.1 acks = 0
- 行为:发送后不等待 Broker 确认。
- 优点:延迟最低,吞吐最高。
- 缺点:Broker 没收到消息也不会知道,最容易丢消息。
8.2 acks = 1
- 行为:Leader 写入成功即返回确认。
- 风险:Leader 刚写入成功但尚未同步给 Follower 时宕机,消息可能丢失。
8.3 acks = all(或 -1)
- 行为:Leader 写入成功,且 ISR 中所有副本都同步完成后才返回确认。
- 优点:可靠性最高,生产环境重要业务首选。
- 配合参数:
min.insync.replicas(最小同步副本数),例如设为 2,表示 ISR 中至少要有 2 个副本同步成功。
8.4 消息丢失的三端防护
| 环节 | 风险 | 解决方案 |
|---|---|---|
| Producer 端 | 网络异常导致发送失败 | acks=all + 配置重试机制 |
| Broker 端 | 单点宕机导致数据丢失 | 提高 replication.factor,合理配置 min.insync.replicas |
| Consumer 端 | 先提交 Offset 后处理业务,业务失败则消息丢失 | 业务成功后手动提交 Offset |
第9章 Offset 提交与三种消费语义
9.1 自动提交 vs 手动提交
自动提交(enable.auto.commit=true):
获取消息 → 自动提交 Offset → 业务处理失败
↓
Kafka 认为已消费,业务实际失败 → 消息丢失手动提交(推荐用于重要业务):
获取消息 → 处理业务 → 数据库成功 → commit offset
│
└── 失败 → 不提交 Offset → 消息可重新消费try {
handleMessage($message);
$consumer->commit($message);
} catch (\Throwable $e) {
// 不提交 offset,消息会再次投递
error_log($e->getMessage());
}9.2 三种消息语义
| 语义 | 中文 | 特点 | 典型实现 |
|---|---|---|---|
| At Most Once | 最多一次 | 可能丢失,不会重复 | 先提交 Offset,再处理业务 |
| At Least Once | 至少一次 | 不会丢失,可能重复 | 先处理业务,成功后再提交 Offset(最常用) |
| Exactly Once | 精确一次 | 不丢且不重复 | Kafka 层可通过事务和幂等 Producer 实现;业务层必须通过幂等保证 |
业务开发核心认知:即使 Kafka 做到了链路级的 Exactly Once,业务中的 MySQL、Redis、HTTP API 仍可能重复执行。因此业务幂等比追求 Exactly Once 更重要。
第10章 幂等设计
10.1 为什么必须幂等
Consumer 可能重复消费同一消息。例如:
{
"event_id": "evt_10001",
"order_id": 10001,
"type": "order_paid"
}若处理逻辑是「增加 100 积分」,消费两次会导致积分错误增加。
10.2 数据库唯一键方案(推荐)
建立消费记录表:
CREATE TABLE consumer_message_log (
event_id VARCHAR(64) NOT NULL,
consumer_name VARCHAR(64) NOT NULL,
created_at DATETIME NOT NULL,
PRIMARY KEY (event_id, consumer_name)
);处理流程:
开始事务
│
▼
INSERT INTO consumer_message_log (event_id, consumer_name, created_at)
│
├── 成功 → 执行业务逻辑 → 提交事务 → 提交 Kafka Offset
│
└── 主键冲突(DuplicateKeyException)→ 回滚 → 说明已处理过 → 直接提交 Kafka Offset$db->beginTransaction();
try {
DB::table('consumer_message_log')->insert([
'event_id' => $eventId,
'consumer_name' => 'score_consumer',
'created_at' => date('Y-m-d H:i:s'),
]);
DB::table('user_score')
->where('user_id', $userId)
->increment('score', 100);
$db->commit();
$consumer->commit($message);
} catch (DuplicateKeyException $e) {
$db->rollBack();
// 已处理过,直接提交 Offset,避免重复消费阻塞
$consumer->commit($message);
} catch (\Throwable $e) {
$db->rollBack();
// 不提交 Offset,等待重试
}10.3 其他幂等手段
- 业务唯一索引:如
order_id + event_type联合唯一键。 - Redis SETNX:作为前置去重,但需注意 Redis 与数据库的一致性。
- 状态机:如订单状态只能从
created → paid → shipped,paid → paid直接幂等拒绝。
第11章 消费积压排查
11.1 Consumer Lag(消费延迟)
Kafka 最新 Offset: 100000
Consumer 当前 Offset: 80000
─────────────────────────────────
Lag(积压量): 2000011.2 积压常见原因
- Producer 速度 > Consumer 处理速度
- Consumer 业务逻辑太慢(复杂 SQL、同步 HTTP 调用)
- 数据库/Redis/第三方接口阻塞
- Consumer 数量不足,或 Partition 数量不足
- Consumer 异常退出或频繁 Rebalance
11.3 处理方案
- 检查 Consumer 健康状态:进程是否存活,日志是否有异常。
- 监控 Lag 增长趋势:判断是突发流量还是持续处理能力不足。
- 优化单条处理耗时:优化 SQL、加入缓存、减少同步调用。
扩容 Consumer:
- Partition = 8,Consumer 从 2 增加到 8,可提升并发。
- 注意:若 Partition = 4,Consumer 增加到 20 也不会提升并发(最多 4 个同时工作)。
- 增加 Partition 数量:需要提前规划,已有数据的 Partition 扩容较复杂。
第12章 Kafka 与传统消息队列的区别
| 特性 | Kafka | 传统队列(如 RabbitMQ 部分模式) |
|---|---|---|
| 消费后数据 | 保留,按策略过期 | 通常消费即删除 |
| 读取方式 | 可重复读取(调整 Offset 即可) | 一般只能消费一次 |
| 吞吐能力 | 极高(十万级 QPS) | 中等(万级 QPS) |
| 核心设计 | 分布式日志流 | 任务队列 |
| 适用场景 | 日志、事件流、大数据、行为采集、数据同步 | 传统任务队列、复杂路由、优先级队列 |
不是谁更高级,而是业务模型不同。
第二部分:PHP + Kafka 实战
第13章 PHP 连接 Kafka
13.1 技术栈
Kafka Broker
│
▼
librdkafka(C/C++ 客户端库)
│
▼
php-rdkafka(PHP 扩展)
│
▼
PHP 业务代码安装扩展(以 Ubuntu 为例):
# 安装 librdkafka
apt-get install librdkafka-dev
# 安装 php-rdkafka 扩展
pecl install rdkafka
# 在 php.ini 中添加:extension=rdkafka.so13.2 连接配置要点
$conf = new RdKafka\Conf();
$conf->set('metadata.broker.list', 'kafka1:9092,kafka2:9092,kafka3:9092');第14章 Producer 封装与最佳实践
14.1 基础 Producer 示例
<?php
$conf = new RdKafka\Conf();
$conf->set('metadata.broker.list', '127.0.0.1:9092');
$producer = new RdKafka\Producer($conf);
$topic = $producer->newTopic('order_created');
$data = [
'event_id' => uniqid('evt_', true),
'order_id' => 10001,
'user_id' => 20001,
'created_at' => time(),
];
// 参数:partition, msgflags, payload, key
$topic->produce(
RD_KAFKA_PARTITION_UA, // 让 Kafka 自动选择 Partition
0,
json_encode($data, JSON_UNESCAPED_UNICODE),
(string) $data['order_id'] // Key:相同 order_id 进入同一 Partition
);
$producer->poll(0);
// 确保消息发送成功(重要!)
for ($i = 0; $i < 10; $i++) {
$result = $producer->flush(1000);
if ($result === RD_KAFKA_RESP_ERR_NO_ERROR) {
break;
}
}14.2 生产环境封装
不要直接在业务代码中 new Producer(),建议封装为服务:
<?php
class KafkaProducerService
{
private RdKafka\Producer $producer;
public function __construct()
{
$conf = new RdKafka\Conf();
$conf->set('metadata.broker.list', getenv('KAFKA_BROKERS'));
$conf->set('acks', 'all'); // 最高可靠性
$conf->set('retries', 3); // 发送失败重试
$conf->set('compression.type', 'snappy'); // 压缩
$this->producer = new RdKafka\Producer($conf);
}
public function send(
string $topicName,
array $payload,
?string $key = null
): void {
$topic = $this->producer->newTopic($topicName);
$topic->produce(
RD_KAFKA_PARTITION_UA,
0,
json_encode($payload, JSON_UNESCAPED_UNICODE),
$key
);
$this->producer->poll(0);
// 同步 flush,确保发送成功
$this->producer->flush(5000);
}
}业务层调用:
$kafkaProducer->send(
'order_paid',
[
'event_id' => $eventId,
'order_id' => $orderId,
'pay_time' => time(),
],
(string) $orderId // 保证同一订单有序
);第15章 Consumer 实现与生产环境结构
15.1 基础 Consumer 示例
<?php
$conf = new RdKafka\Conf();
$conf->set('metadata.broker.list', '127.0.0.1:9092');
$conf->set('group.id', 'order-consumer-group');
$conf->set('enable.auto.commit', 'false'); // 手动提交
$conf->set('auto.offset.reset', 'earliest'); // 从头消费(新 Group)
$consumer = new RdKafka\KafkaConsumer($conf);
$consumer->subscribe(['order_created']);
while (true) {
$message = $consumer->consume(1000);
switch ($message->err) {
case RD_KAFKA_RESP_ERR_NO_ERROR:
try {
$data = json_decode($message->payload, true, 512, JSON_THROW_ON_ERROR);
handleOrderCreated($data);
$consumer->commit($message); // 业务成功后提交
} catch (\Throwable $e) {
error_log($e->getMessage());
// 不提交 Offset,消息会再次投递
}
break;
case RD_KAFKA_RESP_ERR__PARTITION_EOF:
// 当前 Partition 无新消息
break;
case RD_KAFKA_RESP_ERR__TIMED_OUT:
// 消费超时
break;
default:
throw new RuntimeException($message->errstr(), $message->err);
}
}15.2 生产环境项目结构
不建议将所有逻辑写在 while(true) 中,推荐分层架构:
app/
└── Kafka/
├── Producer/
│ └── KafkaProducer.php
├── Consumer/
│ └── KafkaConsumer.php
├── Message/
│ └── KafkaMessage.php
├── Handler/
│ ├── OrderCreatedHandler.php
│ ├── OrderPaidHandler.php
│ ├── OrderCancelledHandler.php
│ └── RefundCreatedHandler.php
├── Router/
│ └── MessageRouter.php
└── Exception/
├── RetryableException.php
└── NonRetryableException.php职责划分:
KafkaConsumer:负责拉取消息、异常捕获、Offset 管理。MessageRouter:根据event_type路由到对应 Handler。Handler:执行业务逻辑,抛出RetryableException或NonRetryableException。KafkaConsumer根据异常类型决定是否重试或转入 DLQ。
第16章 消息格式设计
不要发送无结构的裸 JSON,推荐统一事件信封格式:
{
"event_id": "01JXXXXX",
"event_type": "order.paid",
"event_version": "1.0",
"occurred_at": 1720000000,
"source": "order-service",
"trace_id": "trace_xxx",
"data": {
"order_id": 10001,
"user_id": 20001,
"pay_money": 99.9
}
}好处:
- 追踪:通过
trace_id串联全链路。 - 幂等:通过
event_id实现去重。 - 版本升级:
event_version支持格式演进。 - 故障排查:明确
source和occurred_at。
第17章 MySQL 与 Kafka 一致性(Outbox 模式)
17.1 问题场景
DB::table('orders')->insert($order);
$kafka->send('order_created', $event);可能出现:
- DB 成功,Kafka 失败 → 订单已创建,但下游消费者未收到消息。
- Kafka 成功,DB 回滚 → 下游消费者处理了不存在的订单。
17.2 Transactional Outbox(事务发件箱)
核心思想:将「业务操作」和「待发送事件」放在同一个本地事务中。
BEGIN;
INSERT INTO orders (...);
INSERT INTO outbox_events (
topic,
payload,
created_at
) VALUES (
'order_created',
'{"order_id": 10001, ...}',
NOW()
);
COMMIT;后台 Worker 轮询:
Outbox Worker
│
▼
SELECT * FROM outbox_events WHERE sent = 0
│
▼
发送 Kafka
│
▼
UPDATE outbox_events SET sent = 1 WHERE id = ?保证:订单和事件记录要么同时成功,要么同时失败,最终通过 Worker 保证消息一定发出。
第18章 重试机制与死信队列
18.1 区分错误类型
| 错误类型 | 示例 | 策略 |
|---|---|---|
| 可重试(Retryable) | 数据库超时、网络抖动、第三方 API 限流 | 延迟重试 |
| 不可重试(Non-Retryable) | 数据格式错误、业务规则校验失败 | 直接转入 DLQ |
18.2 Retry Topic 设计
order_created
│
├── 处理成功 → 提交 Offset
│
└── 处理失败(可重试)
│
▼
order_created_retry_1m (延迟 1 分钟)
│
└── 再次失败
│
▼
order_created_retry_5m (延迟 5 分钟)
│
└── 再次失败
│
▼
order_created_retry_30m (延迟 30 分钟)
│
└── 仍然失败
│
▼
order_created_dlq (死信队列)实现方式:
- 独立 Consumer 消费 Retry Topic,通过
sleep或定时任务实现延迟。 - 或使用 Kafka 的
delivery.timeout.ms配合独立重试服务。
18.3 Dead Letter Queue(DLQ)
DLQ 存储最终处理失败的消息,用于:
- 人工排查:分析失败原因,修复后重新投递。
- 自动修复:针对特定错误类型编写补偿脚本。
- 监控告警:DLQ 有消息进入时触发告警。
第19章 PHP Consumer 长期运行方案
19.1 PHP 的特殊性
PHP Web 模式是「请求 → 启动 → 执行 → 退出」,而 Kafka Consumer 需要长期驻留内存。
解决方案:CLI 模式 + 进程管理工具。
19.2 Supervisor 配置示例
[program:kafka_order_consumer]
command=php /var/www/project/artisan kafka:consume order_created
directory=/var/www/project
autostart=true
autorestart=true
numprocs=4
redirect_stderr=true
stdout_logfile=/var/log/kafka-order-consumer.lognumprocs=4:启动 4 个 Consumer 进程(需确保 Partition ≥ 4)。autorestart=true:进程异常退出后自动重启。
19.3 优雅退出(Graceful Shutdown)
避免 kill -9 直接终止,应捕获信号完成当前任务后再退出:
pcntl_async_signals(true);
$running = true;
pcntl_signal(SIGTERM, function () use (&$running) {
$running = false;
});
pcntl_signal(SIGINT, function () use (&$running) {
$running = false;
});
while ($running) {
$message = $consumer->consume(1000);
if ($message->err === RD_KAFKA_RESP_ERR_NO_ERROR) {
try {
handleMessage($message);
$consumer->commit($message);
} catch (\Throwable $e) {
error_log($e->getMessage());
}
}
}
// 收到终止信号后,正常结束当前任务、提交 Offset、退出19.4 内存管理策略
PHP 长期运行进程容易出现内存泄漏(静态变量、大数组、未释放的数据库连接等)。
生产环境常见策略:
$processedCount = 0;
$maxMessagesBeforeRestart = 10000;
while ($running) {
// ... 消费逻辑 ...
if (++$processedCount >= $maxMessagesBeforeRestart) {
// 处理一定数量后主动退出,Supervisor 会自动重启
break;
}
}第三部分:业务场景设计
第20章 订单与退款系统
20.1 订单事件流设计
Topic 设计:
├── order.created → 消费者:库存、优惠券、通知、统计
├── order.paid → 消费者:积分、会员、营销、推荐、短信、数据仓库
├── order.shipped → 消费者:物流跟踪、用户通知
├── order.completed → 消费者:售后窗口计算、用户等级更新
└── order.cancelled → 消费者:库存回滚、优惠券返还、退款触发20.2 退款系统
业务服务 → 创建退款记录 → refund.created Topic
│
▼
Refund Consumer Worker
│
▼
调用第三方支付接口
│
▼
更新退款状态 → 发送 refund.completed适合 Kafka 的特征:
- 耗时操作(调用第三方接口)
- 可能失败(需要重试)
- 需要削峰(退款高峰期缓冲)
- 需要可追溯(通过事件日志审计)
第21章 日志与用户行为采集
21.1 日志系统架构
PHP 应用 A/B/C
│
▼
Kafka (log_topic)
│
▼
Log Consumer / Logstash
│
▼
Elasticsearch
│
▼
Kibana21.2 用户行为采集
行为事件:page_view / product_click / add_cart / search / order_created
│
▼
Kafka
│
├───▶ 实时计算(Flink/Spark)
├───▶ 推荐系统
├───▶ 数据分析(BI)
└───▶ 数据仓库(离线 T+1 计算)第四部分:高级原理与排查
第22章 存储结构与 Retention
22.1 Partition 的物理存储
Partition 在磁盘上被切分为多个 Segment:
Partition 0/
├── 00000000000000000000.log ← 消息数据
├── 00000000000000000000.index ← Offset 索引
├── 00000000000000000000.timeindex ← 时间索引
├── 00000000000000368792.log
├── 00000000000000368792.index
└── ...- .log:实际消息数据。
- .index:Offset 到物理位置的映射,加速消费定位。
- .timeindex:时间戳索引,支持按时间查找。
22.2 Retention(保留策略)
| 策略 | 配置 | 说明 |
|---|---|---|
| 时间保留 | retention.ms = 604800000 | 保留 7 天 |
| 大小保留 | retention.bytes = 107374182400 | Partition 最大 100GB |
| Log Compaction | cleanup.policy = compact | 保留同一 Key 的最新值,适合状态同步 |
Log Compaction 示例:
Key=user_10001, Value={"name":"Tom"}
Key=user_10001, Value={"name":"Jack"}
Key=user_10002, Value={"name":"Alice"}
Compaction 后:
Key=user_10001, Value={"name":"Jack"} ← 保留最新
Key=user_10002, Value={"name":"Alice"}第23章 集群高可用与 Partition 分配
23.1 集群部署示例
Broker 1 Broker 2 Broker 3
───────── ───────── ─────────
P0 Leader P0 Follower P0 Follower
P1 Follower P1 Leader P1 Follower
P2 Follower P2 Follower P2 Leader- Leader 分散在不同 Broker,实现负载均衡。
- 任一 Broker 宕机,其 Partition 的 Follower 可晋升为 Leader,服务不中断。
23.2 Producer 分区策略
| 场景 | 策略 |
|---|---|
| 指定 Key | hash(key) % partition_num,相同 Key 进入同一 Partition |
| 未指定 Key | 客户端按轮询或粘性策略分配 |
| 指定 Partition | 直接写入指定 Partition |
业务设计建议:需要顺序保证的业务实体(如订单)必须设置 Key。
第24章 常见问题排查手册
24.1 消息没消费?排查清单
- Topic 中是否有消息?(
kafka-console-consumer验证) - Producer 是否发送成功?(检查
flush返回值) - Consumer 进程是否在运行?(
ps/ Supervisor 状态) - Consumer Group 是否正确?(是否被其他消费者抢占了 Partition)
- Consumer 是否订阅了正确的 Topic?
- Partition 是否正确分配给了 Consumer?(查看日志)
- Consumer 当前 Offset 在哪里?(是否跳过了消息)
- Consumer Lag 是多少?(是否积压严重)
- Consumer 日志是否有异常?(数据库连接失败、内存溢出等)
- 下游依赖(DB/Redis/API)是否阻塞了消费?
24.2 消息重复消费?
认知纠正:Kafka 的 At Least Once 语义本身就允许重复,这不是 Bug。
解决方案:业务层做幂等,优先使用数据库唯一索引(event_id + consumer_name)。
24.3 消息乱序?
- 检查同一业务实体的消息是否设置了相同的 Key。
- 检查是否因重试导致乱序(重试 Topic 可能破坏顺序)。
- 若严格顺序要求,考虑单 Partition(吞吐会降低)。
24.4 消费越来越慢?
1. 查看 Lag 增长趋势
│
▼
2. 计算单条消息平均处理耗时
例:200ms/条 × 4 Consumer = 20条/s
若 Producer 100条/s → 必然积压
│
▼
3. 优化方案:
- 优化 SQL、加缓存、减少同步 HTTP 调用
- 增加 Partition 和 Consumer
- 改为批量处理(一次消费多条)第五部分:面试重点
第25章 高频面试题精讲
Q1:Kafka 为什么快?
标准回答要点:
- 顺序写磁盘:追加写入,避免随机 IO,磁盘顺序写性能接近内存。
- Page Cache:利用操作系统缓存,热数据无需磁盘读取。
- 零拷贝:数据从 Page Cache 直接通过
sendfile发送到网卡,减少 CPU 拷贝次数。 - Batch 批量处理:多条消息打包发送/写入,减少网络 RTT 和系统调用。
- Partition 并行:多 Partition 支持多线程并行读写。
- 消息压缩:减少网络和磁盘开销。
Q2:Kafka 如何保证消息不丢?
三端防护:
- Producer:
acks=all+ 重试机制。 - Broker:提高副本因子
replication.factor,配置min.insync.replicas。 - Consumer:业务处理成功后手动提交 Offset。
Q3:Kafka 如何保证消息顺序?
- Kafka 只保证单个 Partition 内有序。
- 业务需要顺序时,将业务实体 ID 作为 Key(如
order_id),确保同一实体进入同一 Partition。 - 全局有序需单 Partition,但会牺牲吞吐。
Q4:如何避免重复消费?
- 严格来说,Consumer 不能假设消息永不重复。
- 业务层必须做幂等:数据库唯一键(
event_id)、业务状态机、Redis SETNX 等。 - 优先使用数据库唯一约束,因为 Redis 与 DB 之间仍可能存在一致性窗口。
Q5:Partition 和 Consumer 的关系?
- 同一 Consumer Group 内,一个 Partition 同一时间最多分配给一个 Consumer。
- 因此 最大并发数 ≤ Partition 数量。
- Consumer 数量超过 Partition 数量时,多余 Consumer 空闲。
Q6:Kafka 和 RabbitMQ 怎么选?
| 维度 | Kafka | RabbitMQ |
|---|---|---|
| 核心模型 | 分布式日志流 | 传统消息队列 |
| 吞吐 | 极高(10万级 QPS) | 高(万级 QPS) |
| 消息保留 | 持久化、可重放 | 消费后通常删除 |
| 复杂路由 | 较弱(主要按 Topic) | 强(Exchange 灵活路由) |
| 优先级队列 | 不支持 | 支持 |
| 适用场景 | 日志、事件流、大数据、行为采集 | 任务队列、复杂路由、优先级 |
结论:不是谁更高级,而是业务模型不同。高吞吐事件流选 Kafka,复杂路由任务队列选 RabbitMQ。
第六部分:学习路径与实战建议
第26章 推荐学习顺序
第一阶段:基础概念(1-2 天)
掌握:Producer → Broker → Topic → Partition → Consumer → Consumer Group → Offset。
目标:能画出完整的数据流图。
第二阶段:核心原理(2-3 天)
掌握:消息存储、Partition 并行、Rebalance、Leader/Follower、Replica、ISR、ACK 机制。
目标:能回答「Kafka 为什么快」「为什么能扩展」「为什么不轻易丢数据」。
第三阶段:可靠性设计(3-5 天)⭐ 最重要
掌握:消息丢失场景、重复消费、乱序、手动 Commit、At Least Once、幂等、Retry、DLQ。
目标:能设计一个高可靠的消息消费方案。
第四阶段:PHP 实战(3-5 天)
掌握:php-rdkafka 安装、Producer 封装、Consumer 长驻进程、手动 Commit、异常处理。
目标:能写出生产可用的 Producer 和 Consumer。
第五阶段:生产级设计(持续)
掌握:Outbox 模式、Retry Topic、DLQ、Supervisor 管理、优雅退出、Lag 监控、消息追踪。
目标:能独立设计并落地 Kafka 事件流系统。
第六阶段:高级原理(进阶)
掌握:Segment、Index、Page Cache、Zero Copy、Log Compaction、KRaft(Kafka 3.x 去 ZooKeeper)。
第27章 PHP Kafka 实战项目建议
建议实现一个完整的 Kafka 订单事件系统,包含以下内容:
Topic 设计
order.created
order.paid
order.cancelled
refund.created
notification.send必须实现的功能
| 功能 | 要求 |
|---|---|
| Producer 封装 | 统一发送接口,支持 Key 路由,配置 acks=all |
| Consumer Group | 多进程并行消费,合理分配 Partition |
| 手动 Commit | 业务成功后提交 Offset |
| 幂等消费 | 基于 event_id + DB 唯一键 实现 |
| Retry 机制 | 可重试错误进入 Retry Topic(1m/5m/30m) |
| DLQ | 最终失败消息进入死信队列 |
| Supervisor | 管理 PHP Consumer 进程,配置 autorestart |
| 优雅退出 | 捕获 SIGTERM/SIGINT,完成当前任务后退出 |
| Lag 监控 | 监控 Consumer Lag,超过阈值告警 |
| Outbox 模式 | 解决 MySQL 与 Kafka 的一致性 |
项目结构参考
project/
├── app/
│ ├── Kafka/
│ │ ├── ProducerService.php
│ │ ├── ConsumerCommand.php
│ │ ├── MessageRouter.php
│ │ └── OutboxPublisher.php
│ └── Handlers/
│ ├── InventoryHandler.php
│ ├── ScoreHandler.php
│ ├── NotificationHandler.php
│ └── RefundHandler.php
├── database/
│ └── migrations/
│ ├── create_orders_table.php
│ └── create_outbox_events_table.php
├── supervisor/
│ └── kafka_consumers.conf
└── README.md附录:Kafka 核心知识图谱
Kafka
│
┌───────────────────────┼───────────────────────┐
│ │ │
写入 存储 消费
│ │ │
Producer Topic Consumer
│ │ │
ACK Partition Consumer Group
│ │ │
Retry Offset Rebalance
│ │
Batch Segment
│ │
Compression ┌────────┴────────┐
│ │
Leader Follower
│
ISR
┌─────────────────────────────────────────────────────────────┐
│ 业务可靠性层 │
├──────────┬──────────┬──────────┬──────────┬─────────────────┤
│ 消息丢失 │ 重复消费 │ 乱序 │ Retry │ DLQ │
├──────────┴──────────┴──────────┴──────────┴─────────────────┤
│ 幂等 │
│ DB Unique Key │
└─────────────────────────────────────────────────────────────┘
┌─────────────────────────────────────────────────────────────┐
│ PHP 工程层 │
├──────────────┬──────────────┬──────────────┬────────────────┤
│ php-rdkafka │ Supervisor │ Outbox │ Monitor │
│ │ │ │ │
│ Producer/ │ 进程管理 │ MySQL+Kafka │ Lag 监控 │
│ Consumer │ 自动重启 │ 一致性 │ 告警通知 │
└──────────────┴──────────────┴──────────────┴────────────────┘最后建议:不要一次性硬背所有内容。按章节逐步深入,重点掌握 Partition、Consumer Group、Offset、ACK、幂等、Outbox、Retry/DLQ、PHP Consumer 长驻进程 这 8 个核心主题,即可应对绝大多数生产场景和面试问题。