goletter / hyperf-queue
Millisecond-precision async queue driver for Hyperf, compatible with hyperf/async-queue.
Requires
- php: >=8.1
- hyperf/async-queue: ~3.1.0
- hyperf/codec: ~3.1.0
- hyperf/collection: ~3.1.0
- hyperf/command: ~3.1.0
- hyperf/context: ~3.1.0
- hyperf/contract: ~3.1.0
- hyperf/coroutine: ~3.1.0
- hyperf/process: ~3.1.0
- hyperf/redis: ~3.1.0
- hyperf/support: ~3.1.0
- psr/container: ^1.0 || ^2.0
Requires (Dev)
- friendsofphp/php-cs-fixer: ^3.0
- mockery/mockery: ^1.0
- phpstan/phpstan: ^1.0
- phpunit/phpunit: ^10.0
- swoole/ide-helper: dev-master
README
Millisecond-precision async queue for Hyperf, built on top of hyperf/async-queue.
Official RedisDriver stores delay scores with time() (second precision). This package provides RedisMsDriver that:
- stores delayed / reserved scores in milliseconds
- exposes
pushMs()/Queue::laterMs()/dispatch_ms() - runs a lightweight mover coroutine (default every 5ms) so due jobs are pushed into
waitingwithout waiting forBRPOPtimeout - uses a Lua script to move due jobs atomically (safe under multiple consumers)
Requirements
- PHP >= 8.1
- Hyperf 3.1.x
hyperf/async-queue- Redis
Install
composer require goletter/hyperf-queue
Path repository (monorepo):
{
"repositories": [
{
"type": "path",
"url": "packages/goletter/hyperf-queue",
"options": { "symlink": true }
}
],
"require": {
"goletter/hyperf-queue": "*"
}
}
Publish config (or merge into existing config/autoload/async_queue.php):
php bin/hyperf.php vendor:publish goletter/hyperf-queue
Configuration
Keep the official default pool unchanged, and add a separate ms pool:
use Goletter\Queue\Driver\RedisMsDriver; use Hyperf\AsyncQueue\Driver\RedisDriver; return [ 'default' => [ 'driver' => RedisDriver::class, // ... official second-based config ], 'ms' => [ 'driver' => RedisMsDriver::class, 'redis' => [ 'pool' => 'default', ], 'channel' => '{queue-ms}', 'timeout' => 2, 'retry_milliseconds' => [100, 500, 1000, 3000], 'handle_timeout' => 10, 'move_interval_ms' => 5, 'move_batch' => 200, 'processes' => 1, 'concurrent' => [ 'limit' => 10, ], ], ];
- Official jobs: existing
AsyncQueueConsumer(defaultpool) - Millisecond jobs: package process
Goletter\Queue\Process\MsQueueConsumer(mspool, auto-registered via#[Process])
Usage
use App\Job\SendLetterJob; use Goletter\Queue\Queue; use Goletter\Server\Service\QueueService; use function Goletter\Queue\dispatch_ms; use function Hyperf\AsyncQueue\dispatch; $job = new SendLetterJob(...); // official second-based pool (delay = seconds) dispatch($job); dispatch($job, 5); $queueService->push($job, 'default', 5); // 5 seconds // millisecond pool (delay = milliseconds) $queueService->push($job, 'ms', 10); // 10ms $queueService->push($job, 'ms', 200); // 200ms Queue::laterMs(150, $job); dispatch_ms($job, 150);
Convention: on RedisMsDriver (ms pool), DriverInterface::push($job, $delay) treats $delay as milliseconds, so existing QueueService::push($job, 'ms', 10) works without API changes. On official RedisDriver, $delay remains seconds.
Jobs are normal Hyperf\AsyncQueue\Job classes — no special base class required.
Ops
php bin/hyperf.php queue:ms-info php bin/hyperf.php queue:ms-info ms php bin/hyperf.php queue:ms-info default
Production notes
- Do not share
channelbetweenRedisDriverandRedisMsDriver. Score units differ (seconds vs milliseconds). - Prefer Redis Cluster hash tags in channel names, e.g.
{queue-ms}, so related keys stay in one slot. move_interval_mstrades CPU vs delay accuracy.5is a good default (typical wake latency ≈ 5–20ms + Redis RTT).- On
mspool,push($job, $delay)delay unit is milliseconds; ondefaultit is seconds. - Retry after failure uses
retry_milliseconds(orretry_seconds * 1000). - Multi-tenant: put
tenantIdon the Job and restore context inhandle(); rate-limit at push time if needed.
How it works
pushMs(150, job)
→ ZADD {channel}:delayed score=now_ms+150
mover coroutine (every move_interval_ms)
→ Lua: due members → LPUSH {channel}:waiting
Consumer BRPOP waiting
→ reserved (ms score) → handle → ack / retry(ms) / fail
License
MIT