PHP 微服务与消息队列:RabbitMQ、Kafka 与事件驱动

系统覆盖 PHP 微服务与消息队列:RabbitMQ/Kafka/Redis Stream 客户端选型、生产者消费者模式、异步任务与队列、事件驱动架构、幂等与消息可靠性、Laravel 队列生态。

引言

PHP 单体变大的第一个解药不是「拆微服务」,而是先引入消息队列做异步解耦——发邮件、扣库存、更新搜索索引这些「可稍后做的事」扔进队列,主请求秒回。本文先讲队列的基础模型(生产/消费/确认/死信),再给 RabbitMQ、Kafka、Redis Stream 的 PHP 落地与选型,最后落到底层可靠性——幂等、乱序、消息丢失。

前置:/php-swoole-async/(常驻内存/协程)、/php-laravel-internals/(队列服务)。分布式见 [[distributed-systems]]。


目录


1. 为什么先上消息队列

痛点:同步请求里做重活(邮件、报表、第三方回调)拖垮主链路。

队列的价值:

价值说明
异步解耦主请求秒回,重活后台做
削峰填谷突发流量进队列,消费端匀速处理
故障隔离下游挂了消息还在,恢复后继续
可扩展加消费者实例水平扩容
// 同步(慢:邮件+通知都阻塞请求)
public function checkout(Order $order): void {
    $this->mailer->send($order->customerEmail(), '订单成功');
    $this->notifier->send($order->customerEmail(), '订单更新');
    return ['ok' => true];   // 用户等了 500ms+
}

// 异步(快:只扔队列,立即返回)
public function checkout(Order $order): void {
    OrderEvents::dispatch(new OrderCreated($order));   // 入队
    return ['ok' => true];   // 用户立即收到响应
}

记忆:队列把「必须做」和「可稍后做」分开——主链路只留必须做的。


2. 消息队列核心模型

通用术语(跨 RabbitMQ/Kafka/Redis 一致):

概念含义
Producer生产消息
Consumer消费消息
Queue队列(消息暂存)
Ack(确认)消费成功后告诉 broker 删消息
Dead Letter死信(多次失败的消息)
Requeue失败重入队

消费流程(必须可靠):

Producer → [Queue] → Consumer
                          ├ 成功 → Ack(删除)
                          ├ 失败 → Nack(重试/死信)
                          └ 崩溃 → 消息未 Ack,broker 重新投递

关键可靠性选择:

确认:手动 Ack 而非自动(自动 = 崩溃即丢)
重试:指数退避,别风暴重投
死信:重试到上限进死信队列人工处理

铁律:消费处理必须「先做业务,再 Ack」——先 Ack 后业务,崩溃就丢消息。


3. RabbitMQ:AMQP 与 PHP 客户端

RabbitMQ 基于 AMQP,核心是 Exchange + Queue 绑定(路由灵活)。

// 安装:composer require php-amqplib/php-amqplib
use PhpAmqpLib\Connection\AMQPStreamConnection;
use PhpAmqpLib\Message\AMQPMessage;

// 生产者
$conn = new AMQPStreamConnection('rabbitmq', 5672, 'guest', 'guest');
$ch = $conn->channel();
$ch->queue_declare('orders', durable: true);

$msg = new AMQPMessage(json_encode(['orderId' => 1]), [
    'delivery_mode' => AMQPMessage::DELIVERY_MODE_PERSISTENT,   // 持久化
]);
$ch->basic_publish($msg, exchange: '', routing_key: 'orders');

// 消费者(常驻进程运行)
$ch->basic_qos(null, 1, null);                     // 每次取一条
$ch->basic_consume('orders', callback: function ($msg) {
    try {
        process(json_decode($msg->getBody(), true));
        $ch->basic_ack($msg->getDeliveryTag());    // 成功才 Ack
    } catch (Throwable $e) {
        $ch->basic_nack($msg->getDeliveryTag(), requeue: false);  // 进死信
    }
});
while ($ch->is_consuming()) { $ch->wait(); }

Exchange 类型:

Exchange路由规则适用
direct精确 routing_key点对点任务
topic通配符 # *按主题广播/过滤
fanout广播到所有绑定队列事件广播

RabbitMQ 适合复杂路由、任务分配——点对点 + 主题广播都成熟。


4. Kafka:日志与高吞吐

Kafka 是分布式日志流(不删消息,按 offset 消费),适合海量事件、数据管道。

// 安装:composer require rdkafka(pecl rdkafka)
use RdKafka\Producer;
use RdKafka\KafkaConsumer;
use RdKafka\Message;

// 生产者
$producer = new Producer();
$producer->addBrokers('kafka:9092');
$topic = $producer->newTopic('user-events');
$topic->produce(0, 0, json_encode(['userId' => 1, 'action' => 'signup']));
$producer->flush(1000);

// 消费者(按消费者组 + offset 消费)
$conf = new RdKafka\Conf();
$conf->set('group.id', 'analytics-group');
$conf->set('auto.offset.reset', 'earliest');
$consumer = new KafkaConsumer($conf);
$consumer->subscribe(['user-events']);

while (true) {
    $msg = $consumer->consume(12000);
    if ($msg->err === RD_KAFKA_RESP_ERR_NO_ERROR) {
        handleEvent(json_decode($msg->payload, true));
        $consumer->commit($msg);                 // 提交 offset
    }
}

Kafka vs RabbitMQ:

维度RabbitMQKafka
模型队列(消费即删)日志(按 offset 重放)
吞吐中极高
延迟低中高
重放不支持支持(审计/数据管道)
典型任务队列事件流、日志、指标

记忆:任务分发用 RabbitMQ,海量事件流/数据管道用 Kafka。


5. Redis Stream:轻量队列

已有 Redis 时最轻的队列方案——Redis Stream(5.0+)提供持久化流:

// Redis Stream 生产
$redis->xAdd('task:email', '*', ['order_id' => 1, 'user_email' => 'a@x.com']);

// 消费(消费者组)
$redis->xGroup('CREATE', 'task:email', 'email-workers', '0', true);
$msgs = $redis->xReadGroup('email-workers', 'worker-1', 'task:email', count: 1, block: 5000);

foreach ($msgs as $msg) {
    process($msg['data']);
    $redis->xAck('task:email', 'email-workers', $msg['id']);   // 确认
}
方案适合不适合
Redis Stream中小流量、已有 Redis复杂路由/持久审计
RabbitMQ复杂路由、任务可靠性超高吞吐
Kafka海量事件流、数据管道简单小任务

Redis 的坑:内存型(持久化靠 RDB/AOF),丢消息容忍度低就别用纯 Redis——但 Stream 比 List 可靠(有消费者组与确认)。


6. Laravel 队列生态

Laravel 内置队列——同一条业务代码,换驱动不改业务:

// 定义 Job(命令模式,见 /php-oop-design-patterns/)
class SendOrderEmail implements ShouldQueue {
    public function __construct(private Order $order) {}
    public function handle(OrderMailer $mailer): void {
        $mailer->send($this->order->customerEmail(), '订单成功');
    }
}

// 分发(同步变异步)
OrderCreated::dispatch($order);                     // 事件
SendOrderEmail::dispatch($order)->onQueue('mail'); // 指定队列

驱动对比(config/queue.php):

驱动场景
database起步/无基础设施
redis生产主流(配 Redis Stream 或 List)
sqsAWS 生态
kafka(扩展)大流量事件

失败处理:

// 重试与退避
class SendOrderEmail implements ShouldQueue {
    public int $tries = 3;
    public array $backoff = [5, 30, 120];   // 依次等待
    public function failed(Order $order, Throwable $e): void {
        // 人工介入:标记、通知
    }
}
// 死信:超过 tries 进 failed_jobs 表

Laravel 让队列「配置化」——写业务的人不关心底层是 Redis 还是 SQS,运维切驱动即可。


7. 事件驱动与 CQRS 延伸

事件驱动架构(EDA):服务通过事件协作,不直接调用:

订单服务 ──订单已创建──▶ 通知服务
                 └──▶ 库存服务
                 └──▶ 分析服务
// 订单服务发事件
OrderCreated::dispatch($order);

// 各订阅者各自处理(互不知晓彼此)
// NotificationSubscriber / InventorySubscriber / AnalyticsSubscriber

CQRS 延伸:写用事务模型、读用查询模型(缓存/搜索引擎),靠事件同步——适合读多写复杂场景。

模式解决成本
事件驱动服务解耦一致性靠最终
CQRS读写分离两套模型维护
Saga分布式事务补偿逻辑复杂

路径建议:先事件驱动解耦,别一上来 CQRS/Saga——多数团队事件驱动就够。


8. 可靠性:幂等、乱序与丢失

三大可靠性问题:

① 消息重复 → 幂等消费:

// 消费者必须幂等(同一消息处理两次结果相同)
function processOrderCreated(array $data): void {
    // 用唯一键做去重(如订单号)
    $key = 'processed:order:' . $data['orderId'];
    if ($redis->get($key)) { return; }
    $redis->setex($key, 86400, 1);
    // ... 业务处理
}
// 或数据库唯一约束兜底

② 乱序:同一实体的多条消息可能乱序到达——给消息带 version/sequence,消费时校验。

③ 丢失:

生产端:确认发送成功(publisher confirm)
broker:持久化(durable + persist)
消费端:手动 Ack(业务成功后再 ack)

可靠性检查表:

环节措施
生产持久化消息 + confirm
broker副本(Kafka) / 持久化(Rabbit)
消费手动 Ack、先业务后 Ack
失败重试退避 + 死信队列
重复幂等键去重
乱序序列号校验

记忆:「恰好一次」是理想,工程上「至少一次 + 幂等」最稳——幂等是消息系统的必修课。


9. 从单体到服务拆分的路径

别急着拆微服务——按这个节奏演进:

第 1 步:单体 + 队列 → 慢操作异步化(本专题主题)
第 2 步:模块化单体 → 代码边界清晰、共享数据库
第 3 步:按「变更频率」拆 → 高频独立模块先拆成服务
第 4 步:服务间用事件驱动 → 不互相直接调用

拆分信号:

□ 一个模块每周发布拖累其他模块
□ 团队沟通成本随模块数陡增
□ 独立扩展需求明显(如搜索服务要单独扩容)

拆分的坑:分布式事务、跨服务查询、最终一致性的心智负担——没到痛点别拆。


10. 速查表

需求方案
任务异步化Laravel Queue / Redis Stream
复杂路由任务RabbitMQ(direct/topic)
海量事件流Kafka
已有 Redis 轻量队列Redis Stream
事件解耦事件驱动 + 订阅者
消息去重幂等键(Redis/唯一约束)
消息不丢持久化 + 手动 Ack
失败重试退避 + 死信队列
多实例消费消费者组

一句话记忆:队列先异步化,任务用 RabbitMQ、事件流用 Kafka、轻量用 Redis Stream;消费先业务后 Ack,幂等键防重复,死信接管烂消息;拆分从模块化开始,别一上来就微服务。


延伸阅读

  • /php-swoole-async/ — 常驻内存与协程消费端
  • /php-laravel-internals/ — 队列服务在框架内的接入
  • /php-oop-design-patterns/ — 观察者/命令模式的消息抽象
  • [[distributed-systems]] — 分布式一致性基础
  • [[kafka]] — Kafka 深入专题

继续阅读

探索更多技术文章

浏览归档,发现更多关于系统设计、工具链和工程实践的内容。

全部文章 返回首页

「php」更多文章

  1. PHP 面向对象与设计模式:SOLID、常用模式与 Laravel 实践
  2. PHP 静态分析与代码质量:PHPStan、Psalm、Rector 与 CI 门禁
  3. PHP 部署运维实战:Nginx、PHP-FPM、Docker 与 CI/CD