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点击查看免费下载异步队列Async Queue是 Hyperf 中用于异步处理与延迟异步处理的轻量级组件它不像RabbitMQ、Kafka那样提供严格的消息持久化与完整的 ACK 机制而是聚焦于把耗时的、可异步化的业务逻辑拆出去在独立的常驻进程中消费。本文将以hyperf/async-queue组件为线索结合仓库源码src/async-queue深入讲解其安装配置、Job与注解两种投递方式、多队列多进程架构、事件监听、内置命令行工具、执行流程以及安全退出等完整实战方案。读完本文你将能够在一个 Hyperf 项目中独立搭建异步任务系统并掌握其底层队列流转与容错原理。一、组件定位与适用场景在开始之前先明确该组件的边界它不是消息中间件如 RabbitMQ、Kafka 的替代品它提供的是pemrosesan asynchronous异步处理与pemrosesan asynchronous tertunda延迟异步处理能力它不保证消息的严格持久化也不支持完整的 ACK 机制。因此它最典型的应用场景是将发送短信、发送邮件、更新缓存、调用第三方接口、执行耗时计算等对实时性要求不高、失败可重试的任务从请求链路中剥离放入异步队列由后台进程消费。若你的业务对消息可靠性要求极高如支付、订单状态流转则应选择专业的消息队列中间件而不是本组件。组件当前仅支持 Redis Driver这也是官方文档与源码中唯一可用的驱动实现src/async-queue/src/Driver/RedisDriver.php。二、安装与配置2.1 安装组件composer require hyperf/async-queue2.2 发布配置文件配置文件位于config/autoload/async_queue.php。若该文件不存在可通过以下命令发布php bin/hyperf.php vendor:publish hyperf/async-queue发布后的默认配置文件源码可在仓库中查看src/async-queue/publish/async_queue.php其内容与下文配置示例一致。2.3 核心配置项配置项类型默认值说明driverstringHyperf\AsyncQueue\Driver\RedisDriver::class队列驱动channelstringqueue队列前缀channel 名redis.poolstringdefaultRedis 连接池名称timeoutint2阻塞获取消息的超时时间秒retry_secondsint,array5失败后的重试间隔秒handle_timeoutint10消息处理超时时间秒processesint1消费进程数量concurrent.limitint10并发处理消息的数量max_messagesint0进程重启前最多处理的消息数0 表示不重启完整配置示例?php return [ default [ 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也支持使用数组根据重试次数动态调整重试间隔例如第 1 次失败等待 1 秒、第 2 次失败等待 5 秒、依次类推超出数组长度后取数组最后一个值?php return [ default [ driver Hyperf\AsyncQueue\Driver\RedisDriver::class, channel queue, retry_seconds [1, 5, 10, 20], processes 1, ], ];2.4 配置项的底层实现在源码 RedisDriver.php 的构造函数中各配置项被解析如下$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; $this-channel make(ChannelConfig::class, [channel $channel]);timeout对应RedisDriver::pop()中brPop的阻塞等待时间见 RedisDriver.php阻塞超时后返回空并继续循环因此该值也间接决定消费者进程的空闲轮询频率retry_seconds在getRetrySeconds()中处理若为数组则按下标$attempts - 1取值越界时取end($this-retrySeconds)见 RedisDriver.phpconcurrent.limit在抽象驱动 Driver.php 中用于构建Hyperf\Coroutine\Concurrent并发控制对象消费消息时通过$this-concurrent-create($callback)限制同时处理的消息数见 Driver.phpprocesses在消费进程ConsumerProcess中通过$this-config[processes] ?? 1决定进程数量见 ConsumerProcess.phpmax_messages在consume()循环中控制进程处理到指定条数后自动break退出从而借助进程管理器实现处理 N 条后重启进程的内存释放策略见 Driver.php。三、工作方式与执行流程3.1 消费进程与任务的定义ConsumerProcess是异步消费进程它根据用户创建的Job或标记了#[AsyncQueueMessage]的方法来执行消费逻辑。无论是Job还是#[AsyncQueueMessage]本质上都是需要被投递并执行的任务——数据与消费逻辑都定义在任务内部Job类中的成员变量是待消费的数据handle()方法中则是消费逻辑使用#[AsyncQueueMessage]注解的方法方法参数是待消费的数据方法体是消费逻辑。工作流程可概括为下图3.2 五种队列状态与流转顺序任务的执行过程主要涉及以下队列key 均以 channel 为前缀具体拼接逻辑见 ChannelConfig.php队列名说明waiting等待被消费的队列reserved正在被消费的队列delayed延迟消费的队列failed消费失败的队列timeout消费超时的队列超时并不代表一定失败可能已执行成功对应的 Redis 数据结构waiting、failed、timeout使用listdelayed、reserved使用zset见 RedisDriver.php。队列流转顺序如下3.3 消费循环的核心逻辑在抽象驱动 Driver.php 的consume()中消费进程持续执行如下循环调用pop()取出一条消息若消息无效则ack掉否则派发BeforeHandle事件并执行$message-job()-handle()根据handle()的返回结果Result::ACK / REQUEUE / RETRY / DROP决定确认删除、重新入队重试、或丢弃处理异常时若attempts()还有剩余次数则派发RetryHandle并延迟重试否则派发FailedHandle并将消息移入failed队列同时调用$message-job()-fail($ex)见 Driver.php每处理 500 条消息派发一次QueueLength事件lengthCheckCount 500用于队列长度监控见 Driver.php。其中pop()的实现细节见 RedisDriver.php先通过move()将delayed队列中已到期的消息迁入waiting将reserved中超过handle_timeout的消息迁入timeout再通过brPop从waiting阻塞取出消息取出后立刻zadd写入reservedscore 为time() handleTimeout作为超时判定的依据。四、配置消费进程4.1 使用内置消费进程组件内置了消费进程只需在config/autoload/processes.php中注册即可?php return [ Hyperf\AsyncQueue\Process\ConsumerProcess::class, ];4.2 自定义消费进程注解方式当然也可以将以下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 { }在源码 ConsumerProcess.php 中可以看到ConsumerProcess从DriverFactory中取出$this-pool默认default对应的驱动与配置并将进程名设置为queue.{$pool}进程数量取自processes配置。此外组件的ConfigProvider会自动注册RegisterConsumerProcessesListener见 ConfigProvider.php该监听器会根据配置自动注册消费进程这也是发布配置文件中存在enable字段默认true的原因——用于控制是否自动注册消费进程。五、多配置多队列的使用一些开发者会针对特殊场景配置多个队列。例如将高优先级消息放在不那么繁忙的队列中。配置示例如下?php return [ default [ 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 [ 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配置因此使用其他配置名时需要自行创建新的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; }注意ConsumerProcess中实际生效的成员是protected string $pool默认defaultDriverFactory通过get($this-pool)获取驱动、通过getConfig($this-pool)获取配置见 ConsumerProcess.php。因此自定义消费进程指向哪个配置取决于你重写的是pool配置名而非queue。六、消息投递Producing Messages6.1 传统方式定义 Job 并 push在传统模式下Job对象会被直接序列化后存入 Redis。因此为了控制序列化后的大小尽量不要将Container、Config等对象设置为成员变量。以下Job定义是不推荐的使用#[Inject]同理由于Job会被序列化成员变量中不能包含无法序列化的内容如匿名函数。如果你不确定哪些内容不可序列化建议改用注解方式。?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应该只包含需要处理的数据其他相关数据可以在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.php 中定义Job抽象类内置protected int $maxAttempts 0并提供setMaxAttempts()/getMaxAttempts()方法。结合 Driver.php 的异常处理逻辑可知handle()抛出的异常会先消耗一次attempts()剩余次数归零后消息才进入failed队列并触发fail()回调。正确定义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); } }在DriverFactory的源码中见 DriverFactory.php构造时会遍历async_queue配置逐个校验driver类是否存在且实现了DriverInterface未通过校验会抛出InvalidDriverException然后缓存驱动实例。get($name)用于按配置名获取对应驱动__get魔术方法也支持$driverFactory-default的写法。发送消息——接下来只需调用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; } }push()的底层实现见 RedisDriver.phpdelay 0时通过lPush将序列化后的消息写入waitinglistdelay 0时通过zAdd以time() delay作为 score 写入delayedzset实现延迟投递。延迟消息的到期转移由消费侧的pop()中的move()完成。6.2 注解方式#[AsyncQueueMessage]除了传统投递方式框架还提供了注解方式。注解方式在消费环境之外时会自动将消息投递到队列因此若在队列内部使用它消息不会被重复投递而是在当前消费进程中直接执行。 若你仍需要从队列内部发送消息请使用传统方式。这一行为在 AsyncQueueAspect.php 中实现切面首先检查Environment::isAsyncQueue()若当前处于异步队列环境则直接执行原方法不投递否则解析注解元数据中的pool、delay、maxAttempts参数构造AnnotationJob并push到指定驱动。重写上面的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; } }AsyncQueueMessage注解本身支持三个参数见 AsyncQueueMessage.php#[Attribute(Attribute::TARGET_CLASS | Attribute::TARGET_METHOD)] class AsyncQueueMessage extends AbstractAnnotation { public function __construct( public string $pool default, // 投递到哪个队列配置 public int $delay 0, // 延迟秒数 public int $maxAttempts 0 // 最大重试次数 ) { } }注解方式在消费端的执行由AnnotationJob完成见 AnnotationJob.php从容器中取出目标类实例还原参数CompressInterface对象会被压缩/解压并通过闭包绑定调用目标方法。仓库测试 AsyncQueueAspectTest.php 与桩类 FooProxy.php 中覆盖了注解切面的参数收集与变长参数variadic场景。七、内置命令行工具所有命令均要求第一个参数queue_name队列配置名默认default并通过可选项-Q/--channel_name指定具体队列名如失败队列failed、超时队列timeout。7.1 查看队列当前状态$ php bin/hyperf.php queue:info {queue_name}该命令对应 InfoCommand.php底层调用Driver::info()见 RedisDriver.php输出waiting、delayed、failed、timeout四个队列的消息数。7.2 将所有失败/超时消息重新载入待处理队列php bin/hyperf.php queue:reload {queue_name} -Q {channel_name}对应Driver::reload($queue)见 RedisDriver.php默认将failed队列中的消息通过rpoplpush逐条迁回waiting当-Q指定timeout时则迁移timeout队列。非failed/timeout的队列名会抛出InvalidQueueException。7.3 销毁所有失败/超时消息php bin/hyperf.php queue:flush {queue_name} -Q {channel_name}对应Driver::flush($queue)见 RedisDriver.php默认删除failed队列-Q指定后可删除timeout等其他队列。八、事件系统8.1 事件一览事件名触发时机说明BeforeHandle处理消息之前AfterHandle处理消息之后FailedHandle消息处理失败之后RetryHandle消息处理重试之前QueueLength每处理 500 条消息触发一次用户可监听该事件判断 failed/timeout 队列是否有堆积对应源码位于 src/async-queue/src/Event 目录。其中QueueLength事件由 Driver.php 的checkQueueLength()派发携带驱动实例、队列 key 与长度。8.2 QueueLengthListener队列长度日志框架自带一个用于记录队列长度的监听器默认不启用。如需使用可自行添加到listeners配置中?php declare(strict_types1); return [ Hyperf\AsyncQueue\Listener\QueueLengthListener::class ];其实现见 QueueLengthListener.php按队列长度分级打日志长度小于 10 记debug、小于 50 记info、小于 500 记warning达到或超过 500 记error。8.3 ReloadChannelListener超时消息自动恢复当消息执行超时或项目重启导致消息执行中断时消息最终会被转移到timeout队列。只要你能够保证消息执行的幂等性同一消息执行一次或多次结果一致就可以启用以下监听器框架会自动将timeout队列中的消息重新迁回waiting队列再次消费。该监听器监听QueueLength事件默认每处理 500 条消息触发一次。?php declare(strict_types1); return [ Hyperf\AsyncQueue\Listener\ReloadChannelListener::class ];其实现见 ReloadChannelListener.php只关注timeout通道当长度大于 0 时调用$event-driver-reload($event-key)完成迁移并输出日志。此外仓库还提供了事件日志监听器 QueueHandleListener.php监听 BeforeHandle/AfterHandle/FailedHandle/RetryHandle 四个事件并输出queue日志通道[时间] Processing/Processed/Failed/Retried xxx.可作为排查消费异常的参考。九、配置多个 Async Queue多队列隔离如果你需要将高频与低频消费、或其他类型的消息分开处理可以同时配置多个队列。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());十、安全退出Safe Shutdown异步队列停止时如果仍有正在执行的消费逻辑可能会导致错误。框架提供了ProcessStopHandler来安全地停止异步队列进程。当前信号处理器尚未适配 CoroutineServer如有需要请自行实现。安装信号处理器相关组件composer require hyperf/signal composer require hyperf/process添加autoload/signal.php配置?php declare(strict_types1); return [ handlers [ Hyperf\Process\Handler\ProcessStopHandler::class, ], timeout 5.0, ];十一、不同异步驱动的差异说明组件当前提供的驱动为Hyperf\AsyncQueue\Driver\RedisDriver::class该驱动会将整个JOB序列化投递到即时队列时lpush到list结构投递到延迟队列时zadd到zset结构。正因如此如果Job参数完全相同后投递的延迟消息会覆盖先投递的同参数延迟消息因为 zset 以序列化后的消息作为 memberscore 相同时后写入会覆盖。如果你不希望延迟消息被覆盖请为Job增加一个唯一的uniqid或在注解方式的方法参数中增加一个uniqid输入参数。这一行为也可以从源码直接验证延迟投递使用zAdd($this-channel-getDelayed(), time() $delay, $data)其中 member 即完整序列化数据相同数据的 score 与 member 完全相同第二次写入自然覆盖第一次见 RedisDriver.php。十二、小结与延伸阅读通过本文你已经完整掌握了 Hyperf Async Queue 的以下能力安装配置发布config/autoload/async_queue.php理解每个配置项的语义与底层影响两种投递方式传统Job push()数据序列化入 Redis与#[AsyncQueueMessage]注解切面自动投递支持pool/delay/maxAttempts消费架构ConsumerProcess常驻进程 waiting/reserved/delayed/failed/timeout五队列流转 并发控制与进程重启策略运维手段queue:info、queue:reload、queue:flush命令QueueLengthListener与ReloadChannelListener监听器以及基于hyperf/signal的安全退出方案。如需深入源码可重点阅读以下文件配置发布模板src/async-queue/publish/async_queue.php驱动实现src/async-queue/src/Driver/RedisDriver.php、src/async-queue/src/Driver/Driver.php驱动工厂与通道配置src/async-queue/src/Driver/DriverFactory.php、src/async-queue/src/Driver/ChannelConfig.php消费进程与注解切面src/async-queue/src/Process/ConsumerProcess.php、src/async-queue/src/Aspect/AsyncQueueAspect.php事件监听器src/async-queue/src/Listener/QueueLengthListener.php、src/async-queue/src/Listener/ReloadChannelListener.php测试用例src/async-queue/tests/AsyncQueueAspectTest.php、src/async-queue/tests/RedisDriverTest.php结合 Hyperf 进程组件 与 Redis 连接池 的使用方式你可以在实际项目中组合出符合业务需求的高可用异步处理方案。赞分享后端微服务【免费下载链接】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实战指南从配置、投递到任务流转与源码原理 本指南基于 Hyperf 官方文档与仓库源码系统讲解 hyperf后端微服务Hyperf Async Queue 异步队列组件实战指南从安装配置到源码级原理Hyperf Async Queue 异步队列组件实战指南从安装配置到源码级原理 Hyperf 的 hyperf/async queue 组件提供了一套基于后端Web框架微服务RPC框架异步编程Hyperf Async Queue 异步队列组件实战原理、配置与源码解析Hyperf Async Queue 异步队列组件实战原理、配置与源码解析 Async Queue 是 Hyperf 官方提供的一套轻量级异步处理组件它不依后端Web框架微服务RPC框架异步编程上一篇Gatsby Node Model 深度解析context.nodeModel 数据层查询 API 实战指南下一篇在 llama.cpp 中集成 BLISBLAS 后端编译安装与多线程调优指南创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表