ThinkPHP8集成think-worker实现多协议设备通信:TCP、WebSocket与MQTT统一接入方案
基于ThinkPHP8的混合协议通信架构设计与实现
在构建现代物联网平台时,常需同时支持多种通信协议。前端Web应用依赖WebSocket实现实时交互,传统硬件设备使用TCP长连接传输数据,而大量智能终端则通过MQTT进行轻量级消息上报。面对这种多协议共存的需求,如何在单一服务中高效、稳定地处理不同类型的连接成为关键挑战。
本文将介绍一种基于ThinkPHP8和think-worker的统一通信解决方案,通过合理配置GatewayWorker进程模型,实现TCP、WebSocket和MQTT三类协议的并行处理,并确保与业务逻辑无缝集成。
技术选型分析:为何选择think-worker
尽管Swoole或Node.js也能胜任高并发场景,但考虑到团队对ThinkPHP生态的熟悉度以及后期维护成本,采用think-worker更具优势。它是Workerman在ThinkPHP框架中的官方适配版本,具备以下特性:
- 直接复用TP8的数据库、缓存、日志等组件
- 支持热重载与平滑重启
- 内置GatewayWorker机制,便于扩展分布式部署
其底层依赖关系如下:
Workerman(网络引擎)
├── think-worker(TP8封装层)
└── GatewayWorker(连接管理框架)
环境搭建与依赖安装
首先创建项目并引入必要扩展:
# 初始化项目
composer create-project topthink/think tp-multi-proto
cd tp-multi-proto
# 安装核心通信扩展
composer require topthink/think-worker workerman/gatewayclient workerman/mqtt
确保PHP环境已启用pcntl和posix扩展,这两个扩展为Linux下多进程控制提供支持。推荐运行环境为PHP 8.0+以获得最佳性能表现。
理解双进程架构:Gateway + BusinessWorker
think-worker沿用GatewayWorker的经典双进程模式:
- Gateway进程:负责物理连接的建立与维持,处理Socket读写操作,不执行任何业务代码。
- BusinessWorker进程:接收来自Gateway的数据请求,执行实际业务逻辑,如数据解析、状态更新、消息广播等。
这种分离设计使得网络I/O与计算任务解耦,提升了系统稳定性与可伸缩性。更重要的是,一个Gateway实例可以监听多个端口,每个端口绑定不同协议,从而实现多协议共存。
多协议配置实现
在config/gateway_worker.php中定义多端口监听策略:
<?php
return [
'gateway' => [
// TCP设备接入端口
'device_tcp' => [
'protocol' => 'text', // 使用文本协议接收原始数据包
'listen' => '0.0.0.0:9091',
'name' => 'TcpGateway',
'count' => 4,
'lanIp' => '127.0.0.1',
'startPort' => 4000,
'pingInterval' => 30,
'pingData' => '{"cmd":"heartbeat"}',
],
// WebSocket网页通信端口
'web_ws' => [
'protocol' => 'websocket',
'listen' => '0.0.0.0:9092',
'name' => 'WsGateway',
'count' => 4,
'lanIp' => '127.0.0.1',
'startPort' => 5000,
'pingInterval' => 55,
'pingData' => '{"type":"ping"}',
],
],
'businessWorker' => [
'name' => 'AppWorker',
'count' => 4,
'eventHandler' => 'app\listener\ProtocolHandler',
],
'register' => [
'address' => '127.0.0.1:1236',
],
'pidFile' => runtime_path() . 'worker.pid',
];
上述配置让同一服务监听两个独立端口:9091用于接收TCP设备数据,9092供浏览器端建立WebSocket连接。
事件处理器开发
创建事件处理类app/listener/ProtocolHandler.php:
<?php
namespace app\listener;
use GatewayWorker\Lib\Gateway;
use Workerman\Worker;
use Workerman\Mqtt\Client as MqttClient;
use think\facade\Log;
class ProtocolHandler
{
private static ?MqttClient $mqtt = null;
public static function onWorkerStart(Worker $worker): void
{
Log::info("BusinessWorker {$worker->id} 已启动");
// 建立MQTT客户端连接
self::connectToMqtt();
}
public static function onConnect(string $clientId): void
{
$connection = Gateway::getConnect($clientId);
$port = $connection->getRemotePort();
if ($port === 9091) {
Log::info("TCP设备接入: {$clientId}");
} elseif ($port === 9092) {
Log::info("Web客户端连接: {$clientId}");
}
}
public static function onMessage(string $clientId, $data): void
{
$connection = Gateway::getConnect($clientId);
$remotePort = $connection->getRemotePort();
switch ($remotePort) {
case 9091:
self::handleTcpDeviceData($clientId, $data);
break;
case 9092:
self::handleWebMessage($clientId, $data);
break;
}
}
private static function handleTcpDeviceData(string $clientId, string $rawData): void
{
// 解析设备报文,例如JSON格式
$packet = json_decode($rawData, true);
if (!$packet || !isset($packet['device_id'])) {
Gateway::sendToClient($clientId, json_encode(['error' => 'invalid_format']));
return;
}
$deviceId = $packet['device_id'];
$action = $packet['action'] ?? 'data';
// 模拟业务处理
Log::info("收到设备{$deviceId}数据: " . json_encode($packet));
// 向MQTT代理转发
self::publishToDeviceTopic("device/{$deviceId}/input", $rawData);
// 回复确认
Gateway::sendToClient($clientId, json_encode(['status' => 'ok']));
}
private static function handleWebMessage(string $clientId, string $msg): void
{
$request = json_decode($msg, true);
if (!$request) return;
switch ($request['type']) {
case 'subscribe':
$deviceId = $request['device_id'];
Gateway::joinGroup($clientId, "device_{$deviceId}");
Gateway::sendToClient($clientId, ['type' => 'subscribed', 'device' => $deviceId]);
break;
case 'command':
$targetId = $request['target'];
$cmd = $request['command'];
self::publishToDeviceTopic("device/{$targetId}/control", json_encode($cmd));
break;
}
}
private static function connectToMqtt(): void
{
$config = [
'username' => env('MQTT_USER', 'admin'),
'password' => env('MQTT_PASS', 'public'),
'client_id' => 'tp_server_' . uniqid(),
];
try {
self::$mqtt = new MqttClient("mqtt://".env('MQTT_HOST', '127.0.0.1').":1883", $config);
self::$mqtt->onConnect = function () {
echo "MQTT broker connected\n";
self::$mqtt?->subscribe('device/+/report');
};
self::$mqtt->onMessage = function ($topic, $content) {
// 处理从设备发来的MQTT消息
$parts = explode('/', $topic);
if (isset($parts[1])) {
$devId = $parts[1];
$groupKey = "device_{$devId}";
Gateway::sendToGroup($groupKey, [
'source' => 'mqtt',
'device' => $devId,
'data' => json_decode($content, true)
]);
}
};
self::$mqtt->connect();
} catch (\Exception $e) {
Log::error('MQTT连接失败: ' . $e->getMessage());
}
}
private static function publishToDeviceTopic(string $topic, string $payload): void
{
if (self::$mqtt && self::$mqtt->isConnected()) {
self::$mqtt?->publish($topic, $payload);
}
}
}
该处理器根据客户端来源端口判断协议类型,并调用相应处理函数。同时内建MQTT客户端用于与外部消息中间件交互,形成闭环通信链路。
启动与管理命令
通过Artisan命令控制服务生命周期:
# 启动服务
php think worker:start -d
# 停止服务
php think worker:stop
# 查看状态
php think worker:status
建议配合Supervisor或systemd进行进程守护,保障服务持续可用。
常见问题与优化建议
- 端口冲突:确保各协议监听端口无占用,避免与Nginx、MySQL等服务冲突。
- 心跳设置:TCP和WebSocket应配置合理的ping间隔,防止空闲连接被防火墙切断。
- 资源隔离:若某类协议负载极高,可将其拆分为独立GatewayWorker实例部署。
- 安全加固:生产环境应启用TLS加密,限制非法连接尝试。
通过以上配置,即可在一个ThinkPHP8应用中同时支撑三种主流通信协议,满足复杂物联网系统的接入需求。