必须采用运行时参数注入与多队列路由机制实现动态投递:构造函数传入队列名和延迟时间,手动通过redisdriver按配置写入指定channel;需在async_queue.php中定义各队列配置,注册多个consumerprocess并显式绑定queue属性;重试策略可通过retryafter()在handle()中动态控制。

要在 Hyperf 3.1 中实现异步任务的动态投递(比如根据用户行为实时选择队列名、延迟时间或重试策略)并确保消费端能准确识别和执行,必须绕过静态配置硬编码,改用运行时参数注入与多队列路由机制。
动态投递:运行时指定队列名与延迟时间
Hyperf 的 Job 类默认绑定配置文件中的 default 队列,但实际业务中常需按场景分流——例如高优先级订单走 urgent 队列,普通通知走 notify 队列,且部分任务需延迟 5 分钟执行。
方法一:通过构造函数传入队列标识与延迟秒数
在 Job 类中声明可变属性,并在 handle() 执行前动态切换队列驱动上下文:
编辑 app/Job/DynamicNotifyJob.php:
public $queue; public $delay; public function __construct($params, string $queue = 'default', int $delay = 0) { $this->params = $params; $this->queue = $queue; $this->delay = $delay; }
关键点在于投递时不调用 $job->push(),而是用 Hyperf\AsyncQueue\Driver\RedisDriver 实例手动写入指定 channel:
【必须先获取对应队列配置,否则投递到不存在的 channel 将静默失败】
在 Service 中编写投递逻辑:
$driver = make(RedisDriver::class, ['config' => config('async_queue.' . $job->queue)]); $driver->push($job, $job->delay);
这一步不能省略 config 注入——Hyperf 不会自动为非 default 队列加载配置,config('async_queue.urgent') 必须在 async_queue.php 中明确定义该键。
多队列消费进程注册
单个 Consumer 进程只能监听一个 channel,若要同时消费 urgent 和 notify 两个队列,需注册多个独立进程。
第一步:在 config/autoload/async_queue.php 中补全多队列配置块:
'urgent' => [ 'driver' => RedisDriver::class, 'redis' => ['pool' => 'default'], 'channel' => 'queue:urgent', 'timeout' => 3, 'retry_seconds' => 2, 'processes' => 2, ], 'notify' => [ 'driver' => RedisDriver::class, 'redis' => ['pool' => 'default'], 'channel' => 'queue:notify', 'timeout' => 5, 'retry_seconds' => 10, 'processes' => 1, ]
第二步:创建专用消费者类,每个类绑定一个队列配置名:
新建 app/Process/UrgentConsumer.php:
#[Process(name: 'urgent-consumer')] class UrgentConsumer extends ConsumerProcess { protected string $queue = 'urgent'; }
第三步:将新进程加入 config/autoload/processes.php:
return [ Hyperf\AsyncQueue\Process\ConsumerProcess::class, App\Process\UrgentConsumer::class, App\Process\NotifyConsumer::class, ];
注意:ConsumerProcess 默认只处理 default 队列;自定义类必须显式设置 $queue 属性,否则仍读 default 配置。
动态重试策略控制
某些任务失败后应立即重试(如网络抖动),另一些则需指数退避(如第三方接口限流),不能共用全局 retry_seconds。
方法一:在 Job 中覆盖父类 $maxAttempts 并结合 handle() 内部逻辑判断是否重试
public function handle() { try { // 执行核心逻辑 } catch (ApiRateLimitException $e) { // 触发重试,但下次延迟加倍 $this->retryAfter(60); return; } catch (ConnectionException $e) { // 立即重试 $this->retryAfter(0); return; } }
【retryAfter(0) 表示下一秒立刻重试,不是立即执行;Hyperf 异步队列不支持真正“同步重试”】
方法二:投递时携带重试上下文,由 Consumer 在执行前解析策略
在 Job 构造函数中存入重试规则数组:
$this->retryPolicy = ['type' => 'exponential', 'base' => 2, 'max_delay' => 300];
然后在 handle() 开头解析该策略,调用 $this->retryAfter($calculatedDelay)。
验证动态投递是否生效
启动服务后,执行以下三步验证:
① 查看进程列表确认两个消费者均已拉起:ps aux | grep 'urgent-consumer\|notify-consumer'
② 使用 redis-cli 监控对应 channel 是否有数据写入:redis-cli -c monitor | grep 'queue:urgent'
③ 投递一条带 delay 的任务后,检查 Redis ZSET 中 score 值是否等于当前时间戳 + delay 秒:zrange queue:urgent 0 -1 WITHSCORES
如果 score 显示为整数时间戳(如 1754479200),说明 delay 单位是秒;若为小数(如 1754479200.123),则需检查 score_precision 配置是否被意外修改。











