MQTT数据推送服务负责从数据库配置中读取MQTT连接列表,定时获取指定设备的数据并按照standard格式推送到MQTT服务器。
- 配置驱动: 从数据库设置中读取MQTT配置列表
- 多连接支持: 支持同时管理多个MQTT连接
- 定时推送: 根据配置的
uploadPeriod定时推送设备数据 - 数据格式化: 支持standard格式的数据格式化,可扩展其他格式
- 动态重载: 支持运行时重新加载配置
- 错误处理: 单个连接失败不影响其他连接继续工作
- 自动重连: MQTT连接断开时自动重连
MQTT配置存储在数据库的mqtt_config_list设置中,JSON格式如下:
[
{
"serverAddress": "mqtt.example.com",
"serverPort": 8883,
"username": "demo-user",
"password": "demo-password",
"useSSL": true,
"insecureSkipVerify": false,
"connectTimeout": 30,
"reconnectInterval": 30,
"keepAliveTimeout": 60,
"serviceStandard": "standard",
"allowControl": false,
"enabled": true,
"deviceIds": ["pylon_bms"],
"rewriteChannel": false,
"pushChannel": "111",
"subscribeChannel": "222",
"uploadPeriod": 60
},
{
"serverAddress": "127.0.0.1",
"serverPort": 1883,
"useSSL": false,
"connectTimeout": 30,
"reconnectInterval": 30,
"keepAliveTimeout": 60,
"serviceStandard": "standard",
"allowControl": false,
"enabled": false,
"deviceIds": ["led_green"],
"rewriteChannel": false,
"pushChannel": "",
"subscribeChannel": "",
"uploadPeriod": 60
}
]serverAddress: MQTT服务器地址serverPort: MQTT服务器端口username: MQTT用户名(可选,用于认证)password: MQTT密码(可选,用于认证)useSSL: 是否使用SSL/TLS连接insecureSkipVerify: 是否跳过SSL证书验证connectTimeout: 连接超时时间(秒)reconnectInterval: 重连间隔时间(秒)keepAliveTimeout: 保活超时时间(秒)serviceStandard: 服务标准(目前支持"standard")allowControl: 是否允许控制(本期不实现)enabled: 是否启用该配置deviceIds: 要推送数据的设备ID列表rewriteChannel: 是否重写通道pushChannel: 自定义推送通道(topic)subscribeChannel: 订阅通道(本期不实现)uploadPeriod: 上传周期(秒)
- 用户名和密码:
username和password字段用于MQTT服务器认证 - 可选认证:如果不提供用户名,将使用匿名连接
- 安全建议:生产环境中建议使用强密码,并定期更换
- 连接日志:系统会记录连接时使用的认证方式(用户名认证或匿名连接)
- SSL开关:
useSSL字段控制是否使用SSL/TLS加密连接 - 证书验证:
insecureSkipVerify字段控制是否跳过SSL证书验证 - 端口配置:
- SSL连接通常使用8883端口
- 非SSL连接通常使用1883端口
- 安全建议:
- 生产环境建议启用SSL(
useSSL: true) - 生产环境建议启用证书验证(
insecureSkipVerify: false) - 开发环境可以跳过证书验证(
insecureSkipVerify: true)
- 生产环境建议启用SSL(
- 连接日志:系统会记录连接时使用的协议(TCP或SSL/TLS)
- 连接超时:
connectTimeout设置连接MQTT服务器的超时时间(秒) - 重连间隔:
reconnectInterval设置连接断开后的重连间隔时间(秒) - 保活超时:
keepAliveTimeout设置MQTT保活心跳的超时时间(秒) - 推荐值:
connectTimeout: 30秒(网络较慢时可适当增加)reconnectInterval: 30秒(避免频繁重连)keepAliveTimeout: 60秒(标准MQTT保活时间)
推送的数据格式为:
{
"sn": "设备ID",
"time": 时间戳(毫秒),
"data": {
"pointA": 数值,
"pointB": 数值
}
}默认topic格式:lems/{system_number}/info
其中{system_number}会被替换为系统序列号(基于机器硬件ID生成的10位唯一标识符)。
如果配置了rewriteChannel=true且pushChannel不为空,则使用pushChannel作为topic。
import "s_mqtt"
// 初始化MQTT服务
s_mqtt.Init()// 启动MQTT服务
err := s_mqtt.StartMqtt(ctx)
if err != nil {
log.Fatalf("启动MQTT服务失败: %v", err)
}// 停止MQTT服务
err := s_mqtt.StopMqtt(ctx)
if err != nil {
log.Errorf("停止MQTT服务失败: %v", err)
}// 重新加载MQTT配置
err := s_mqtt.ReloadMqtt(ctx)
if err != nil {
log.Errorf("重新加载MQTT配置失败: %v", err)
}// 获取MQTT服务状态
isRunning, clientCount, clientStatusList := s_mqtt.GetMqttStatus()
fmt.Printf("服务运行状态: %v, 客户端数量: %d\n", isRunning, clientCount)
// 遍历每个客户端的状态
for i, status := range clientStatusList {
fmt.Printf("客户端 %d: 连接状态=%v, Topic=%s, 设备数量=%d\n",
i, status.IsConnected, status.Topic, status.DeviceCount)
}- SMqttConfig: MQTT配置结构体
- IDataFormatter: 数据格式化器接口
- SStandardFormatter: Standard格式实现
- SMqttClient: 单个MQTT连接管理
- SMqttManager: MQTT管理器
- 数据格式扩展: 实现
IDataFormatter接口可支持新的数据格式 - 服务标准扩展: 在
createFormatter方法中添加新的格式化器
- 单个MQTT连接失败不影响其他连接
- 连接失败时自动重连(配置AutoReconnect)
- 发布失败记录日志但不中断定时器
- 设备无数据时跳过推送
服务会记录以下日志:
- MQTT连接成功/失败
- 数据推送成功/失败
- 配置加载和重载
- 服务启动和停止
github.com/eclipse/paho.mqtt.golang: MQTT客户端库common: 通用库(设备管理、日志等)s_db: 数据库服务(配置读取)