Search by

lax1024git / webman-mq

lax1024git

High-throughput local message queue broker plugin for Webman (WAL + sharded brokers, no Redis).

Package info

github.com/lax1024git/webman-mq

pkg:composer/lax1024git/webman-mq

Statistics

Installs: 15

Dependents: 0

Suggesters: 0

Stars: 1

Open Issues: 0

1.0.6 2026-08-04 10:22 UTC

This package is auto-updated.

Last update: 2026-09-04 10:56:58 UTC


README

基于 Webman 的单机高性能消息队列插件:内存队列 + WAL 分片 Broker,无需 Redis。

仓库:https://github.com/lax1024git/webman-mq

特性

  • Direct / Topic / Fanout 交换机
  • 多 Broker 分片(TCP / Unix Socket)
  • WAL + checkpoint + CRC;大消息 Blob 外置
  • 延时投递 / 失败重试 / 死信
  • 业务用法对齐 webman/redis-queueClient::send / Consumer

要求

  • PHP >= 8.1
  • Webman ^2.1

安装

composer require lax1024git/webman-mq

安装后会复制配置到:

  • config/plugin/webman/mq/app.php
  • config/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