ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

Hyperf 异步队列(async-queue)实战指南:从配置、投递到任务流转与源码原理

Hyperf 异步队列(async-queue)实战指南:从配置、投递到任务流转与源码原理 后端微服务【免费下载链接】hyperf A coroutine framework that focuses on hyperspeed and flexibility. Building microservice or middleware with ease.项目地址https://gitcode.com/gh_mirrors/hy/hyperf点击查看免费下载本指南基于 Hyperf 官方文档与仓库源码系统讲解hyperf/async-queue异步队列组件的完整使用链路。它不同于RabbitMQ、Kafka等完整消息队列只提供异步处理与异步延时处理能力适合在协程框架内做削峰填谷、耗时任务后台化等场景。读完本文你将掌握组件安装配置、两种消费进程启动方式、传统 Job 与注解两种投递方式、多队列隔离、失败/超时消息运维命令、事件监听与安全关闭并理解 Redis 驱动下的消息流转与重试原理。一、组件定位与边界Hyperf 异步队列hyperf/async-queue提供异步处理和异步延时处理两种能力核心思想是把耗时逻辑从请求链路中剥离交给独立的消费进程执行。需要注意它的能力边界它不能严格保证消息的持久化依赖底层 Redis 的持久化策略它不支持完备的 ACK 应答机制如 RabbitMQ 的 channel ack、Kafka 的 offset 提交。因此它适合允许少量消息丢失、注重吞吐与低延迟的内部异步任务场景例如邮件发送、日志落库、缓存预热、定时结算等业务强一致性、必须不丢消息的场景应选用完整 MQ 方案。二、安装与配置2.1 安装composer require hyperf/async-queue2.2 发布配置文件配置文件位于config/autoload/async_queue.php。若文件不存在可通过以下命令发布php bin/hyperf.php vendor:publish hyperf/async-queue当前版本暂时只支持Redis Driver驱动。2.3 配置项说明配置类型默认值备注enableboolfalse是否自动创建消费进程driverstringHyperf\AsyncQueue\Driver\RedisDriver::class队列驱动类channelstringqueue队列前缀channelredis.poolstringdefaultredis 连接池名称timeoutint2pop消息的超时时间秒即阻塞弹出等待时间retry_secondsint,array5失败后重新尝试间隔秒可为数组按重试次数递增handle_timeoutint10消息处理超时时间秒processesint1消费进程数concurrent.limitint10单进程内同时处理消息数协程并发上限max_messagesint0进程重启所需最大处理消息数默认 0 表示不重启一个完整的默认配置示例如下?php return [ default [ enable true, driver Hyperf\AsyncQueue\Driver\RedisDriver::class, redis [ pool default ], channel queue, timeout 2, retry_seconds 5, handle_timeout 10, processes 1, concurrent [ limit 10, ], max_messages 0, ], ];retry_seconds也支持传入数组根据重试次数动态调整重试间隔下标即第几次重试的等待秒数?php return [ default [ enable true, driver Hyperf\AsyncQueue\Driver\RedisDriver::class, channel queue, retry_seconds [1, 5, 10, 20], processes 1, ], ];2.4 配置项在源码中的实际作用从 RedisDriver 构造函数 可以看到channel、redis.pool、timeout、retry_seconds、handle_timeout均在此处被读取并注入驱动实例$channel $config[channel] ?? queue; $this-redis $container-get(RedisFactory::class)-get($config[redis][pool] ?? default); $this-timeout $config[timeout] ?? 5; $this-retrySeconds $config[retry_seconds] ?? 10; $this-handleTimeout $config[handle_timeout] ?? 10;其中timeout用于brPop的阻塞等待秒数handle_timeout决定消息进入reserved队列时的时间戳time() handleTimeout超时未 ACK 会被判定为超时retry_seconds数组取值的逻辑见 getRetrySeconds按下标attempts - 1取值超出数组长度则取最后一个值。而concurrent.limit在 Driver 基类构造函数 中创建了Hyperf\Coroutine\Concurrent并发控制器max_messages则在 consume 方法 中控制进程重启时机。三、工作原理ConsumerProcess是异步消费进程它会根据用户创建的Job或使用#[AsyncQueueMessage]注解的代码块执行消费逻辑。Job与#[AsyncQueueMessage]都是待投递且待执行的任务即数据与消费逻辑都定义在任务本身中Job类的成员变量即待消费的数据handle()方法即消费逻辑#[AsyncQueueMessage]注解的方法其入参即待消费数据方法体即消费逻辑。从源码看消费进程的启动入口是 ConsumerProcess::handle()它直接调用$this-driver-consume()而consume()见 Driver.php是一个常驻循环pop()取消息 → 创建协程回调执行 → 每处理 500 条消息检查一次队列长度 → 达到max_messages则退出循环配合进程管理器实现处理 N 条后重启的内存释放策略。四、配置异步消费进程组件提供进程配置与参数配置两种方式启动消费进程。4.1 参数配置推荐上文配置文件中enable true即可由组件自动创建消费进程。该能力由 RegisterConsumerProcessesListener 在进程启动阶段注册实现。4.2 进程配置组件已内置默认异步消费进程只需将其配置到config/autoload/processes.php中?php return [ Hyperf\AsyncQueue\Process\ConsumerProcess::class, ];也可以把以下Process加入自己的项目中配置方式与注解方式二选一即可?php declare(strict_types1); namespace App\Process; use Hyperf\AsyncQueue\Process\ConsumerProcess; use Hyperf\Process\Annotation\Process; #[Process(name: async-queue)] class AsyncQueueConsumer extends ConsumerProcess { }4.3 消费进程与配置的绑定关系默认的 ConsumerProcess 中protected string $pool default构造时通过DriverFactory获取对应 pool 的驱动并以config[processes] ?? 1作为进程实例数nums进程名固定为queue.{pool}。因此要消费其他队列必须继承该类并覆写$pool。五、如何使用多个配置某些场景需要创建多个队列配置比如把需要优先处理的消息投递到更清闲的队列中。示例如下?php return [ default [ enable true, driver Hyperf\AsyncQueue\Driver\RedisDriver::class, redis [ pool default ], channel queue, timeout 2, retry_seconds 5, handle_timeout 10, processes 1, concurrent [ limit 5, ], ], fast [ enable true, driver Hyperf\AsyncQueue\Driver\RedisDriver::class, redis [ pool default ], channel {queue:fast}, timeout 2, retry_seconds 5, handle_timeout 10, processes 1, concurrent [ limit 5, ], ], ];默认的Hyperf\AsyncQueue\Process\ConsumerProcess只会处理default配置因此需要新建一个消费fast队列的Process?php declare(strict_types1); namespace App\Process; use Hyperf\AsyncQueue\Process\ConsumerProcess; use Hyperf\Process\Annotation\Process; #[Process(name: async-queue)] class AsyncQueueConsumer extends ConsumerProcess { protected string $queue fast; }注意上述示例中覆写的成员变量名在部分版本中为$pool其语义是队列配置名即async_queue配置数组的 key。若你的组件版本提示$pool请对应使用protected string $pool fast;。这里提醒一个细节channel决定了 Redis 中实际使用的键前缀见 ChannelConfig 构造函数不同配置的channel必须不同否则多个队列会落到同一组 Redis 键上互相干扰。例如{queue:fast}对应键为{queue:fast}:waiting、{queue:fast}:delayed等其中花括号写法利用了 Redis Cluster 的 hash tag使同一队列的多个键落在同一 slot便于集群场景使用。六、生产消息投递任务6.1 传统方式定义 Job这种模式会把对象直接序列化后存入 Redis 队列因此为了保证序列化后的体积尽量不要将Container、Config等对象设置为成员变量。下面这个Job定义是不可取的同理#[Inject]注入的属性也不应作为成员变量?php declare(strict_types1); namespace App\Job; use Hyperf\AsyncQueue\Job; use Psr\Container\ContainerInterface; class ExampleJob extends Job { public $container; public $params; public function __construct(ContainerInterface $container, $params) { $this-container $container; $this-params $params; } public function handle() { // 根据参数处理具体逻辑 var_dump($this-params); } } $job make(ExampleJob::class);因为Job会被序列化所以成员变量不要包含匿名函数等无法被序列化的内容如果不清楚哪些内容无法序列化尽量使用下文注解方式。正确的Job应该只携带需要处理的数据其他依赖在handle()方法内重新获取?php declare(strict_types1); namespace App\Job; use Hyperf\AsyncQueue\Job; class ExampleJob extends Job { public $params; /** * 任务执行失败后的重试次数即最大执行次数为 $maxAttempts1 次 */ protected int $maxAttempts 2; public function __construct($params) { // 这里最好是普通数据不要使用携带 IO 的对象比如 PDO 对象 $this-params $params; } public function handle() { // 根据参数处理具体逻辑 // 通过具体参数获取模型等 // 这里的逻辑会在 ConsumerProcess 进程中执行 var_dump($this-params); } }$maxAttempts的含义在 Job 基类 中定义默认为 0表示不限制重试次数设置为 N 时任务最多执行 N1 次首次执行 N 次重试。基类还实现了CompressInterface/UnCompressInterfacecompress/uncompress在序列化前会对成员变量中可压缩的对象做压缩反序列化后自动解压。正确定义Job后需要写一个专门投递消息的Service?php declare(strict_types1); namespace App\Service; use App\Job\ExampleJob; use Hyperf\AsyncQueue\Driver\DriverFactory; use Hyperf\AsyncQueue\Driver\DriverInterface; class QueueService { protected DriverInterface $driver; public function __construct(DriverFactory $driverFactory) { $this-driver $driverFactory-get(default); } /** * 生产消息. * param $params 数据 * param int $delay 延时时间 单位秒 */ public function push($params, int $delay 0): bool { // 这里的 ExampleJob 会被序列化存到 Redis 中所以内部变量最好只传入普通数据 // 同理如果内部使用了注解 Value 会把对应对象一起序列化导致消息体变大。 // 所以这里也不推荐使用 make 方法来创建 Job 对象。 return $this-driver-push(new ExampleJob($params), $delay); } }投递消息时直接调用QueueService即可?php declare(strict_types1); namespace App\Controller; use App\Service\QueueService; use Hyperf\Di\Annotation\Inject; use Hyperf\HttpServer\Annotation\AutoController; #[AutoController] class QueueController extends AbstractController { #[Inject] protected QueueService $service; /** * 传统模式投递消息 */ public function index() { $this-service-push([ grouphyperf.io, https://doc.hyperf.io, https://www.hyperf.io, ]); return success; } }6.2 传统方式的底层投递逻辑从 RedisDriver::push() 可以看到两种投递路径delay 0时lPush到{channel}:waiting这个 Redis List 中即时消费delay 0时zAdd到{channel}:delayed这个 Redis ZSet 中score 为time() delay延时消费到期后由消费者移入 waiting。消息对象会先经 JobMessage 包装再被packer默认PhpSerializerPacker见 Driver 构造函数序列化后写入 Redis这也是Job 成员变量决定消息体积的根源。6.3 注解方式#[AsyncQueueMessage]框架除了传统方式外还提供注解方式投递消息。注解方式会在非消费环境下自动投递消息到队列因此在队列消费环境中调用带注解的方法时不会再次投递到队列而是直接在本消费进程中执行。若仍需要在队列中投递消息可在队列内部改用传统模式投递。重写上述QueueService将ExampleJob的逻辑搬到example方法中并加上AsyncQueueMessage注解?php declare(strict_types1); namespace App\Service; use Hyperf\AsyncQueue\Annotation\AsyncQueueMessage; class QueueService { #[AsyncQueueMessage] public function example($params) { // 需要异步执行的代码逻辑 // 这里的逻辑会在 ConsumerProcess 进程中执行 var_dump($params); } }投递消息时与平常调用方法完全一致?php declare(strict_types1); namespace App\Controller; use App\Service\QueueService; use Hyperf\Di\Annotation\Inject; use Hyperf\HttpServer\Annotation\AutoController; #[AutoController] class QueueController extends AbstractController { #[Inject] protected QueueService $service; /** * 注解模式投递消息 */ public function example() { $this-service-example([ grouphyperf.io, https://doc.hyperf.io, https://www.hyperf.io, ]); return success; } }6.4 注解方式的底层原理#[AsyncQueueMessage]注解定义于 AsyncQueueMessage.php支持三个参数pool默认default指定使用的队列配置名delay默认0延时投递秒数maxAttempts默认0最大重试次数。拦截逻辑在 AsyncQueueAspect::process()通过Environment::isAsyncQueue()判断当前是否处于异步队列消费环境若处于消费环境则直接执行原方法不再投递否则收集被调用方法的参数支持可变参数isVariadic读取方法级或类级注解上的pool/delay/maxAttempts构造 AnnotationJob记录类名、方法名、参数列表调用$driver-push($job, $delay)投递。AnnotationJob::handle()见 AnnotationJob.php#L33-L51在消费进程中从容器取出目标类实例、解压参数并通过Environment::setAsyncQueue(true)标记消费环境后反射调用原方法——这正是注解方式在消费进程中直接执行、不再二次投递的机制来源。七、默认脚本命令行运维组件提供以下命令用于查看与维护队列消息状态。Arguments 中queue_name为队列配置名默认defaultOptions 中-Q {channel_name}指定目标子队列如失败队列failed、超时队列timeout。7.1 展示当前队列的消息状态php bin/hyperf.php queue:info {queue_name}该命令底层调用 RedisDriver::info()分别统计waitinglLen、delayedzCard、failedlLen、timeoutlLen四个子队列的长度。7.2 重载所有失败/超时的消息到待执行队列php bin/hyperf.php queue:reload {queue_name} -Q {channel_name}底层对应 RedisDriver::reload()通过rpoplpush将failed或指定的timeout/failed子队列消息逐个移回waiting。-Q只接受timeout或failed否则抛出InvalidQueueException。7.3 销毁所有失败/超时的消息php bin/hyperf.php queue:flush {queue_name} -Q {channel_name}底层对应 RedisDriver::flush()直接del删除对应子队列键。八、事件系统组件在消费的关键节点派发事件便于开发者做监控、告警与补偿事件名称触发时机备注BeforeHandle处理消息前触发对应 Event/BeforeHandle.phpAfterHandle处理消息后触发对应 Event/AfterHandle.phpFailedHandle处理消息失败后触发对应 Event/FailedHandle.phpRetryHandle重试处理消息前触发对应 Event/RetryHandle.phpQueueLength每处理 500 个消息后触发用户可监听此事件判断失败或超时队列是否有消息积压事件派发位置集中在 Driver::getCallback()执行handle()前派发BeforeHandle成功后按Result分支处理并派发AfterHandle捕获到异常时若attempts()尚未用尽则派发RetryHandle并重试否则派发FailedHandle、把消息移入failed队列并调用$message-job()-fail($ex)。8.1 QueueLengthListener记录队列长度框架自带一个记录队列长度的监听器默认不开启需要时自行添加到listeners配置中?php declare(strict_types1); return [ Hyperf\AsyncQueue\Listener\QueueLengthListener::class ];8.2 ReloadChannelListener自动重载超时消息当消息执行超时或项目重启导致消息执行被中断消息最终都会被移动到timeout队列中。只要你能保证消息执行是幂等的同一消息执行一次或多次最终表现一致就可以开启以下监听器框架会自动将timeout队列中的消息移动到waiting队列等待下次消费?php declare(strict_types1); return [ Hyperf\AsyncQueue\Listener\ReloadChannelListener::class ];该监听器监听QueueLength事件默认执行 500 次消息后触发一次。从 ReloadChannelListener::process() 可见它只关注timeout通道且仅在队列长度大于 0 时调用$driver-reload(timeout)并输出日志因此不会对空队列产生无效操作。九、任务执行流转流程任务执行流转主要涉及以下几个子队列键前缀均为{channel}:队列名备注waiting等待消费的队列reserved正在消费的队列delayed延迟消费的队列failed消费失败的队列timeout消费超时的队列虽然超时但可能执行成功队列流转顺序如下上述流转在 RedisDriver::pop() 中有完整的对应实现先把delayed中 score 已到期的元素move到waiting每次最多 100 条见 move()再把reserved中超过handle_timeout的元素move到timeoutbrPop阻塞弹出waiting最长等待timeout秒拿到消息后zadd写入reservedscore 为time() handleTimeout并反序列化返回。消费完成后按 Result 枚举 的四种结果分支处理ACKack($data)即从reserved移除默认行为handle()未返回Result实例时视为ACKREQUEUE从reserved移除并重新投递到delayed重新排队RETRY从reserved移除、attempts()计数递增后重试DROP仅从reserved移除丢弃。十、配置多个异步队列当需要区分消费高频、低频或其他种类的消息时可以配置多个队列。1. 添加配置?php return [ default [ driver Hyperf\AsyncQueue\Driver\RedisDriver::class, channel {queue}, timeout 2, retry_seconds 5, handle_timeout 10, processes 1, concurrent [ limit 2, ], ], other [ driver Hyperf\AsyncQueue\Driver\RedisDriver::class, channel {other.queue}, timeout 2, retry_seconds 5, handle_timeout 10, processes 1, concurrent [ limit 2, ], ], ];2. 添加消费进程?php declare(strict_types1); namespace App\Process; use Hyperf\AsyncQueue\Process\ConsumerProcess; use Hyperf\Process\Annotation\Process; #[Process] class OtherConsumerProcess extends ConsumerProcess { protected string $queue other; }3. 调用use Hyperf\AsyncQueue\Driver\DriverFactory; use Hyperf\Context\ApplicationContext; $driver ApplicationContext::getContainer()-get(DriverFactory::class)-get(other); return $driver-push(new ExampleJob());DriverFactory在构造时遍历async_queue配置并校验每个driver类是否存在且实现DriverInterface见 DriverFactory.php随后按 key 缓存驱动实例get($name)取出对应驱动getConfig($name)取出对应配置供消费进程读取processes等参数。十一、安全关闭异步队列进程在终止时如果正在进行消费逻辑可能会导致错误。框架提供了ProcessStopHandler让异步队列进程安全关闭。当前信号处理器并不适配于 CoroutineServer如有需要请自行实现。安装信号处理器所需组件composer require hyperf/signal composer require hyperf/process添加配置config/autoload/signal.php?php declare(strict_types1); return [ handlers [ Hyperf\Process\Handler\ProcessStopHandler::class, ], timeout 5.0, ];安全关闭配合前文max_messages参数共同构成了优雅退出 定期重启释放内存的完整运维方案ProcessStopHandler收到终止信号后等待正在执行的协程完成timeout秒内max_messages则在进程处理满 N 条消息后主动退出由hyperf/process的进程管理器重启。十二、异步驱动之间的区别当前组件只有Hyperf\AsyncQueue\Driver\RedisDriver::class一个驱动实现其核心特征投递即时消息时整个JOB对象被序列化后lpush到 List 结构{channel}:waiting投递延时消息时序列化后zadd到 ZSet 结构{channel}:delayed以time() delay作为 score。因此存在一个需要注意的边界如果Job的参数完全一致在延时队列中后投递的消息会覆盖前面投递的消息ZSet 中 score 不同但 member 相同的元素会被更新 score 而非新增。若不想出现延时消息覆盖只需在Job里增加一个唯一的uniqid或者在使用#[AsyncQueueMessage]注解的方法上增加一个uniqid入参。十三、总结至此你已完整掌握 Hyperf 异步队列的使用闭环配置层config/autoload/async_queue.php中的enable、channel、timeout、retry_seconds、handle_timeout、processes、concurrent.limit、max_messages均能在 Driver 基类 与 RedisDriver 中找到一一对应的实现投递层传统Job继承 Job注意成员变量可序列化与#[AsyncQueueMessage]注解由 AsyncQueueAspect 拦截内部包装为AnnotationJob两种方式消费层ConsumerProcess常驻循环 协程并发 事件钩子通过waiting → reserved → delayed/failed/timeout五队列完成状态流转运维层queue:info、queue:reload、queue:flush三条命令配合QueueLength事件监听器可监控与治理积压消息。结合延时消息按 score 覆盖的边界认知与幂等设计你可以在自己的 Hyperf 项目中放心地把耗时任务迁入异步队列实现请求链路与后台任务的有效解耦。赞分享后端微服务【免费下载链接】hyperf A coroutine framework that focuses on hyperspeed and flexibility. Building microservice or middleware with ease.项目地址https://gitcode.com/gh_mirrors/hy/hyperf点击查看免费下载相关推荐Hyperf Async Queue 异步队列组件实战原理、配置与源码解析Hyperf Async Queue 异步队列组件实战原理、配置与源码解析 Async Queue 是 Hyperf 官方提供的一套轻量级异步处理组件它不依后端Web框架微服务RPC框架异步编程easy-vibe 后端实战异步任务队列原理与可靠投递指南easy vibe 后端实战异步任务队列原理与可靠投递指南 导读 用户点击导出报表后盯着转圈动画干等 30 秒这是后端架构中典型的体验灾难。本文围绕教程文档Hyperf消息队列与异步任务处理Hyperf消息队列与异步任务处理 本文全面介绍了Hyperf框架在消息队列和异步任务处理方面的强大功能。内容涵盖AMQP/RabbitMQ的深度集成、异步任务后端Web框架微服务RPC框架异步编程上一篇WeChatMsg本地导出微信聊天记录的免费工具下一篇PlotJuggler零门槛时间序列可视化的 5 个核心能力创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表