Redis延时队列消费端:用Swoole多进程实现高可靠订单超时关闭
简介基于Swoole构建的Redis延时消息队列多进程消费端实现面向需要处理延迟任务与高并发消息的PHP中高级开发者。项目利用Redis Sorted Set存储待处理消息以时间戳作为score结合ZRANGEBYSCORE精确取出到期任务同时借助Swoole多进程模型让子进程并行消费队列避免单进程阻塞显著提升系统吞吐量与稳定性。设计上涵盖消息入队、延时轮询、进程调度、失败重试等完整环节适用于电商订单超时、定时任务、延迟通知等常见业务场景。压缩包约1.66MB包含Serve项目核心源码、Swoole服务配置及相关说明代码结构清晰便于读者深入理解Redis与Swoole的配合方式并在此基础上进行二次开发或迁移到自有项目。目前已有1268人学习该资源适合希望掌握PHP高性能队列开发、熟悉Swoole常驻内存与异步非阻塞IO特性的工程师参考实践。1. 先把结论放前面Redis延时消息队列的消费端swoole多进程是更省心的选择做过订单超时自动关闭的人大概都经历过cron五分钟扫一次订单表量小够用量一上来就开始报警超时误差大、慢查询多。后来我们把方案迁到Redis延时消息队列上——用ZSET的score存到期时间戳到点再把任务捞出来消费端用Swoole起多进程常驻处理。这个组合没有引入额外MQ组件吞吐比单进程轮询高一个量级异常恢复也直观。下面按落地顺序讲清楚为什么这么选、最小代码怎么跑通、参数怎么调、以及踩过的几个坑。它适合有PHP技术栈、任务量在每秒几十到几千、想少维护一套中间件又能把延迟任务稳定跑起来的团队。2. 用ZSET存延时任务、用Swoole进程池消费选型逻辑与边界2.1 为什么Redis数据类型里只有ZSET适合做延时队列Redis基础数据类型里有List、Set、Hash、StringList最接近队列但List是纯FIFO做延时只能靠消费端sleep硬等Set/Hash/String更不贴。ZSET为每个member挂一个score我们约定score就是任务的到期毫秒时间戳于是“查谁到期了”变成一次范围查询再把到期member用ZREM删掉语义就是“取走并处理”。Redis是单线程执行引擎把范围查询和删除打包进一段Lua脚本后这两步从取出到删掉是原子的不会被其他worker的命令插队。和RabbitMQ延迟插件、应用内时间轮相比这个方案最大的好处是不引入额外组件。任务就摆在业务Redis里用任何redis可视化工具都能直接看积压和到期情况代价是可靠性依赖Redis持久化配置任务量大了要分桶。适用的规模大致是每秒几千到几万条入队超过之后单实例的CPU和内存会先到瓶颈。把方案拆开看就是三件事生产端写ZSET、消费端按score取任务、任务执行失败后重投标题里“消费端”是工作量最重的一块后面重点写它。方案组件成本运维难度可靠性Redis ZSET Swoole复用业务Redis低数据可查依赖持久化天然“至少一次”RabbitMQ延迟插件额外MQ集群中自带ACK与重投应用内时间轮纯代码高要自己落盘重启即失忆需额外持久化如果团队已经为Kafka消费端多线程如何保证消息顺序性问题专门维护过分区路由会更容易理解Redis这里“无共享”带来的便利对比之下ZSET方案最贴近中小团队的运维能力。2.2 消费端为什么做多进程而不是多线程PHP侧的现实选择PHP做常驻消费的第一道坎就是“脚本跑完就退出”的原生惯性。Swoole的Process\Pool正好给出固定worker数的进程池模型Manager负责fork worker、监控退出并拉起worker各自循环去Redis拉任务本质上一套现成的linux多进程通信框架。多进程模型最大的好处是无共享每个worker有独立的内存空间不需要为临界区加锁某worker异常退出Manager重新拉起一个干净的进程继续消费不至于像单进程一样整个消费端停摆。这和Kafka消费端多线程场景很不一样。Kafka多线程要面对消息顺序性同一个分区的记录落到不同线程时要么线程池内做hash路由要么引入队列按序投递Redis这边用多进程天然按进程隔离。如果确实需要顺序处理同一业务单子把jobId做哈希、固定路由到某个worker即可比如crc32($jobId) % workerNum用对worker编号的方式而不是加锁。实际项目中我一般不额外搭IPCworker之间不通信各自连Redis取任务出问题最好排查。2.3 这个方案在什么规模下会撑不住边界先想清楚在给代码之前补一句边界避免读者按文章搭完踩到天花板。第一单key积压到百万级member时ZRANGEBYSCORE虽然走跳跃表但一次性返回的成员多Redis单线程处理时间会拉长必须按业务分桶。第二“取走即删除”的语义决定了它是至少一次投递取出之后、业务执行完成之前进程崩溃任务就丢了无法接受丢失的业务要增加“处理中登记”这一点在第四章展开。第三单条任务执行耗时到秒级时worker会长时间占用batch里的任务批次必须调小否则大量worker忙在旧任务上新任务排队时间拉长。这些边界决定了后面每个参数怎么给。3. 从入队到消费跑通一个最小可复用的swoole多进程消费端3.1 入队侧与score计算毫秒时间戳才是统一单位先给出入队最小脚本生产端和消费端共用同一套命名与单位后面才能少出幺蛾子。?php $redis new Redis(); $redis-connect(127.0.0.1, 6379); $delayKey delay:queue:order_close; $payloadKey delay:queue:order_close:payload; $jobId uniqid(job_, true); // member 用业务ID不要塞大JSON $payload json_encode([ order_id 123456, op auto_close, created_at date(Y-m-d H:i:s) ]); // score 统一用毫秒时间戳30秒后到期 $score (int) (microtime(true) * 1000) 30 * 1000; $redis-zAdd($delayKey, $score, $jobId); $redis-hSet($payloadKey, $jobId, $payload);这段代码有两个约定要固化下来ZSET的member只放短jobId完整任务内容放同级Hash避免ZSET成员过长导致range结果占用大量内存score必须统一毫秒不要在这里用time()秒值、在消费端用microtime(true)*1000单位不一致会引发任务秒出或永久不触发。入队的命名我习惯用delay:queue:{业务}:payload这样的成对结构一个业务一套key不会混。3.2 消费端骨架用Swoole Process Pool创建固定数量worker最小可跑的消费端长这样?php use Swoole\Process\Pool; $pool new Pool(4); // 4个worker后续按CPU和任务耗时调 $pool-on(WorkerStart, function (Pool $pool, int $workerId) { // 关键每个worker在子进程内自己连接Redis绝不能在fork前共享连接 $redis new Redis(); if (!$redis-connect(127.0.0.1, 6379)) { fprintf(STDERR, worker %d redis connect failed\n, $workerId); return; } $delayKey delay:queue:order_close; $payloadKey delay:queue:order_close:payload; $batchSize 20; while (true) { $nowMs (int) (microtime(true) * 1000); // 单条Lua范围内取出到期任务并删除保证取和删是原子操作 $script LUA local jobs redis.call(ZRANGEBYSCORE, KEYS[1], -inf, ARGV[1], LIMIT, 0, ARGV[2]) if #jobs 0 then return {} end redis.call(ZREM, KEYS[1], unpack(jobs)) return jobs LUA; $jobs $redis-eval($script, [$delayKey, $nowMs, $batchSize], 1); if (empty($jobs)) { usleep(500000); // 没有到期任务空转500ms降CPU占用 continue; } foreach ($jobs as $jobId) { $payload $redis-hGet($payloadKey, $jobId); if (false $payload) { continue; // payload不存在宁可跳过也不能拿残缺数据去执行业务 } handleJob($jobId, json_decode($payload, true)); $redis-hDel($payloadKey, $jobId); } } }); $pool-start();逻辑说明Pool(4)在manager进程里先建好WorkerStart回调在fork出的子进程里执行代码里专门在每个worker内重新new Redis并connect这是后面踩坑章的重点。循环里先用Lua取出最多20条到期任务并删除取到空就sleep 500ms避免CPU空转。eval第3个参数传1表示KEYS数组长度是1phpredis 4.x以上支持这种写法旧版本要用rawCommand拼命令。handleJob是业务处理入口比如改订单状态、发通知它抛出的异常不应该让while循环退出否则整个worker就死了所以实际项目里要包try/catch并按第四章重投。3.3 取出即删的原子脚本为什么必须用Lua而不是两步走很多第一次写的同学会用两步先ZRANGEBYSCORE把到期jobId拉到PHP再逐条ZREM删除。问题在于两步之间有网络往返两个worker可能同时把同一条jobId看成自己的两边都ZRANGE到了两边都拿去做业务通知就发了两次。把取出和删除放进一段LuaRedis保证整段脚本执行期间不会被其他命令插入从逻辑上杜绝了“同一条任务被两个worker同时取走”。Lua里LIMIT 0后面那个20就是批次上限防止一次取出过多任务把PHP内存打爆score用ARGV传入而不是写死在脚本里脚本可以跨多个业务key复用。到这里消费端能跑但“可靠性”还不够任务取出后、handleJob前进程崩溃任务会丢下一章补重试和死信。注意不要用ZPOPMIN替代上面的Lua。ZPOPMIN取的是全局score最小的成员它不一定已经到期也没有按业务批量过滤的能力语义对延时队列不成立。3.4 用systemd把消费端守护起来常驻进程不能裸奔进程池本身不解决“整台机器重启”和“worker反复崩溃”的问题我习惯让systemd托管[Unit] DescriptionRedis Delay Queue Consumer Afternetwork-online.target redis.service [Service] Userwww-data ExecStart/usr/local/bin/php /data/app/consumer.php Restartalways RestartSec3 [Install] WantedBymulti-user.targetRestartalways表示异常退出后3秒拉起避免人工半夜补进程。值得注意的是拉起的是整个进程池新的manager会重新fork指定数量的worker之前卡在业务执行中途的任务不会“接着跑”所以重试逻辑必须放在业务侧而非依赖进程内状态。4. 消费端参数与Redis配置worker数、轮询间隔、批量取出怎么调4.1 核心参数表把四类最常改的参数列成一张表给出我的经验值参数经验取值设计依据worker_num4起步最高到CPU核数的2~4倍任务多为IO等待型可以多开但别超过Redis maxclients和业务DB连接池上限batch_size20~100任务耗时短取大值耗时长取小值避免一个worker长期占着一批任务empty_sleep500~2000ms空转sleep降CPU一旦取到任务本批次处理完前不sleepretry_delay_base1s按次数指数退避防止失败任务集中在同一秒重试把Redis打出一波尖峰max_retry3~5超限进死信人工介入不无限重试worker数的判断标准不是“核数越多越快”而是单worker的CPU占用。如果handleJob里全是Redis操作和curlCPU占用低可以往4倍开如果涉及本地大量计算超过核数会导致上下文切换。batch_size和任务耗时强相关一次直播营的过期通知任务只要几百毫秒100条没问题一个秒级执行的重型任务批次取20条都可能让其他worker饿肚子。空转sleep也别盯太死Redis里没有到期任务时消费端本质在做polling1秒一次已经足够应对绝大多数延时业务再短只是把CPU堆在空转上。4.2 Redis侧配置持久化、淘汰策略与主从选型consumer代码只解决“怎么取”Redis本身的配置更决定数据丢不丢。给出我常用的redis.conf关键项maxmemory 2gb maxmemory-policy noeviction appendonly yes appendfsync everysec tcp-keepalive 60延时队列的key绝不能被内存淘汰策略清掉所以maxmemory-policy用noeviction内存满时写入直接报错而不是默默丢任务配合监控告警才可控。appendonly yes everysec能在最坏情况下丢1秒以内的写入比默认RDB触发式快照稳得多“能不能容忍分钟级丢失”这个判断决定要不要开AOF纯RDB在Redis进程异常退出后会丢不少任务。tcp-keepalive是为了让长时间空闲的worker连接被Redis及时发现避免服务端半开连接堆积。主从与哨兵方面哨兵模式和集群模式有区别哨兵场景下读写都走主节点消费端通过哨兵客户端拿到当前主节点地址即可Cluster场景下同一段Lua脚本涉及多个key时key必须落在同一slot所以按业务分桶的key要写成{order_close}:delay和{order_close}:payload这种hash tag形式否则直接报CROSSSLOT错误。切换主从后worker里已建立的连接会断开下一轮循环hGet报错捕获后重连即可不要试图在循环外做长连接复用。注意Cluster环境下涉及多个key的Lua必须保证key在同一slot命名上使用hash tag是常见解法。4.3 失败重试与死信消费端不只要会取任务还要会还任务消费端最容易被忽略的是异常路径handleJob抛异常时任务已经离开原ZSET直接丢弃会导致订单永远没人关。我的做法是重投递到延时队列按重试次数给它新的scoretry { handleJob($jobId, json_decode($payload, true)); $redis-hDel($payloadKey, $jobId); } catch (Throwable $e) { $retryCnt (int) $redis-hIncrBy(stat:job_retry:order_close, $jobId, 1); if ($retryCnt 5) { $deadKey dead:queue:order_close: . date(Ymd); $redis-zAdd($deadKey, (int)(microtime(true) * 1000), $jobId); error_log(sprintf(job dead %s reason%s payload%s, $jobId, $e-getMessage(), $payload)); $redis-hDel($payloadKey, $jobId); // 死信可以落Redis给人看payload从工作hash移除 } else { $delayMs 1000 * (1 $retryCnt); // 1s、2s、4s、8s…指数退避 $redis-zAdd($delayKey, (int)(microtime(true) * 1000) $delayMs, $jobId); } }这里的要点是重投时覆盖原ZSET的score原score是第一次的到期时间不覆盖的话下轮循环立即再次取到重试就失去了退避意义。重试次数记在Hash里以jobId为字段不污染ZSET成员超过max_retry后进入以天为粒度的死信key方便第二天看板。payload从工作hash移除要在死信写入之后避免死信只有jobId没有内容。重试机制和幂等是配合关系网络异常、Redis主从切换都可能导致同一条任务出现在两个阶段业务侧要有唯一键兜底否则重试越坚定重复生产事故越严重。5. 多进程消费端的5个踩坑记录现象、原因、解决方案连踩坑都整理出来了照着检查一遍能省不少线上时间。下面五条按现象、原因、解决展开。5.1 先取再删导致两个worker重复消费现象订单关闭任务上线后同一笔订单的关闭逻辑被跑了两次业务表多出两条重复的变更记录。原因最初的消费端先ZRANGEBYSCORE取出到期任务再发ZREM删除两步之间存在网络往返两个worker可能同时ZRANGE到同一条jobId各自都认为自己取到了。解决按3.2改用Lua脚本把ZRANGEBYSCORE和ZREM封装成一步从源头消除窗口。另外提醒一句不要试图再用Redis分布式锁去锁jobId锁的过期时间和持有者崩溃带来的新窗口会让问题更复杂业务表做唯一键、做状态机幂等才是最后的后悔药。5.2 fork前建立的redis连接被进程池共用后报错现象消费端刚跑时正常几分钟后日志里刷“Connection closed”和协议错乱数据张冠李戴。原因连接对象是在Pool-start()之前new出来的fork把父进程的TCP socket文件描述符复制给了每个子进程多个worker共用同一个连接发的命令互相抢读响应。解决Redis连接的建立必须放在WorkerStart回调里面每个worker独立connect同理Worker内不能继承父进程的其他长连接资源比如MySQL连接、AMQP连接凡是进程池里要用的都在WorkerStart里重建。5.3 score单位不统一任务要么秒出要么卡死现象一部分任务入队后立刻被消费另一部分任务过了几个小时都没动静。原因入队端用time()返回的秒数score是3672xxx消费端用microtime(true)*1000score是1672xxxxxxx差了一千倍范围查询把“未到期”判成“已到期”或相反。解决全链路约定score一律毫秒生产端和消费端共用一个常量函数生成上线自检时入队一条延迟1秒的任务看它是否约1秒后被消费把这个冒烟检查写进部署脚本。这个坑看起来小实测里是最容易集体翻车的一类。5.4 百万成员的ZSET大key把Redis阻塞到可视化工具打不开现象某个促销活动的任务全塞进同一个ZSET积压到百万级Redis操作耗时出现秒级毛刺用redis可视化工具查这个key时界面转半天。原因单key过大ZRANGEBYSCORE返回大量member序列化传输和内存分配都压在单线程Redis上大key的rehash和内存碎片进一步放大延迟。解决任务按业务域分桶比如order_close、order_notify、point_expire各自独立key同一个业务内如果积压规模继续增长再按天分桶比如delay:queue:order_close:20250112消费端按当天桶依次扫描。分桶之后每个key的体量可控range和ZREM都不会成为瓶颈。5.5 Redis重启后任务全丢持久化策略没跟上现象Redis实例因机器重启而重启恢复后延迟队列里空空如也大量订单没被自动关闭。原因默认配置下Redis只做RDB触发式快照两次快照之间的数据只存在内存里进程异常退出就丢了消费端的AOF/持久化策略没有同步补齐。解决按4.2开AOF everysec并设置noeviction把丢失窗口压到秒级同时给“已取出正在处理”的任务做一个DB登记取出时插一条处理中记录处理完删除Redis恢复后扫DB把仍处于处理中的任务重新投递。这是把延时队列从一个“可能会丢任务的Redis脚本”升级成“可恢复的服务”的关键一步。6. 用压测数据验证消费端以及我保留的三个调参习惯6.1 用批量入队脚本验证端到端延迟与重复率压测不要只看消费速度。我常用下面这个脚本入队5万条任务延迟统一设为1秒然后看两件事端到端延迟和重复率。#!/bin/bash KEYdelay:queue:order_close for i in $(seq 1 50000); do NOW_MS$(date %s%3N) SCORE$((NOW_MS 1000)) redis-cli ZADD $KEY $SCORE job_$i /dev/null done echo enqueued 50000 jobs验证时在handleJob里写一条日志记录jobId、预期到期时间和实际消费时间统计P99延迟重复率在消费端对jobId做SADD去重SADD返回0就说明这条任务不是第一次被处理。四个验收指标可以参考下表验收项目标怎么查端到端延迟P99设定延时500ms以内业务日志时间差重复消费率0%消费时SADD jobId判断返回值积压趋势ZCARD单调下降redis-cli --stat轮询重启恢复任务不丢或按死信重投手动kill消费端后观察恢复6.2 我留在部署脚本里的三个自检习惯第一每次发版前自动入队一条延迟1秒任务验证score单位没被改坏第二把dead key的ZCARD和业务错误率接到告警死信只增不减时立刻停止压测看代码第三新任务类型先放10%流量跑一周确认积压和延迟P99达标后再全量。把消费端当成一个独立常驻服务来运维而不是“跑一个php脚本”你会少走很多弯路。希望帮到你。本文还有配套的精品资源点击获取
上一篇/下一篇内容由系统自动关联
返回资讯列表 →