公司动态

think‑queue Redis 消息队列完全指南

📅 2026/8/25 23:26:37
think‑queue Redis 消息队列完全指南
核心概念 配置参数 执行原理 实战案例 生产部署一、概述与核心概念think‑queue 是 ThinkPHP 官方提供的消息队列扩展包用于将耗时任务(如邮件发送、短信通知、文件处理、日志写入、数据同步等)从 Web 请求中剥离出来异步执行从而缩短接口响应时间并在流量高峰时起到削峰填谷的作用。其架构遵循生产者—消费者模型业务代码(生产者)把任务投递到队列后台常驻进程(消费者)循环从队列中取出任务并执行对应的 Job 类。think‑queue 具备以下核心能力多种驱动官方提供 Redis、Database、Sync(同步执行多用于测试)、Topthink 等连接器其中 Redis 驱动性能最好、功能最完整是生产环境的首选。延迟任务支持 later() 投递延迟任务可实现“30 分钟后未支付自动取消订单”等场景。失败重试任务执行失败可重新放回队列支持限制最大重试次数并可监听失败事件。超时保护任务被消费者取出后若超时未确认会自动重新入队避免任务丢失。多队列与优先级支持按业务划分多个队列并可对队列设置消费权重。版本对应关系think‑queue 3.x 适配 ThinkPHP 6/8think‑queue 2.x 适配 ThinkPHP 5.15.0 使用 thinkphp‑queue 1.x。本文以 3.x Redis 驱动为准其余版本用法基本一致。二、安装与配置2.1 安装扩展包composer require topthink/think-queue2.2 配置文件 config/queue.php安装后若框架未自动生成配置文件可手动创建 app 目录同级 config/queue.php?php return [ //默认连接器(驱动),可选redis/database/sync defaultredis, connections[ sync[ typesync, ], database[ typedatabase, queuedefault, tablejobs, ], redis[ typeredis, queuedefault,//默认队列名 host127.0.0.1,//Redis主机 port6379,//Redis端口 password,//Redis密码,无密码留空 select0,//Redis库编号 timeout0,//连接超时时间(秒) expire60,//任务执行超时时间(秒) ], ], //失败任务的记录连接(可选) failed[ typenone, tablefailed_jobs, ], ];2.3 Redis 配置参数详解参数说明type连接器类型填写驱动名或完整的连接器类名(如 \think\queue\connector\Redis)。参数默认值说明queuedefault默认队列名称。投递任务时未指定队列则使用该队列Redis 中的 key 形如 queues:队列名。host127.0.0.1Redis 服务器地址可填 IP 或域名。port6379Redis 服务器端口。passwordRedis 认证密码未设置 requirepass 时留空字符串。select0Redis 数据库编号(0~15)建议为队列业务单独分配一个库。timeout0建立 Redis 连接的超时时间(秒)0 表示不限制。expire60任务执行超时时间(秒)。消费者取出任务后若超过该时间仍未删除任务任务会被判定为“执行超时”并重新放回队列可能被其他消费者再次执行。设为 -1 表示永不超时。expire 的取值必须大于业务代码中单个任务的最长执行时间否则任务会在尚未执行完时被重新入队造成同一任务被多个进程并发执行。expire 也不会取代命令行的 --timeout 参数两者需要配合使用。三、消费命令与参数详解think‑queue 提供两个内置命令用于消费队列queue:listen与queue:work。#监听模式(推荐)常驻进程循环取任务框架每次执行完任务后重新加载可配合--once调试 php think queue:listen --queue email #工作模式取到任务后由子进程执行主进程负责调度 php think queue:work --queue email --tries33.1 queue:listen 常用参数参数默认值说明--queuenull要监听的队列名多个队列用逗号分隔可为队列指定权重(优先级)如--queueemail,order:2表示 order 队列被消费的优先级更高。--oncefalse仅处理一个任务后退出通常用于调试。--delay0任务执行失败后重新入队的延迟时间(秒)。--memory128进程内存上限(MB)超过后进程退出(便于配合进程管理器重启释放内存)。--timeout60单个任务执行超时时间(秒)超时后进程被杀掉任务因未确认而重新入队。--sleep3队列空闲时两次轮询之间的休眠时间(秒)降低空转对 CPU 的占用。--tries0任务最大尝试次数。超过次数后任务进入失败流程。需在 Job 中配合 attempts() 使用。3.2 queue:work 常用参数参数默认值说明--queuenull同 queue:listen指定消费的队列及权重。--oncefalse仅处理一个任务后退出。--daemonfalse以守护模式运行进程常驻框架只启动一次。性能高但代码更新后必须重启进程才能生效。--delay0任务失败后重新入队的延迟秒数。--memory128内存上限(MB)超出后退出。--timeout60单个任务执行超时时间(秒)。--sleep3空闲时轮询间隔(秒)。--tries0最大尝试次数。listen 与 work 的选择建议queue:listen每次处理任务前都会重新引导框架代码更新后自动生效适合频繁发版的环境但吞吐量较低。queue:work --daemon进程常驻内存、性能更高适合生产环境缺点是发版后必须手动或通过进程管理器重启。两种模式下都应通过 Supervisor 等工具守护进程避免进程意外退出后队列停止消费。四、核心 API 说明4.1 投递任务(生产者)use think\facade\Queue; //立即执行投递到默认队列 Queue::push(app\job\SendEmail::class,$data); //立即执行投递到指定队列 Queue::push(app\job\SendEmail::class,$data,email); //延迟执行30秒后执行 Queue::later(30,app\job\OrderCancel::class,$data,order); //延迟执行指定具体时间戳 Queue::later(time()1800,app\job\OrderCancel::class,$data,order);参数说明第一个参数Job 类名(推荐写完整类名常量::class)也可以是类实例对象。第二个参数$data传给任务的业务数据只支持标量或数组(会被序列化为 JSON 存入 Redis)不要传入对象、闭包或资源句柄。第三个参数队列名称省略时使用配置中的默认队列。later() 第一个参数延迟秒数(相对时间)或绝对时间戳。4.2 Job 任务类(消费者)Job 类中通过 fire() 方法执行业务逻辑必须正确处理任务生命周期?php namespace app\job; use think\queue\Job; class SendEmail { /** * fire是队列调用任务时的入口方法 * param Job $job 任务句柄对象 * param array $data 投递时传入的业务数据 */ public function fire(Job $job,$data) { //防御性判断超过最大重试次数则直接删除并记录失败 if($job-attempts()3){ $this-logFailed($data); $job-delete(); return; } try{ $this-send($data);//真正的业务逻辑 $job-delete();//成功删除任务结束生命周期 }catch(\Throwable $e){ //失败延迟10秒后重新入队再次执行 $job-release(10); } } /**任务最终失败时的回调*/ public function failed($data) { //记录失败日志、通知运维等 } }Job 对象(think\queue\Job)的常用方法方法说明fire($job, $data)队列进程调用任务时执行的方法签名固定为fire(Job $job, $data)。delete()删除任务表示执行成功。任务执行成功后必须调用否则任务会因超时被重复执行。release($delay)将任务重新放回队列$delay为延迟秒数用于失败重试。attempts()返回当前任务已被执行的次数(含本次)常用于控制最大重试次数。failed($data)任务达到最大重试次数仍失败后的回调方法用于记录日志或告警。getJobId()获取任务 ID。关键纪律fire() 中每条执行路径都必须以 delete() 或 release() 收尾。既不删除也不释放的任务会在 expire 超时后自动重新入队被再次执行导致业务重复处理。五、Redis 驱动底层原理5.1 Redis 中的数据结构以队列名 default 为例Redis 中涉及三个 keykey类型说明queues:defaultlist就绪队列存放待消费的任务。push 时 LPUSH 入队消费时 RPOP/BRPOP 出队先进先出。queues:default:delayedzset延迟队列。later() 投递的任务先存放在此成员为任务 payloadscore 为到期时间戳。queues:default:reservedzset已取出(执行中)的任务。pop 时任务从 list 移入此集合score 为“取出时间 expire”。5.2 任务流转过程立即投递(push)将任务 payload 序列化后 LPUSH 到 queues:队列名等待消费。延迟投递(later)ZADD 到 delayed 集合score 为“当前时间 延迟秒数”。消费(pop)每次取任务前先把 delayed 与 reserved 中已到期的任务迁移回 list然后从 list 尾部 RPOP(配置了 expire 时取出后同时 ZADD 到 reserved 并设置过期时间)。成功确认(delete)从 reserved 集合中 ZREM 删除该任务生命周期结束。失败重试(release)任务重新回到 list(或 delayed)等待再次执行。超时回收若任务在 reserved 中停留超过 expire 秒仍未删除下一次 pop 时会被迁移回 list 重新消费保证消费者进程崩溃(来不及 delete)时任务不丢失。5.3 任务 payload 结构任务以 JSON 字符串存储主要字段包括{ job:app\\job\\SendEmail, maxTries:3, timeout:60, data:{...业务数据...} }jobJob类名maxTries最大尝试次数(--tries传入)timeout任务超时时间(--timeout传入)data投递时传入的$data理解这一结构有助于排查问题当任务“不执行”或“重复执行”时可直接用 redis‑cli 的 LRANGE、ZRANGE、ZRANGEBYSCORE 命令观察三个 key 的内容与 score定位任务卡在哪个环节。六、实战案例6.1 案例一异步发送注册欢迎邮件场景用户注册后接口需要调用 SMTP 发送欢迎邮件同步发送耗时 2~3 秒拖慢注册接口。改造为异步后接口立即返回邮件由队列进程发送。第一步编写 Job 类app/job/SendWelcomeMail.php?php namespace app\job; use think\facade\Log; use think\queue\Job; class SendWelcomeMail { public function fire(Job $job,$data) { try{ $to $data[email]; $subject 欢迎注册; $body 亲爱的{$data[nickname]},欢迎加入我们!; $this-smtpSend($to,$subject,$body);//实际的邮件发送逻辑 $job-delete();//成功删除任务 }catch(\Throwable $e){ Log::error(邮件发送失败:.$e-getMessage()); //失败第3次以内延迟10秒重试超过则放弃 if($job-attempts()3){ $job-release(10); }else{ $this-failed($data); $job-delete(); } } } public function failed($data) { Log::error(邮件最终发送失败:.json_encode($data,JSON_UNESCAPED_UNICODE)); //可在此接入告警通知 } private function smtpSend($to,$subject,$body){/*...*/} }第二步在控制器中投递任务(生产者)?php namespace app\controller; use app\job\SendWelcomeMail; use think\facade\Queue; class Register { public function doRegister() { //...完成用户写入数据库等业务 Queue::push(SendWelcomeMail::class,[ emailuserexample.com, nickname张三, ],email);//投递到email队列接口立即返回 return json([code0,msg注册成功]); } }第三步启动消费进程(开发调试用 --once线上交给 Supervisor)php think queue:listen --queueemail --tries3 --timeout306.2 案例二订单 30 分钟未支付自动取消(延迟队列)场景下单后若 30 分钟内未支付则自动取消订单并回收库存。核心是利用 later() 投递延迟任务消费时先判断订单状态。下单时投递延迟任务//创建订单成功后投递1800秒后执行的检查任务 Queue::later(1800,app\job\CheckOrderPaid::class,[ order_id$order-id, ],order);编写延迟检查任务?php namespace app\job; use app\model\Order; use think\queue\Job; class CheckOrderPaid { public function fire(Job $job,$data) { $order Order::find($data[order_id]); if($order $order-status Order::STATUS_UNPAID){ //仍未支付取消订单、回滚库存(注意用事务行锁防并发) $order-transaction(function()use($order){ $order-status Order::STATUS_CANCELED; $order-save(); //stock回滚逻辑... }); } //已支付或订单不存在什么都不做 $job-delete();//无论哪种情况都要删除任务 } }对应的多队列消费命令(order 队列权重设为 2优先消费)php think queue:listen --queueorder:2,email --tries3 --timeout60七、生产环境部署(Supervisor)线上环境不使用 nohup 直接挂起进程而是用 Supervisor 守护队列进程实现开机自启、崩溃自动拉起、优雅重启。安装 Supervisor 后创建配置文件/etc/supervisor/conf.d/think‑queue.conf[program:think-queue-email] commandphp /var/www/app/think queue:work --queueemail --daemon --tries3 --timeout60 --sleep3 directory/var/www/app autostarttrue autorestarttrue startretries3 userwww numprocs2 #启动2个消费进程并行消费 stdout_logfile/var/log/queue-email.log stderr_logfile/var/log/queue-email.err.log stopwaitsecs30 #优雅退出等待时间应大于任务最长执行时间常用管理命令supervisorctl reread #重新读取配置 supervisorctl update #应用新增配置 supervisorctl start think-queue-email #启动 supervisorctl restart all #重启(代码发版后必须重启) supervisorctl status #查看进程状态daemon 模式下框架常驻内存代码更新后进程内仍是旧代码发版流程中必须包含队列进程重启步骤否则会出现“线上代码已修复、队列任务仍报旧 Bug”的诡异现象。八、常见问题与最佳实践8.1 常见问题排查任务不执行检查消费进程是否存活(ps / supervisorctl status)检查投递队列名与消费 --queue 参数是否一致用 redis‑cli LRANGE queues:队列名 查看任务是否堆积。任务重复执行几乎都源于 fire() 中存在既未 delete 也未 release 的分支或 expire / --timeout 设置小于任务实际执行时间导致任务被提前重新入队。延迟任务不执行只有消费进程在运行时才会把 delayed 集合中到期的任务迁移到 list消费进程停止期间延迟任务不会执行恢复后会集中补执行。daemon 模式内存持续上涨常驻进程会累积框架层面的开销业务中应避免在 Job 里定义大静态数组、缓存全表数据同时配置 --memory 让进程超限退出并由 Supervisor 拉起。数据库连接失效daemon 进程长连接可能因 MySQL wait_timeout 断开取数前可先 reconnect 或捕获异常重连。8.2 最佳实践清单按业务拆分队列(email、order、sms⋯)不同业务互不阻塞并可按需调整进程数。$data只传必要的标量字段(如 order_id)任务执行时再查询最新数据不要把大对象塞进队列。业务逻辑本身要幂等任务可能被重复执行(重试、超时重入队)写操作前先判断状态或使用唯一约束。生产端调用 Queue::push 时注意捕获连接异常必要时降级为同步处理避免 Redis 故障拖垮主业务。expire、--timeout 要大于任务最长执行时间--tries 与 attempts() 配合控制重试failed() 中接入日志告警。监控队列长度对 queues:* 的 LLEN 设置告警阈值堆积说明消费能力不足需要加进程或优化任务。本文档基于 think‑queue 3.x(ThinkPHP 6/8)官方扩展整理内容涵盖核心概念、配置参数、命令参数、底层原理、实战案例与生产部署供开发与运维参考。