Kafka 学习文档

Kafka 学习文档:PHP 开发者版

一份围绕 Kafka 核心原理、PHP 接入实践、生产环境问题、面试高频考点 的系统化学习指南。


目录


学习目标

学完本文档后,你应该能够:

  1. 理解核心概念:Producer、Broker、Topic、Partition、Consumer、Consumer Group、Offset、Replica、Leader/Follower。
  2. 理解存储与吞吐原理:顺序写磁盘、Batch、Page Cache、零拷贝、Partition 并行。
  3. 掌握可靠性设计:ACK 机制、ISR、手动 Commit、At Least Once、幂等、Retry、DLQ。
  4. 掌握 PHP 实战:使用 php-rdkafka 实现 Producer 和 Consumer,封装服务层,处理异常与重试。
  5. 掌握生产级方案:Supervisor 进程管理、优雅退出、Outbox 模式、消费积压排查。
  6. 能够设计业务场景:订单事件流、退款异步处理、日志采集、用户行为分析。
  7. 能够应对面试:清晰回答 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 解决的三大核心问题

  1. 解耦(Decoupling)

    • 生产者不依赖消费者,消费者也不依赖生产者。
    • 新增下游服务无需修改上游代码。
  2. 异步(Asynchronous)

    • 将非核心业务(短信、积分、统计)从主流程剥离,降低接口延迟。
  3. 削峰(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-3

2.4 Topic(主题)

用于消息分类,例如:

order_created
order_paid
order_cancelled
refund_created
user_registered

2.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 D

Kafka 不断向 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 → D

5.2 业务级顺序保证方案

如果需要保证同一个订单的状态变更有序(创建 → 支付 → 发货 → 完成),可将 order_id 作为消息 Key。

Kafka 根据 Key 的 Hash 值选择 Partition:

$messageKey = (string) $orderId;
// 同一 order_id 永远进入同一 Partition
order 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, P3

Rebalance 的影响:

  • 期间可能暂停消费
  • 频繁 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 Follower

7.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(积压量): 20000

11.2 积压常见原因

  • Producer 速度 > Consumer 处理速度
  • Consumer 业务逻辑太慢(复杂 SQL、同步 HTTP 调用)
  • 数据库/Redis/第三方接口阻塞
  • Consumer 数量不足,或 Partition 数量不足
  • Consumer 异常退出或频繁 Rebalance

11.3 处理方案

  1. 检查 Consumer 健康状态:进程是否存活,日志是否有异常。
  2. 监控 Lag 增长趋势:判断是突发流量还是持续处理能力不足。
  3. 优化单条处理耗时:优化 SQL、加入缓存、减少同步调用。
  4. 扩容 Consumer:

    • Partition = 8,Consumer 从 2 增加到 8,可提升并发。
    • 注意:若 Partition = 4,Consumer 增加到 20 也不会提升并发(最多 4 个同时工作)。
  5. 增加 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.so

13.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);

可能出现:

  1. DB 成功,Kafka 失败 → 订单已创建,但下游消费者未收到消息。
  2. 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.log
  • numprocs=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
        │
        ▼
    Kibana

21.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 = 107374182400Partition 最大 100GB
Log Compactioncleanup.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 分区策略

场景策略
指定 Keyhash(key) % partition_num,相同 Key 进入同一 Partition
未指定 Key客户端按轮询或粘性策略分配
指定 Partition直接写入指定 Partition

业务设计建议:需要顺序保证的业务实体(如订单)必须设置 Key。


第24章 常见问题排查手册

24.1 消息没消费?排查清单

  1. Topic 中是否有消息?(kafka-console-consumer 验证)
  2. Producer 是否发送成功?(检查 flush 返回值)
  3. Consumer 进程是否在运行?(ps / Supervisor 状态)
  4. Consumer Group 是否正确?(是否被其他消费者抢占了 Partition)
  5. Consumer 是否订阅了正确的 Topic?
  6. Partition 是否正确分配给了 Consumer?(查看日志)
  7. Consumer 当前 Offset 在哪里?(是否跳过了消息)
  8. Consumer Lag 是多少?(是否积压严重)
  9. Consumer 日志是否有异常?(数据库连接失败、内存溢出等)
  10. 下游依赖(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 为什么快?

标准回答要点:

  1. 顺序写磁盘:追加写入,避免随机 IO,磁盘顺序写性能接近内存。
  2. Page Cache:利用操作系统缓存,热数据无需磁盘读取。
  3. 零拷贝:数据从 Page Cache 直接通过 sendfile 发送到网卡,减少 CPU 拷贝次数。
  4. Batch 批量处理:多条消息打包发送/写入,减少网络 RTT 和系统调用。
  5. Partition 并行:多 Partition 支持多线程并行读写。
  6. 消息压缩:减少网络和磁盘开销。

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 怎么选?

维度KafkaRabbitMQ
核心模型分布式日志流传统消息队列
吞吐极高(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 个核心主题,即可应对绝大多数生产场景和面试问题。