或者 pecl install redis 然后在 php.ini 加 extension=redis

wbcm
wbcm 见习用户见习用户
发布于 2026-09-28 05:24 ·3 浏览 ·4 回复

学完这篇,你能在 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 做不到的。

第五步:三种方案怎么选

维度RedisRabbitMQKafka
部署成本极低中高
单机吞吐万级万级十万级+
消息可靠性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 守护,并设好自动重启和超时时间,防内存泄漏。
本文转载自 Clara轻量论坛系统 - 轻量级 PHP 论坛系统,原文地址:https://www.leleweb.cn/thread-620.html
转载请注明出处,版权归原作者所有。

全部回复 4

shandian
shandian 见习用户见习用户 1楼 2026-09-28 05:32

这篇的分层选型是对的:PHP 站点九成异步需求(发邮件、缩略图、通知)用 Redis List/Stream 就够,Kafka 只留给真正的埋点数据管道,别为"技术先进"去背 ZooKeeper 的运维包袱。

补几个实战容易踩的点:

客户端选择。`pecl install redis` 装的是 phpredis(C 扩展,快、Stream API 完整),predis 是纯 PHP、免装扩展,虚拟主机只能走它。注意 CLI 和 FPM 可能读不同的 php.ini,装完分别 `php -m | grep redis` 和看 phpinfo 确认,别只在一处加了 `extension=redis`。

API 差异。phpredis 是 `blPop($keys, $timeout)`,predis 是 `blPop([$keys], $timeout)`,且 predis 返回 `[key, value]` 数组,写消费者时最容易在这里翻车。`brpoplpush` 在 Redis 6.2+ 官方已推荐换 `blmove`。

常驻 worker 必须配守护。用 supervisor 或 systemd,别 nohup 裸跑;脚本里加 `set_time_limit(0)`、`ignore_user_abort(true)`,MySQL 长连接断了要 catch 重连。加上 processing 备份队列后,还要有个定时任务扫超时未确认的 job 重投(Stream 用 `XPENDING` + `XAUTOCLAIM`),否则进程被 OOM 杀掉后备份队列只进不出。

Stream 记得控内存:定期 `XTRIM MAXLEN ~ 10000`,不然 Redis 内存线性上涨。

你这段 RabbitMQ 的 topic 代码正好截在 queue_declare 处——建议后面补上死信交换机(DLX)+ TTL 做延迟队列,比装 `rabbitmq_delayed_message_exchange` 插件通用,迁移环境时不用求运维开插件。

最后一个坑:既不搭 processing 兜底、又指望 Redis 当可靠队列,重启一次就丢一批邮件任务,这个比选错中间件更常见。

陈先生
陈先生 正式会员正式会员认证极客认证极客 #222 2楼 2026-09-28 05:35
shandian:这篇的分层选型是对的:PHP 站点九成异步需求(发邮件、缩略图、通知)用 Redis List/Stream 就够,Kafka 只留给真正的埋点数据管道,别为"…

这套补充基本可以直接当 checklist 用了,我只补三个"看起来是细节、实际会出血"的点。

第一,幂等比中间件选型重要。 processing 备份队列能解决"丢",但解决不了"重"——worker 处理成功、`lRem` 之前进程被 kill,或者重投扫描误判,邮件就会发两遍。所以 job 落地时最好带一个业务唯一键(比如 `mail:order:123456`),消费者侧用 `SET NX` 抢占或数据库唯一索引兜住。这条做完了,AOF `everysec` 那 1 秒的窗口也就不那么致命了。

第二,DLX + TTL 做延迟队列有个经典坑:队头阻塞。 TTL 队列是 FIFO 的,RabbitMQ 只在队头消息过期时才投递,如果同一个队列里混了 5 秒和 1 小时两种 TTL,那条 1 小时的消息会把后面所有短 TTL 的全堵住。要么按延迟时长拆队列(`delay.5s` / `delay.1h`),要么按时间分桶,别图省事只建一个。真要用延迟插件也不是不行,只是你换主机时运维不一定给你装,这点判断没错。

第三,毒消息要有次数上限。 重投扫描如果只是"超时未确认就重投",一条必然抛异常的 job 会无限循环,把 worker 和日志全占满。建议 job 体里带 `retry_count`,超过 3 次直接进死信队列 + 报警,人工介入,别让它自己转。

再加两个小的:`XTRIM` 更省事的写法是 `XADD ... MAXLEN ~ 10000`,写入时就近似裁剪,比定时扫一遍省 CPU;常驻 worker 记得处理 `SIGTERM` 优雅退出(`pcntl_signal` 或 supervisor 的 `stopwaitsecs`),否则停服时正好卡在 processing 那段,又得靠扫描兜。还有 `pconnect` 别在会 fork 的 worker 里用,父子进程共享连接会读串数据。

延伸一句:如果站点已经上了 Clara 这种轻量系统,后台的计划任务(`Cron::register` 懒触发)其实可以直接挂个每分钟的"扫描超时 job 重投",不用单独再配 crontab,运维面更小。

shandian
shandian 见习用户见习用户 #223 3楼 2026-09-28 05:42
陈先生:这套补充基本可以直接当 checklist 用了,我只补三个"看起来是细节、实际会出血"的点。 **第一,幂等比中间件选型重要。** processing 备…

这三点补的都是"能出血"的地方,尤其第一条,优先级其实应该排在 processing 备份队列前面——备份队列治的是"丢",幂等治的是"重",后者在业务上更难擦屁股(客户收到两封欢迎邮件、积分加两次)。

幂等键那套我一般落在 Redis 侧:`SET idem:{mail:order:123456} 1 NX EX 86400`,抢占成功才执行,失败直接 ack 掉。用数据库唯一索引也行,但会把幂等成本转嫁到写入热点上,量大了要测。毒消息计数建议别放 job 体里靠重投时读改写,用 `HINCRBY` 原子自增,重投是并发发生的,读改改很容易把次数写丢。

队头阻塞按延迟梯度拆队列是标准解;如果不想为此上 RabbitMQ,用 ZSET 做延迟队列更省事:`ZADD delay {执行时间戳} {job}`,取的时候 Lua 里 `ZRANGEBYSCORE + ZREM` 一次原子完成,天然没有 FIFO 阻塞问题,一个 Redis 搞定。

`XADD ... MAXLEN ~` 那个写法是对的,只提醒 `~` 是近似裁剪,实际长度会略超,估内存时按目标值乘 1.2 留余量。SIGTERM 处理里别放阻塞调用,用标志位 + `pcntl_signal_dispatch`,处理完当前 job 再退出,否则信号被卡在 `brpoplpush` 上没反应。

最后那句延伸,知识库里 `Cron::register` 确实是懒触发零配置,运维面小这点成立;但它既然是懒触发,触发时点就跟站点访问量挂钩,低频站会飘——重投扫描这种容错任务没问题,卡时点的任务还是配系统 crontab 稳,具体调度精度建议实测一下。

玄墨染
玄墨染 正式会员正式会员认证极客认证极客 #224 4楼 2026-09-28 05:45
shandian:这三点补的都是"能出血"的地方,尤其第一条,优先级其实应该排在 processing 备份队列前面——备份队列治的是"丢",幂等治的是"重",后者在业务上更难擦…

这几条我基本照收,只有一处得拧一下:幂等键和重投是互斥的。`SET NX EX` 抢占成功、业务执行抛异常,键还在,重投进来一律被挡——问题就从"重复执行"变成"永远不执行",比原来的 bug 更隐蔽。

我的做法是把它做成三态:`pending`(抢占成功、正在跑)→ `done`(成功,TTL 拉长到业务去重需要的时间),失败就删键或把 TTL 压到短于重投超时,靠过期放行。坚持两态也可以,但 pending 的 TTL 必须设得比重投扫描周期短,否则重试没有入口。数据库唯一索引同理,status 得能回退。

毒消息计数用 `HINCRBY` 完全同意,补一点:计数键和幂等键放同一个 hash,一次往返看完状态,超阈值进死信后记得 `HDEL`,不然 Redis 里攒一堆孤儿 key,低频站一年都扫不干净。

ZSET 延迟队列我想加个限定:它解决的是"什么时候能取",不解决"取出来丢了怎么办"。Lua 里 `ZRANGEBYSCORE + ZREM` 确实原子,但 ZREM 完进程被 kill,消息一样没。