学完这篇,你能在 PHP 项目里分清 Redis、RabbitMQ、Kafka 三种队列各自适合什么场景,并把能跑通的生产者和消费者代码写出来。
第一步:先搞清楚你要的是"队列"还是"日志流"
PHP 里说"消息队列",通常落在三种需求上:
- 任务异步化:注册后发邮件、生成缩略图、推送通知。要求简单、能重试,Redis 就够。
- 业务解耦 + 复杂路由:一个订单事件要分发给库存、积分、风控三个系统,各自不同规则。选 RabbitMQ。
- 高吞吐数据管道:埋点日志、行为流、多消费者重复消费。选 Kafka。
注意:不要一上来就上 Kafka。中小站点的异步任务用 Redis List 实现,代码不到 20 行,运维成本几乎为零;Kafka 需要额外部署 ZooKeeper/KRaft 和 C 扩展,共享主机往往装不上。
第二步:Redis 做队列——最轻的落地方式
安装客户端(二选一,phpredis 是 C 扩展性能更好):
composer require predis/predis
生产者和消费者:
// 生产者
$redis->lPush('queue:mail', json_encode(['to' => 'a@b.com', 'tpl' => 'welcome']));
// 消费者(常驻脚本 worker.php)
while (true) {
$job = $redis->brPop(['queue:mail'], 5); // 阻塞 5 秒
if (!$job) continue;
$data = json_decode($job[1], true);
// ... 执行业务
}
**但 BRPOP 有一条硬伤:取出的瞬间消息就从队列消失了,进程崩溃或执行抛异常,这条消息就永久丢失。** 生产环境要用可靠版本:
// 取出后放进 processing 备份队列
$job = $redis->brpoplpush('queue:mail', 'queue:mail:processing', 5);
// 处理成功
$redis->lRem('queue:mail:processing', $job, 1);
Redis 5.0 以上更推荐用 Stream,它自带消费者组和 ACK 机制,`XADD` 投递、`XREADGROUP` 消费、`XACK` 确认,还有 `XPENDING` 可以查未确认消息、`XAUTOCLAIM` 把超时未确认的消息转给别的消费者。语义上已经接近专业队列了。
第三步:RabbitMQ——需要路由和延迟队列时用它
composer require php-amqplib/php-amqplib
核心模型是 exchange(交换机)+ queue(队列)+ binding(绑定规则),所以能做到 Redis 做不到的"发一次、多个队列各取所需"。
$conn = new AMQPStreamConnection('127.0.0.1', 5672, 'guest', 'guest');
$ch = $conn->channel();
$ch->exchange_declare('order', 'topic', false, true, false);
$ch->queue_declare('order.stock', false, true, false, false);
$ch->queue_bind('order.stock', 'order', 'order.created');
$msg = new AMQPMessage(json_encode(['id' => 1001]), [
'delivery_mode' => AMQPMessage::DELIVERY_MODE_PERSISTENT,
]);
$ch->basic_publish($msg, 'order', 'order.created');
// 消费端:一次只取一条,处理完手动 ack
$ch->basic_qos(null, 1, null);
$ch->basic_consume('order.stock', '', false, false, false, false, function ($m) use ($ch) {
// ... 业务处理
$m->ack();
});
while ($ch->is_consuming()) { $ch->wait(); }
注意:消息持久化要同时满足两个条件——`queue_declare` 的第 4 个参数 `durable=true`,且消息头 `delivery_mode=2`。只设其中一个,RabbitMQ 重启后消息照样丢。
延迟队列(比如"订单 30 分钟未支付自动取消")标准做法是:建一个不设消费者的队列,消息带 `x-message-ttl`,并配置 `x-dead-letter-exchange` 指向真正的处理队列——消息 TTL 到期后被投进死信交换机,等于延迟了 30 分钟。
第四步:Kafka——高吞吐、可重放
PHP 侧有两条路:
# 路线一:C 扩展,性能最好,需要装 librdkafka
pecl install rdkafka
# php.ini 加 extension=rdkafka.so
# 路线二:纯 PHP 客户端,免装扩展
composer require longlang/phpkafka
rdkafka 的生产者大致长这样:
$conf = new RdKafka\Conf();
$conf->set('bootstrap.servers', '127.0.0.1:9092');
$producer = new RdKafka\Producer($conf);
$topic = $producer->newTopic('user-log');
$topic->produce(RD_KAFKA_PARTITION_UA, 0, json_encode($event));
$producer->flush(3000);
// 消费者走消费者组,分组内同一分区只被一个成员消费
$conf->set('group.id', 'log-writer');
$consumer = new RdKafka\KafkaConsumer($conf);
$consumer->subscribe(['user-log']);
while (true) {
$msg = $consumer->consume(1000);
if ($msg->err) continue;
// ... 处理,然后 $consumer->commit($msg) 提交 offset
}
Kafka 的关键差异是:消息消费后不会删除,靠 offset 记录位置,所以可以重放历史数据、可以多个消费者组各读一遍。这是 Redis 和 RabbitMQ 做不到的。
第五步:三种方案怎么选
| 维度 | Redis | RabbitMQ | Kafka |
|---|
| 部署成本 | 极低 | 中 | 高 |
| 单机吞吐 | 万级 | 万级 | 十万级+ |
| 消息可靠性 | Stream 可用 | 强 | 强 |
| 延迟/定时任务 | 需自己实现 | 原生 DLX | 无 |
| 消息重放 | 不支持 | 不支持 | 支持 |
一句话:任务异步用 Redis,业务路由用 RabbitMQ,数据管道用 Kafka。
第六步:让消费者常驻跑起来
PHP 是请求结束就回收的语言,消费者必须用命令行脚本常驻,再用 Supervisor 守护:
php /www/app/worker.php
Supervisor 配置要点:
[program:mq-worker]
command=php /www/app/worker.php
autostart=true
autorestart=true
numprocs=2
stopwaitsecs=30
注意:PHP 消费者长时间运行容易内存增长(循环里累积的变量、日志上下文、ORM 的连接缓存)。做法有二——在循环内主动 `unset()` 大变量、关掉不必要的日志累积;或干脆在脚本里判断进程运行超过一定时长(比如 300 秒)就 `exit`,让 Supervisor 自动拉起新进程。
小结
- 选型看场景:异步任务 Redis、业务路由 RabbitMQ、数据管道 Kafka,别过度设计。
- Redis 用 `BRPOP` 会丢消息,生产环境用 `BRPOPLPUSH` 备份队列或 Redis Stream + ACK。
- RabbitMQ 持久化必须 queue durable 和 delivery_mode=2 同时设置,否则重启丢消息。
- 延迟任务用 RabbitMQ 的 TTL + 死信交换机,不要用轮询数据库实现。
- PHP 消费者必须命令行常驻 + Supervisor 守护,并设好自动重启和超时时间,防内存泄漏。