引言
PHP 单体变大的第一个解药不是「拆微服务」,而是先引入消息队列做异步解耦——发邮件、扣库存、更新搜索索引这些「可稍后做的事」扔进队列,主请求秒回。本文先讲队列的基础模型(生产/消费/确认/死信),再给 RabbitMQ、Kafka、Redis Stream 的 PHP 落地与选型,最后落到底层可靠性——幂等、乱序、消息丢失。
前置:/php-swoole-async/(常驻内存/协程)、/php-laravel-internals/(队列服务)。分布式见 [[distributed-systems]]。
目录
- 1. 为什么先上消息队列
- 2. 消息队列核心模型
- 3. RabbitMQ:AMQP 与 PHP 客户端
- 4. Kafka:日志与高吞吐
- 5. Redis Stream:轻量队列
- 6. Laravel 队列生态
- 7. 事件驱动与 CQRS 延伸
- 8. 可靠性:幂等、乱序与丢失
- 9. 从单体到服务拆分的路径
- 10. 速查表
- 延伸阅读
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:
| 维度 | RabbitMQ | Kafka |
|---|---|---|
| 模型 | 队列(消费即删) | 日志(按 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) |
| sqs | AWS 生态 |
| 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 深入专题
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。