当前位置:首页 > 技术 > 正文内容

ThinkPHP8集成think-worker实现多协议设备通信:TCP、WebSocket与MQTT统一接入方案

访客 技术 2026年7月27日 3

基于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环境已启用pcntlposix扩展,这两个扩展为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应用中同时支撑三种主流通信协议,满足复杂物联网系统的接入需求。

标签: ThinkPHP8

相关文章

Linux crontab 详解

1) crontab 是什么cron 是 Linux 的定时任务守护进程;crontab 是用来编辑/查看“按时间周期执行命令”的表(cron table)。常见两类:用户 crontab:每个用户一份(crontab -e 编辑)系统级 crontab / cron.d:可指定执行用户(/etc/crontab、/etc/cron.d/*)2) crontab 时间...

富文本里可以允许的 HTML 属性

一、所有标签默认允许的安全属性(极少)class        (可选)id           (通常建议禁用)title️ 注意:id 容易被滥用做锚点注入,很多系统直接禁用class 允许的话最好只允许固定前缀(如 editor-*)二、a 标签允许属性<a href="" t...

Mac 安装 Node.js 指南

方法一:通过官网安装包(最简单,适合初学者)如果你只是想快速安装并开始使用,这是最直接的方法。访问 Node.js 官网。页面会显示两个版本:LTS (Recommended For Most Users):长期支持版,最稳定。建议选这个。Current:最新特性版,包含最新功能但可能不够稳定。下载 .pkg 安装包并运行。按照安装向导点击“下一步”即可完成。方法二:使用 Homebrew 安装(...

Dom\HTML_NO_DEFAULT_NS 的副作用:自动加闭合标签

在使用Dom\HTMLDocument时,Dom\HTML_NO_DEFAULT_NS 将禁止在解析过程中设置元素的命名空间, 此设置是为了与DOMDocument向后兼容而存在的。当使用它时,已知的一个副作用就是:自动加闭合标签例如 </img> 为什么会这样?当你使用:Dom\HTML_NO_DEFAULT_NS文档会变成 无命名空间模式,此时内部更接近 XML...

Laravel 事件和监听器创建

在 Laravel 中,使用 Artisan 命令创建 Events(事件) 和 Listeners(监听器) 是非常高效的。你可以通过以下几种方式来实现:1. 手动创建单个 Event如果你只想创建一个事件类,可以使用 make:event 命令:Bashphp artisan make:event UserRegistered执行后,文件将生成在 app/Even...

自定义域名解析神器 dnsmasq

什么是 dnsmasq?dnsmasq 是一个轻量级、功能强大的网络服务工具,专为小型和中等规模网络设计。它是一个综合的网络基础设施解决方案[1]。dnsmasq 能做什么?功能说明应用场景DNS 转发与缓存将 DNS 查询转发到上游服务器(ISP、Google DNS 等),并在本地缓存结果加快 DNS 查询速度,减少外部 DNS 流量本地 DNS解析本地网络设备的主机名,无需编辑&n...

发表评论

访客

◎欢迎参与讨论,请在这里发表您的看法和观点。