lax1024git / webman-mq
High-throughput local message queue broker plugin for Webman (WAL + sharded brokers, no Redis).
1.0.6
2026-08-04 10:22 UTC
Requires
- php: >=8.1
- workerman/webman-framework: ^2.1
Requires (Dev)
None
Suggests
None
Provides
None
Conflicts
None
Replaces
None
README
基于 Webman 的单机高性能消息队列插件:内存队列 + WAL 分片 Broker,无需 Redis。
仓库:https://github.com/lax1024git/webman-mq
特性
- Direct / Topic / Fanout 交换机
- 多 Broker 分片(TCP / Unix Socket)
- WAL + checkpoint + CRC;大消息 Blob 外置
- 延时投递 / 失败重试 / 死信
- 业务用法对齐
webman/redis-queue:Client::send/Consumer
要求
- PHP >= 8.1
- Webman ^2.1
安装
composer require lax1024git/webman-mq
安装后会复制配置到:
config/plugin/webman/mq/app.phpconfig/plugin/webman/mq/process.php
并自动创建(若不存在):
app/queue/mq/ExampleWebmanMqConsumer.php
按需修改 exchanges / queue_shards,并在 .env 中调整:
# ------------------------------------------- # webman/mq(本地 Windows:TCP;Linux 同机可改 unix) # ------------------------------------------- # 消费进程数(order_queue 单分片,8~16 合理;再大抢同一 shard) QUEUE_CONSUMER_COUNT=8 # Broker 内存(WAL 回放需要,防 OOM) QUEUE_BROKER_MEMORY_LIMIT=512M # 推荐先用 TCP(最稳)。unix 需 run/ 目录存在且 broker 进程成功绑定后才有 .sock QUEUE_USE_UNIX_SOCKET=false # 队列深度上限(ready+unacked+delayed);0=不限制。压测 10 万+ 可调大 QUEUE_MAX_QUEUE_LENGTH=500000 # 每次 BATCH_POP 条数 / 轮询间隔(ms) / 有活时连续排空轮数(过大易把单分片 broker 打满,stats 超时) QUEUE_PREFETCH=128 QUEUE_CONSUME_INTERVAL_MS=5 QUEUE_CONSUME_DRAIN_MAX=64 QUEUE_BROKER_CLIENT_TIMEOUT=10 # 显式 TCP(可选,false 时默认就是这些地址): # QUEUE_SHARD_0_LISTEN=tcp://127.0.0.1:9200 # QUEUE_SHARD_1_LISTEN=tcp://127.0.0.1:9201 # QUEUE_SHARD_2_LISTEN=tcp://127.0.0.1:9202 # QUEUE_SHARD_3_LISTEN=tcp://127.0.0.1:9203 # memory | fsync_batch(推荐)| fsync(最稳最慢) # 高吞吐默认 fsync_batch;关键消息可在代码里单独 confirm=fsync QUEUE_CONFIRM_DEFAULT=fsync_batch # 须大于业务最长处理时间(秒) QUEUE_VISIBILITY_TIMEOUT=300 # WAL / checkpoint(RabbitMQ 式批量刷盘:加大 N / 间隔可减 fsync 次数) QUEUE_WAL_FSYNC_EVERY_N=1024 QUEUE_WAL_FSYNC_INTERVAL_MS=50 QUEUE_CHECKPOINT_WAL_BYTES=67108864 QUEUE_CHECKPOINT_INTERVAL_SEC=120 QUEUE_WAL_MAX_REPLAY_BYTES=134217728 QUEUE_WAL_WARN_BYTES=100663296 # 隔离 WAL 后是否允许空启动(默认 false=只读降级,防静默丢数) QUEUE_QUARANTINE_ALLOW_EMPTY_START=false # 大消息外置(默认开启) QUEUE_BLOB_ENABLED=true QUEUE_BLOB_THRESHOLD=65536
重启 Webman:
php start.php restart
# Windows: php windows.php
快速使用
发布:
\Webman\Mq\Client::send('order_queue', ['id' => 1]); \Webman\Mq\Client::publish('webmanmq_exchange', 'webmanmq.create', ['id' => 1]); // 或 \Webman\Mq\MqQueue::send('order_queue', ['id' => 1]);
消费(app/queue/mq/YourConsumer.php):
namespace app\queue\mq; use Webman\Mq\Consumer; class YourConsumer implements Consumer { public string $queue = 'order_queue'; public function consume($data): void { // 必须幂等;正常返回 → ACK;抛异常 → 延时重试 } }
运维
php webman mq:stats php webman mq:health php webman mq:clear order_queue --status=all
定位说明
当前定位为单机消息队列(at-least-once,业务侧需幂等)。默认 fsync_batch 会在确认前对 WAL 做真实 fsync。
License
MIT