Workerman 4 配合 MQTT 实现异步订阅与发布,需使用 workerman/mqtt v2.0.0+ 建立常驻长连接,启用自动重连,在 onConnect 中订阅、onMessage 中处理消息,发布时用 QoS 1 并监听 onPublishAck,集成 TP6 时应解耦避免混用 Facade。

Workerman 4 配合 MQTT 实现异步订阅与发布,核心在于用 workerman/mqtt 客户端建立长连接、自动重连,并在事件回调中处理业务逻辑。它不依赖 Web 请求生命周期,而是以常驻进程方式运行,适合物联网后台服务。
确认环境与依赖版本
确保使用的是 Workerman 4.x(非旧版 3.x),并安装兼容的 MQTT 客户端:
推荐组合:
- Workerman ≥ v4.1.0
- workerman/mqtt ≥ v2.0.0(适配 PHP 8.0+ 和 MQTT 3.1.1/5.0 协议)
- PHP ≥ 8.0(Workerman 4 要求)
安装命令:composer require workerman/workerman:4.1.* workerman/mqtt:2.0.*
编写独立 MQTT Worker 进程
不要把 MQTT 客户端塞进 HTTP 或 WebSocket Worker 中,应单独定义一个 Worker 类或匿名函数启动 MQTT 连接:
- 每个 Worker 实例对应一个 MQTT 连接,支持多实例横向扩展
- 连接参数需包含
client_id(全局唯一)、username/password(如 Broker 启用认证) - 启用
reconnect_period(如 5–10 秒)保障断网后自动恢复 - 务必在
onConnect回调中完成subscribe(),避免连接未就绪就订阅失败
示例片段:
$mqtt_worker = new Worker('none://0.0.0.0:0');
$mqtt_worker->count = 1;
$mqtt_worker->onWorkerStart = function () {
$client = new \Workerman\Mqtt\Client('mqtt://broker.example.com:1883', [
'client_id' => 'php-backend-01',
'username' => 'app',
'password' => 's3cr3t',
'keepalive' => 60,
'reconnect_period' => 5,
'clean_session' => true
]);
<pre class="brush:php;toolbar:false;">$client->onConnect = function ($mqtt) {
$mqtt->subscribe(['sensor/+/temperature', 'device/+/status'], function ($topic, $result) {
echo "Subscribed to $topic: " . ($result ? 'OK' : 'FAIL') . "\n";
});
};
$client->onMessage = function ($topic, $payload) {
$data = json_decode($payload, true);
// ⚠️ 此处是纯异步上下文,不可直接调用 ThinkPHP 的 Request/Response 对象
// 建议转为事件分发、写入 Redis 队列或调用独立 service 方法处理
handleMqttMessage($topic, $data);
};
$client->onClose = function () {
echo "MQTT connection closed, will auto-reconnect...\n";
};
$client->connect();};
安全发布消息(避免阻塞主线程)
发布操作必须在已连接状态下进行,且不能同步等待响应(MQTT 是异步协议)。关键点:
- 调用
$client->publish()后立即返回,不等 PUBACK;如需确认送达,监听onPublishAck事件(v2.0+ 支持) - 高频发布建议加队列缓冲(如用
Workerman\Timer::add(0.01, ...)批量发送),防止 TCP 拥塞 - 敏感数据(如设备控制指令)建议启用 QoS 1 并校验
message_id
发布示例:
// 在任意地方获取 client 实例(需全局持有或通过引用传递)
$client->publish('device/abc123/control', '{"cmd":"reboot","ts":' . time() . '}', [
'qos' => 1,
'retain' => false
]);
集成到 ThinkPHP6(非必需但常见)
若项目基于 TP6,可将 MQTT Worker 作为独立服务启动,与 HTTP 应用解耦:
- 在
app/command下新增MqttWorkerCommand,复用 TP6 的命令行能力启动 Worker - 通过 Redis 或数据库共享状态:HTTP 接口接收用户请求 → 写入 Redis → MQTT Worker 监听 Redis 变化 → 触发
publish() - 禁止在 MQTT 回调中直接使用
Db::或Log::等 TP6 Facade,应改用原生 PDO / Monolog 或封装无框架依赖的工具类
这样既保持 TP6 的开发体验,又不牺牲 MQTT 的实时性与稳定性。











