浅入物联网实现采集、告警、兼容tcp协议,同时也新增了USB局域网的实现【mqtt、webman、gatway】
·



下面是采集的日志;

后面新增本地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);
}
}
更多推荐
所有评论(0)