简介这份资源面向PHP后端开发者与消息队列学习者提供一套基于Redis延时队列与Swoole多进程模型构建的高并发消费端实现可用于订单超时关闭、定时任务触发等需要延迟处理的业务场景。压缩包为zip格式大小约1.66MB内含项目源码及相关资源文件围绕Redis的Sorted Set延时队列设计、Swoole Server多进程消费配置与消息处理逻辑展开适合具备一定PHP与Swoole基础的中高级开发者参考。目前已有1269人学习下载。通过阅读源码读者可以理解如何用时间戳作为score将消息写入Sorted Set再借助ZRANGEBYSCORE按到期时间取出并消费同时掌握Swoole多子进程并行消费、非阻塞I/O降低资源消耗的具体写法并借鉴其目录组织与消费服务拆分思路快速搭建可扩展的延时消息处理系统。1. 从一次订单超时未关闭说起这套 Redis 延时队列到底解决了什么去年帮一个做社区团购的团队排查问题凌晨两点接到电话一批订单在支付超时后没有自动关闭库存被锁死第二天早高峰直接爆单失败。翻代码发现他们用的是crontab每分钟扫一次 MySQL订单量一上来扫描周期内积压几千条延迟从 30 分钟漂到 40 多分钟。这不是个例很多 PHP 项目在「延时任务」这件事上都踩过同一个坑——用数据库轮询硬扛量小的时候岁月静好量一大就翻车。这套「Redis 延时消息队列 Swoole 多进程消费端」要解决的就是这个场景订单超时关闭、优惠券到期提醒、拼团倒计时成团判定、异步通知重试。核心思路是把延时任务塞进 Redis 的有序集合ZSetscore 存到期时间戳消费端用 Swoole 起多个 worker 进程并发拉取到期任务。相比轮询数据库它把「什么时候该执行」这件事交给 Redis 的内存排序消费端只关心「现在有没有到期的」单机轻松扛住每秒几千次的延时投递。适合谁手上是 PHP 技术栈、已经在用 Redis 做缓存或中间件、又不愿意为了一个延时功能去引 Kafka 或 RabbitMQ 的团队。如果你正好卡在「crontab 不够准、数据库扛不住」的中间地带这套东西值得拆开看。2. 延时队列的数据结构选型为什么是 ZSet 而不是 List2.1 ZSet 的 score 排序与到期判定逻辑延时队列的本质是一个「按时间排序的待办清单」。Redis 里能排序的结构不少但真正适合的只有有序集合 ZSet。它的每个成员带一个 scoreRedis 内部用跳表维护 score 的顺序ZRANGEBYSCORE可以按分数区间取数据复杂度是 O(log N M)。把任务的到期时间戳毫秒或秒当 score任务体当 member取「所有 score 小于等于当前时间」的成员就是所有该执行的任务。这里有个关键点member 必须唯一。如果两个任务内容完全一样ZSet 会去重第二个任务直接覆盖第一个。常见做法是把任务 ID 或 UUID 拼进 member比如order_close:123456保证唯一性。任务体本身比如订单号、业务类型、重试次数序列化成 JSON 存进去消费端取出来再反序列化。# 投递一个 30 分钟后执行的订单关闭任务 # score 当前时间戳 1800 秒 ZADD delay_queue 1735660800 {task_id:order_close:123456,order_id:123456,type:close,retry:0} # 消费端拉取所有已到期的任务score 当前时间戳 ZRANGEBYSCORE delay_queue -inf 1735659000 LIMIT 0 100上面第一条命令的 score 是绝对时间戳不是「延迟多少秒」。很多新手会直接把 1800 当 score 存进去结果任务永远不到期——因为 1800 对应的是 1970 年。投递时必须用time() delay算出绝对到期时间。第二条命令的-inf表示从最小分数开始LIMIT 0 100控制单次拉取数量避免一次拉太多把内存打爆。2.2 为什么不用 List 和 Stream 做延时List 是队列先进先出但它没有「按时间排序」的能力。你没法问 List「哪些元素到期了」只能一个个LPOP出来判断没到期还得塞回去逻辑上就拧巴。有人用两个 List 来回倒腾一个待处理、一个延时中本质是在应用层模拟 ZSet多进程并发时还得加锁得不偿失。Stream 是 Redis 5.0 引入的消息队列结构支持消费者组和 ACK 机制做即时消息队列很合适。但它同样没有延时投递的原生能力XADD进去的消息立刻可被消费想延时还是得靠外部调度。所以「延时」这件事ZSet 是 Redis 原生结构里最顺手的没有之一。选型上还有一层考虑ZSet 的ZRANGEBYSCORE配合ZREM可以做到「取出即删除」但这两步不是原子的。多进程并发时两个 worker 可能同时拉到同一条任务。解决办法有两个一是用 Lua 脚本把「查删」包成原子操作二是用ZPOPMIN这类原子命令。Lua 脚本更灵活下面这段就是常见的原子拉取逻辑。-- 原子地从 ZSet 中取出并删除已到期的任务 -- KEYS[1] 队列名, ARGV[1] 当前时间戳, ARGV[2] 单次拉取上限 local tasks redis.call(ZRANGEBYSCORE, KEYS[1], -inf, ARGV[1], LIMIT, 0, ARGV[2]) if #tasks 0 then redis.call(ZREM, KEYS[1], unpack(tasks)) end return tasks这段脚本的ARGV[1]是消费端传进来的当前时间戳ARGV[2]是批量大小。ZRANGEBYSCORE取出到期任务后立刻ZREM删除整个过程在 Redis 单线程里执行不会被其他 worker 插队。注意unpack在任务数量很大时可能触发 Lua 栈限制一般批量控制在 100 以内没问题。如果任务体很大建议 member 只存任务 ID任务详情放 Hash 或独立 key避免 ZSet 成员过大拖慢操作。3. Swoole 多进程消费端的搭建从进程模型到任务分发3.1 Swoole 进程模型与 worker 数量配置Swoole 的多进程模型比传统 PHP-FPM 更适合常驻消费端。用Swoole\Process手动 fork 出多个 worker每个 worker 独立循环拉取任务互不阻塞。也可以用Swoole\Server的 task worker但那个偏向「请求-任务」模式纯消费场景用Process更直接。worker 数量怎么定不是越多越好。每个 worker 都在跑ZRANGEBYSCORERedis 是单线程处理命令的worker 太多反而让 Redis 的 QPS 成为瓶颈。经验值是 CPU 核数的 2 到 4 倍比如 4 核机器开 8 到 16 个 worker。如果任务里有大量 IO 等待比如调第三方接口可以适当多开如果任务是纯 CPU 计算开太多只会互相抢核。?php // consumer.php - Swoole 多进程消费端骨架 $workerNum 8; // 根据 CPU 核数和任务类型调整 $queueName delay_queue; $redisHost 127.0.0.1; $redisPort 6379; $processes []; for ($i 0; $i $workerNum; $i) { $process new Swoole\Process(function (Swoole\Process $worker) use ($queueName, $redisHost, $redisPort, $i) { $redis new Redis(); $redis-connect($redisHost, $redisPort); $redis-setOption(Redis::OPT_READ_TIMEOUT, -1); // 常驻连接不超时 // 注册 Lua 脚本避免每次传输脚本内容 $script file_get_contents(__DIR__ . /pop_expired.lua); $sha $redis-script(load, $script); while (true) { $now time(); // evalsha 执行原子拉取批量 50 条 $tasks $redis-evalSha($sha, [$queueName, $now, 50], 1); if (empty($tasks)) { usleep(100000); // 没任务时睡 100ms降低 Redis 压力 continue; } foreach ($tasks as $taskJson) { $task json_decode($taskJson, true); // 这里分发到具体业务处理 handleTask($task, $redis); } } }, false, false); // 不启用管道独立进程 $pid $process-start(); $processes[$pid] $process; echo worker {$i} started, pid{$pid}\n; } // 主进程等待回收子进程 foreach ($processes as $process) { Swoole\Process::wait(); }这段代码里几个参数值得说清楚。usleep(100000)是 100 毫秒空轮询时让出 CPU不然 8 个 worker 会把 Redis 的 QPS 打满。evalSha比eval省带宽脚本内容只传一次后续只传 SHA1。OPT_READ_TIMEOUT设为 -1 是让 Redis 连接不因读超时断开常驻进程必须设否则跑一段时间就报RedisException: read error on connection。handleTask是业务处理入口下面单独说。3.2 任务分发与失败重试的落点任务从 ZSet 取出来后不能直接删了就完事。如果业务处理失败任务就丢了这是延时队列最容易被忽略的坑。常见做法是「取出后先不删处理成功再删」但这样又回到并发重复消费的问题。折中方案是Lua 脚本取出后立刻删除同时把任务写入一个「处理中」的 List 或 Hash处理成功后从「处理中」移除如果 worker 崩溃有个兜底进程扫描「处理中」超时的任务重新投递。?php function handleTask(array $task, Redis $redis): void { $processingKey delay_queue:processing; // 取出时已经 ZREM这里记录到处理中集合带时间戳 $redis-hSet($processingKey, $task[task_id], json_encode([ task $task, start_at time(), ])); try { // 按业务类型分发 switch ($task[type]) { case close: closeOrder($task[order_id]); break; case notify: sendNotify($task[order_id]); break; default: throw new Exception(unknown task type: {$task[type]}); } // 处理成功从处理中移除 $redis-hDel($processingKey, $task[task_id]); } catch (Throwable $e) { // 失败重试次数 1延迟 60 秒重新投递 $task[retry] ($task[retry] ?? 0) 1; if ($task[retry] 3) { $redis-zAdd(delay_queue, time() 60, json_encode($task)); } else { // 超过重试上限进死信队列人工介入 $redis-lPush(delay_queue:dead, json_encode($task)); } $redis-hDel($processingKey, $task[task_id]); echo task {$task[task_id]} failed: {$e-getMessage()}\n; } }重试策略里retry字段是关键每次失败加一超过 3 次进死信队列。死信队列用 List 存方便人工捞出来排查。注意重试时重新zAdd的 score 是time() 60也就是延迟 60 秒再试避免失败任务立刻又被拉起来形成死循环。processingKey用 Hash 存field 是 task_idvalue 是任务和处理开始时间兜底进程可以扫这个 Hash 里start_at超过 5 分钟的记录说明 worker 可能挂了重新投递。4. 避坑与排查延时队列上线后最容易翻车的五个点4.1 任务重复消费现象是同一订单被关闭两次现象日志里同一个order_id出现两次关闭操作第二次报「订单状态不允许关闭」。原因通常是 Lua 脚本没用好或者消费端在ZREM之前就处理了任务处理过程中另一个 worker 又拉到了同一条。解决确保「查删」在同一个 Lua 脚本里原子执行处理逻辑放在ZREM之后。如果业务本身不幂等在handleTask开头加一层状态判断比如查一下订单当前状态已关闭的直接跳过。4.2 Redis 连接超时报错 read error on connection现象消费端跑几小时后报RedisException: read error on connection进程退出。原因是 PHP Redis 扩展默认读超时较短常驻连接空闲一段时间后被服务端或中间网络断开。解决$redis-setOption(Redis::OPT_READ_TIMEOUT, -1)关掉读超时同时在循环里加try-catch捕获异常后重连 Redis不要让整个 worker 挂掉。重连逻辑要带退避比如第一次等 1 秒第二次 2 秒避免 Redis 刚重启就被打满。4.3 时间戳精度不一致任务延迟忽大忽小现象投递时明明设了 30 分钟实际执行有时 29 分有时 31 分。原因是投递端和消费端用的时间源不一致或者投递时用了秒级时间戳消费端用毫秒级比较。解决全链路统一用秒级或毫秒级投递和消费用同一台机器的时间或 NTP 同步过的时间。如果对精度要求高score 用毫秒时间戳消费端microtime(true) * 1000取当前毫秒。另外usleep(100000)意味着最坏情况下任务会晚 100 毫秒执行这个精度对绝大多数业务够用但如果你要毫秒级准时得把 sleep 调小或改用事件驱动。4.4 worker 数量过多导致 Redis QPS 打满现象加到 32 个 worker 后Redis 监控显示 QPS 飙升其他业务读写变慢。原因是每个 worker 都在高频执行ZRANGEBYSCORE空轮询时尤其明显。解决worker 数量控制在 CPU 核数的 2 到 4 倍空轮询 sleep 时间适当加大比如 200 毫秒或者用BLPOP配合一个「有新任务」的通知 List让 worker 在没任务时阻塞等待而不是空转。通知机制是投递任务时同时LPUSH一个信号到通知 Listworker 用BLPOP阻塞读读到信号再去 ZSet 拉任务。4.5 任务体过大导致 ZSet 内存膨胀现象Redis 内存持续增长ZRANGEBYSCORE变慢。原因是任务体直接序列化进 ZSet member一个任务几 KB几十万任务就是几百 MB。解决ZSet 的 member 只存任务 ID任务详情放独立的 Hash 或 String key消费端拿到 ID 再去取详情。这样 ZSet 本身很小排序和范围查询都快。任务详情 key 设过期时间比如 7 天避免历史任务占内存。如果任务详情也要持久化落库比放 Redis 更合适。5. 进阶技巧用 Lua 脚本做批量拉取与消费端优雅退出5.1 批量拉取脚本的边界处理前面给的 Lua 脚本是最简版实际用的时候有几个边界要处理。第一unpack在任务数量大时会报「too many results to unpack」Lua 5.1 的栈上限大概是 8000 左右所以批量大小别超过 1000一般 50 到 100 足够。第二如果ZRANGEBYSCORE返回空ZREM不要执行虽然执行了也没事但省一次命令调用。第三脚本里可以加一个「最大拉取时间」保护避免单次执行太久阻塞 Redis。-- 增强版带空结果判断和批量上限保护 local queue KEYS[1] local now ARGV[1] local limit tonumber(ARGV[2]) if limit 200 then limit 200 end -- 硬上限防止误传大值 local tasks redis.call(ZRANGEBYSCORE, queue, -inf, now, LIMIT, 0, limit) if #tasks 0 then return {} end redis.call(ZREM, queue, unpack(tasks)) return tasks这个版本把批量上限硬编码为 200即使调用方传了 1000 也会被截断。tonumber转换是必要的因为 Redis 传进来的 ARGV 都是字符串。返回空表时消费端empty($tasks)判断为真直接 sleep不会有多余的ZREM调用。5.2 消费端优雅退出与信号处理常驻进程最怕的是「kill -9 直接杀」任务处理到一半被中断处理中集合里留下脏数据。Swoole 的Process支持信号监听收到SIGTERM时设置一个退出标志当前循环处理完再退出。?php // 在 worker 闭包开头注册信号处理 Swoole\Process::signal(SIGTERM, function () use ($worker) { echo worker {$worker-pid} received SIGTERM, exiting after current task\n; $GLOBALS[running] false; }); while ($GLOBALS[running] ?? true) { // ... 拉取和处理逻辑 } // 退出前清理把处理中的任务重新投递 cleanupProcessingTasks($redis);SIGTERM是kill默认发的信号SIGKILLkill -9捕获不了所以部署脚本里要用kill而不是kill -9。退出前调cleanupProcessingTasks把当前 worker 处理中的任务重新zAdd回队列score 设为time() 10让其他 worker 很快能接手。这个清理逻辑也可以做成独立命令手动触发。5.3 监控指标与验证方法上线后怎么知道队列健康盯三个指标队列长度ZCARD delay_queue、最老任务的到期时间ZRANGE delay_queue 0 0 WITHSCORES、处理中集合大小HLEN delay_queue:processing。队列长度持续增长说明消费速度跟不上投递速度要么加 worker 要么优化任务处理逻辑。最老任务的 score 如果远小于当前时间说明有任务积压没被及时消费。处理中集合大小如果只增不减说明有 worker 卡死或崩溃兜底进程该上场了。验证延时精度可以写个测试脚本投递 100 个任务延迟分别设 1 到 100 秒记录实际执行时间算偏差。正常情况下偏差应该在 sleep 间隔以内比如 100 毫秒。如果偏差很大检查是不是 worker 太少导致任务排队或者 Redis 负载太高导致命令响应慢。从那以后我每次上线延时队列都强制先跑一遍「投递-消费-重试-死信」的完整链路测试确认每个环节的日志都能对上才敢放量。这套东西不复杂但细节多尤其是原子性和重试这两块偷懒少写一行 Lua后面就得花一晚上排查重复消费。希望帮到你。本文还有配套的精品资源点击获取