下面是采集的日志;


后面新增本地USB串口:2种可以分开使用,不过usb 的实现是用 python的时间调度框架,采集过来直接丢redis,基本可以完美适配原先代码;


“希望大家可以对代码进行优化提出建议”
实现的场景是:
①modbus 485的串口解析;
②使用的是有人云的4G模块,主要用于接收设备传输
③发布用redis队列,因为485串口特性是没有关联关系的,每次需要把队列放第一个,超出时间缓存直接丢弃即可,主要看你对设备采集数据要求怎么样,我这边是采集电力,所以不要求每条数据都存到,所以问题不大
④在线离线的判断是根据每个设备采集的时间 * 3倍,判断为离线,接收到数据即为在线
一共我分为几个线程:

1.发布端 

2.订阅端(主题要一直)

3.接收到的存储到缓存,定时把modbus存到数据库。

4.采集的值用来消费,计算采集存储到中的modbus报文数据,

提供代码:

下面是订阅端
<?php
namespace app\process;

use support\Redis;
use Workerman\Mqtt\Client;
use support\Log;
use Workerman\Timer;
use support\think\Db;

class MqttSubscriber
{
    /**
     * Worker 启动
     * @throws \Exception
     */
    public static function onWorkerStart($worker): void
    {
        $prefix = 'webman';
        $clientId = $prefix . '-' . bin2hex(random_bytes(4));

        $mqtt = new Client('mqtt://你的公网ip:1883', [
            'username'  => 'admin',
            'password'  => 'public',
            'client_id' => $clientId,
            'keepalive' => 60,
        ]);

        // 连接成功
        $mqtt->onConnect = function($mqtt) use ($clientId) {
            Log::debug("[MQTT-$clientId] 连接成功");
            // 查询所有启用设备
            $equipments = Db::table('equipment')
                ->where('collect_status', 1)
                ->select()
                ->toArray();

            if (empty($equipments)) {
                Log::warning("[MQTT-$clientId] 没有可订阅的设备");
                return;
            }

            foreach ($equipments as $equipment) {
                $topic = 'subscriber/' . $equipment['client_id'];
                $mqtt->subscribe($topic);
                Log::info("[MQTT-$clientId] 订阅成功: $topic");
            }
        };

        // 接收到消息
        $mqtt->onMessage = function($topic, $content) use ($clientId) {
            $clientIdPart = substr($topic, strrpos($topic, '/') + 1); // 提取 client_id
            $hex = strtoupper(bin2hex($content));

            // 心跳包判断
            if ($hex === '7777772E7573722E636E') {return;}

            // 从设备队列取最早的请求
            $queueKey = "modbus:pending:publish:{$clientIdPart}";
            $contextJson = Redis::rPop($queueKey);
            if (!$contextJson) {
                Log::warning("[MQTT-$clientId] Redis 上下文不存在,设备 $clientIdPart HEX={$hex}");
                return;
            }
            $context = json_decode($contextJson, true);

            // 更新设备最后活跃时间
            $equipmentId = $context['equipment_id'];
            $lastActiveKey = "mqtt:equipment_last_active:{$equipmentId}";
            Redis::set($lastActiveKey, date('YmdHis'));

            // 写入设备日志队列,供消费端处理
            $logKey = "mqtt:equipment_log:{$clientIdPart}";
            Redis::lPush($logKey, json_encode([
                'equipment_id' => $context['equipment_id'],
                'equipment_name' => $context['equipment_name'],
                'company_id' => $context['company_id'],
                'field_name' => $context['field_name'],
                'field_id' => $context['field_id'],
                'unit' => $context['unit'],
                'read_length' => $context['read_length'],
                'hex' => $hex,
                'request_id' => $context['request_id'],
                'report_time' => date('Y-m-d H:i:s', time()),
            ], JSON_UNESCAPED_UNICODE));
            Log::info("[MQTT-$clientId] 收到设备 {$clientIdPart} 消息 HEX={$hex}, field={$context['field_name']}, request_id={$context['request_id']} 入队成功");
        };

        // 连接错误
        $mqtt->onError = function($exception) use ($mqtt, $clientId) {
            Log::error("[MQTT-$clientId] 连接出错: " . $exception->getMessage());
            // 尝试重连
            Timer::add(5, function() use ($mqtt) {
                $mqtt->connect();
            }, [], false);
        };

        // 连接关闭,5 秒后重连
        $mqtt->onClose = function() use ($mqtt, $clientId) {
            Log::error("[MQTT-$clientId] 连接已关闭,5 秒后重连...");
            Timer::add(5, function() use ($mqtt) {
                $mqtt->connect();
            }, [], false);
        };

        $mqtt->connect();
    }
}
下面是发布端
<?php
namespace app\process;

use support\Redis;
use support\think\Db;
use think\db\exception\DataNotFoundException;
use think\db\exception\DbException;
use think\db\exception\ModelNotFoundException;
use Workerman\Mqtt\Client;
use support\Log;
use Workerman\Timer;

class MqttPublish
{
    protected static array $deviceTimers = [];   // 保存每个设备的定时器ID
    protected static ?Client $mqttClient = null; // MQTT客户端实例
    protected static string $clientId = '';

    /**
     * Worker启动时执行
     * @throws \Exception
     */
    public static function onWorkerStart($worker): void
    {
        self::$clientId = 'webman-' . bin2hex(random_bytes(4));
        $mqtt = new Client('mqtt://你的公网IP:1883', [
            'username'  => 'admin',
            'password'  => 'public',
            'client_id' => self::$clientId,
        ]);
        self::$mqttClient = $mqtt;

        self::registerCallbacks($mqtt);
        $mqtt->connect();
    }

    /**
     * 注册 MQTT 回调
     */
    protected static function registerCallbacks(Client $mqtt): void
    {
        // 连接成功
        $mqtt->onConnect = function(Client $mqtt) {
            Log::info("[MQTT-" . self::$clientId . "] 连接成功");

            // 全局队列清理定时器(只一个)
            self::startGlobalQueueCleaner();

            // 启动 Redis 队列定时器
            self::startTimerTaskQueue();

            // 刷新设备和模板缓存
            self::refreshEquipmentsAndFields();

            // 每 5 分钟刷新一次
            Timer::add(300, function () {
                self::refreshEquipmentsAndFields();
            });

            // 给每个设备创建定时任务
            $equipments = json_decode(Redis::get('equipments'), true) ?? [];
            foreach ($equipments as $equipment) {
                self::createDeviceTimer($mqtt, $equipment);
            }
        };

        // 连接出错
        $mqtt->onError = function(\Exception $exception) use ($mqtt) {
            Log::error("[MQTT-" . self::$clientId . "] 连接出错: " . $exception->getMessage());
            self::scheduleReconnect($mqtt);
        };

        // 连接关闭
        $mqtt->onClose = function() use ($mqtt) {
            Log::warning("[MQTT-" . self::$clientId . "] 连接已关闭,5 秒后重连...");
            self::scheduleReconnect($mqtt);
        };
    }

    /**
     * 调度重连任务(只执行一次,避免重复注册)
     */
    protected static ?int $reconnectTimer = null;
    protected static function scheduleReconnect(Client $mqtt): void
    {
        if (self::$reconnectTimer !== null) {
            return; // 已有定时器,不重复注册
        }

        self::$reconnectTimer = Timer::add(5, function() use ($mqtt) {
            try {
                Log::warning("[MQTT-" . self::$clientId . "] 尝试重新连接...");
                $mqtt->connect();
            } catch (\Exception $e) {
                Log::error("[MQTT-" . self::$clientId . "] 重连失败: " . $e->getMessage());
            } finally {
                // 删除定时器,允许下次再注册
                if (self::$reconnectTimer !== null) {
                    Timer::del(self::$reconnectTimer);
                    self::$reconnectTimer = null;
                }
            }
        }, [], false);
    }

    protected static function startTimerTaskQueue(): void
    {
        Timer::add(300, function() { // 每5分钟执行一次
            $queue = Redis::lRange('mqtt_timer_tasks', 0, -1);
            if (empty($queue)) return;

            // 按设备ID合并任务,取最后一次操作
            $tasksByEquipment = [];
            foreach ($queue as $item) {
                $task = json_decode($item, true);
                if (!$task) continue;
                $tasksByEquipment[(int)$task['equipment_id']] = $task['action'] ?? 'add';
            }

            foreach ($tasksByEquipment as $equipmentId => $action) {
                self::manageDeviceTimer($equipmentId, $action);
            }

            // 处理完清空队列
            Redis::del('mqtt_timer_tasks');
        });
    }

    /**
     * 管理设备定时器(添加 / 删除)
     * @param int $equipmentId 设备ID
     * @param string $action 操作类型:add=添加定时器,del=删除定时器
     */
    public static function manageDeviceTimer(int $equipmentId, string $action = 'add'): void
    {
        // 删除定时器
        if ($action === 'del') {
            if (isset(self::$deviceTimers[$equipmentId])) {
                Timer::del(self::$deviceTimers[$equipmentId]);
                unset(self::$deviceTimers[$equipmentId]);
                Log::info("[MQTT-" . self::$clientId . "] 已清理设备 {$equipmentId} 的定时器");
            }
            return;
        }

        // 添加定时器
        if ($action === 'add') {
            if (isset(self::$deviceTimers[$equipmentId])) {
                Timer::del(self::$deviceTimers[$equipmentId]);
                unset(self::$deviceTimers[$equipmentId]);
            }

            $equipment = Db::table('equipment')->where('id', $equipmentId)->find();
            if (!$equipment) return;

            self::createDeviceTimer(self::$mqttClient, $equipment);
            Log::info("[MQTT-" . self::$clientId . "] 已为设备 {$equipmentId} 创建新的定时器");
        }
    }

    /**
     * 刷新设备和模板到 Redis
     * @throws ModelNotFoundException
     * @throws DataNotFoundException
     * @throws DbException
     */
    public static function refreshEquipmentsAndFields(): void
    {
        $equipments = Db::table('equipment')
            ->where('collect_status', 1)
            ->select()
            ->toArray();

        $typeIds = array_unique(array_column($equipments, 'equipment_type_id'));
        $allFields = Db::table('device_type_modbus')
            ->whereIn('equipment_type_id', $typeIds)
            ->where('status', 1)
            ->order('sort_order')
            ->select();

        $fieldsMap = [];
        foreach ($allFields as $f) {
            $fieldsMap[$f['equipment_type_id']][] = $f;
        }

        Redis::set('equipments', json_encode($equipments));
        Redis::set('fields_map', json_encode($fieldsMap));

        Log::debug("设备和模板已刷新,设备数量:" . count($equipments));
    }

    /**
     * 为设备创建定时发送任务
     */
    protected static function createDeviceTimer(Client $mqtt, array $equipment): void
    {
        $topic = 'publish/' . $equipment['client_id'];
        $fieldsMap = json_decode(Redis::get('fields_map'), true) ?? [];
        $fields = $fieldsMap[$equipment['equipment_type_id']] ?? [];

        if (empty($fields)) {
            Log::debug("[MQTT-" . self::$clientId . "] 设备 {$equipment['client_id']} 没有启用的 Modbus 模板字段");
            return;
        }

        $interval = max(1, (int)$equipment['collect_interval']);
        $n = 0;

        $queueKey = "modbus:pending:publish:{$equipment['client_id']}";
        $ttl = 600; // 超时时间秒

        $timerId = Timer::add($interval, function() use ($mqtt, $topic, &$n, $equipment, $fields, $queueKey, $ttl) {
            $field = $fields[$n % count($fields)];
            $identifier = $field['field_name'];
            $requestId = uniqid('msg_'.date('YmdHis'), true);

            $pdu = pack('C', $equipment['slave_id'])
                . pack('C', $field['function_code'] ?? 3)
                . pack('n', $field['register_address'])
                . pack('n', $field['read_length']);

            $crc = self::crc16($pdu);
            $frame = $pdu . $crc;
            $hexData = strtoupper(bin2hex($frame));
            $binaryData = hex2bin($hexData);

            // 入队
            $data = json_encode([
                'equipment_name' => $equipment['name'],
                'equipment_id' => $equipment['id'],
                'company_id' => $equipment['company_id'],
                'field_name' => $identifier,
                'field_id' => $field['id'],
                'request_hex' => $hexData,
                'read_length' => $field['read_length'],
                'unit' => $field['unit'],
                'send_time' => date('Y-m-d H:i:s', time()),
                'request_id' => $requestId,
            ]);
            Redis::lPush($queueKey, $data);
            Redis::lTrim($queueKey, 0, 100); // 保留最新 100 条

            // 发布 MQTT
            try {
                $mqtt->publish($topic, $binaryData, 1);
            } catch (\Exception $e) {
                Log::error("[MQTT-" . self::$clientId . "] 消息发布失败([$identifier],[$requestId]): " . $e->getMessage());
            }

            Log::info('【'.$equipment['name'] . "】发送第 " . ($n + 1) . " 条消息 ($topic:$identifier:$requestId): $hexData");
            $n++;
        });

        self::$deviceTimers[$equipment['id']] = $timerId;
        // 打印定时器ID和设备信息
        Log::info("为设备 {$equipment['id']} 创建定时器,定时器ID:$timerId");
    }

    protected static function startGlobalQueueCleaner(): void
    {
        Timer::add(60, function() {
            $keys = Redis::keys('modbus:pending:publish:*');
            $now = time();
            $ttl = 600;
            foreach ($keys as $queueKey) {
                $list = Redis::lRange($queueKey, 0, -1);
                foreach ($list as $item) {
                    $ctx = json_decode($item, true);
                    if (!$ctx) continue;
                    $sendTime = strtotime($ctx['send_time']);
                    if ($sendTime !== false && $now - $sendTime > $ttl) {
                        Redis::lRem($queueKey, 0, $item);
                    }
                }
            }
        });
        Log::info("[MQTT-" . self::$clientId . "] 已启动全局队列清理定时器");
    }

    /**
     * 解析 Modbus CRC16
     */
    protected static function crc16(string $data): string
    {
        $crc = 0xFFFF;
        for ($i = 0; $i < strlen($data); $i++) {
            $crc ^= ord($data[$i]);
            for ($j = 0; $j < 8; $j++) {
                if ($crc & 0x0001) {
                    $crc = ($crc >> 1) ^ 0xA001;
                } else {
                    $crc >>= 1;
                }
            }
        }
        return pack('v', $crc);
    }
}
Logo

北京人形旗下天工造物具身智能开源社区,聚焦具身天工与慧思开物两大平台

更多推荐