Files
lserp_cs_6.0/其他程序/MQTTClient/MQTTClient/MqttConnect.cs
T
cyf ab56a9bcf7 基线 SVN r240
SVN-Revision: r240
2025-02-06 06:46:06 +00:00

308 lines
12 KiB
C#

using DevExpress.XtraGrid;
using DevExpress.XtraGrid.Views.Grid;
using MQTTnet;
using MQTTnet.Client;
using Newtonsoft.Json;
using System;
using System.Collections.Generic;
using System.Data;
using System.Data.SqlClient;
using System.Linq;
using System.Text;
using System.Threading.Tasks;
namespace MQTTClient
{
public class MqttConnect
{
/// <summary>
/// 主界面GridControl
/// </summary>
public static GridControl MainGridControl;
/// <summary>
/// 主界面GridView
/// </summary>
public static GridView MainGridView;
/// <summary>
/// 启动连接标识
/// </summary>
public bool MqttSubscribeStart;
/// <summary>
/// 终止连接标识
/// </summary>
public bool MqttSubscribeStop;
/// <summary>
/// MQTT主配置表
/// </summary>
//private readonly string mainTableName = "Mqtt_SubscribeMainTab";
/// <summary>
/// MQTT字段配置表
/// </summary>
private readonly string fieldTableName = "Mqtt_SubscribeFieldsTab";
/// <summary>
/// 数据库连接类
/// </summary>
private SqlHelper sqlHelper;
/// <summary>
/// MQTT客户端唯一Id
/// </summary>
private readonly string clientId = Guid.NewGuid().ToString().Substring(0, 13);
/// <summary>
/// MQTT客户端对象
/// </summary>
private MqttClient mqttClient = null;
/// <summary>
/// MQTT连接属性行数据
/// </summary>
private DataRow mqttDataRow = null;
/// <summary>
/// MQTT连接属性
/// </summary>
private MqttConnectModel mqttConnectModel = null;
/// <summary>
/// MQTT连接保存字段配置
/// </summary>
private List<MqttDataFieldModel> mqttDataFields = new List<MqttDataFieldModel>();
/// <summary>
/// MQTT保存业务表字段表
/// </summary>
private DataTable tempTableFields = new DataTable();
public MqttConnect(DataRow dataRow)
{
Update(dataRow);
}
public void Update(DataRow dataRow)
{
mqttDataRow = dataRow;
MainGridControl.Invoke((new Action(() =>
{
mqttDataRow["clientId"] = clientId;
})));
mqttConnectModel = new MqttConnectModel(dataRow);
sqlHelper = new SqlHelper(new SqlConnection(GetSqlConnection()));
if (!string.IsNullOrWhiteSpace(mqttConnectModel.TempTable))
{
string selectSql = $"select name from sys.columns where object_id=object_id('{mqttConnectModel.TempTable}')";
tempTableFields = sqlHelper.ExecuteDataTable(selectSql);
}
string selectDataFields = $"select * from {fieldTableName} where id = '{mqttConnectModel.Id}'";
DataTable dataFieldsTable = sqlHelper.ExecuteDataTable(selectDataFields);
mqttDataFields.Clear();
foreach (DataRow dataFieldRow in dataFieldsTable.Rows)
{
MqttDataFieldModel mqttDataFieldModel = new MqttDataFieldModel(dataFieldRow);
if (!string.IsNullOrWhiteSpace(mqttDataFieldModel.FieldName) && tempTableFields.Rows.Cast<DataRow>().Where(n => (n["name"] + "").Equals(mqttDataFieldModel.FieldName)).Count() > 0)
{
mqttDataFields.Add(mqttDataFieldModel);
}
}
}
#region 方法
/// <summary>
/// 连接mqtt服务器
/// </summary>
/// <returns></returns>
private async Task<bool> ConnectionMqttServerAsync()
{
bool isConnect = false;
if (mqttClient == null)
{
var mqttFactory = new MqttFactory();
mqttClient = mqttFactory.CreateMqttClient() as MqttClient;
mqttClient.ConnectedAsync += OnConnectedAsync;
mqttClient.DisconnectedAsync += OnDisconnectedAsync;
mqttClient.ApplicationMessageReceivedAsync += OnApplicationMessageReceivedAsync;
}
try
{
var options = new MqttClientOptions
{
ClientId = clientId,
ProtocolVersion = MQTTnet.Formatter.MqttProtocolVersion.V500,
ChannelOptions = new MqttClientTcpOptions
{
Server = mqttConnectModel.MqttServer,
Port = mqttConnectModel.MqttPort,
}
};
options.CleanSession = true;
options.Credentials = new MqttClientCredentials(mqttConnectModel.MqttUserName, Encoding.UTF8.GetBytes(mqttConnectModel.MqttPassword));
await mqttClient.ConnectAsync(options);
isConnect = true;
}
catch (Exception ex)
{
Console.WriteLine(ex.Message);
isConnect = false;
}
return isConnect;
}
/// <summary>
/// 订阅主题
/// </summary>
public async Task Subscribe(bool isConnect = false)
{
MqttSubscribeStart = true;
MqttSubscribeStop = false;
if (!isConnect)
{
MainGridControl.Invoke(new Action(() =>
{
mqttDataRow["connectState"] = "";
mqttDataRow["msg"] = "正在连接Mqtt服务器";
}));
isConnect = await ConnectionMqttServerAsync();
}
if (isConnect)
{
MainGridControl.Invoke(new Action(() =>
{
mqttDataRow["connectState"] = "已连接";
}));
string topic = mqttConnectModel.Topic;
if (string.IsNullOrEmpty(topic))
{
MainGridControl.Invoke(new Action(() =>
{
mqttDataRow["msg"] = "订阅主题不能为空";
}));
MqttSubscribeStart = false;
return;
}
await mqttClient.SubscribeAsync(topic);
//订阅成功
MainGridControl.Invoke(new Action(() =>
{
mqttDataRow["msg"] = $"订阅主题{topic}成功";
}));
}
}
/// <summary>
/// 停止连接
/// </summary>
public async Task Stop()
{
MqttSubscribeStop = true;
if (mqttClient.IsConnected)
{
bool isDisconnect = await mqttClient.TryDisconnectAsync();
if (isDisconnect)
{
MqttSubscribeStart = true;
MainGridControl.Invoke((new Action(() =>
{
mqttDataRow["msg"] = "已断开连接";
})));
}
}
else
{
MainGridControl.Invoke((new Action(() =>
{
mqttDataRow["msg"] = "正在停止连接,请稍候";
})));
}
}
/// <summary>
/// 获取数据库连接
/// </summary>
/// <returns></returns>
private string GetSqlConnection()
{
return string.Format("Server={0};Database={1};Persist Security Info=True;User ID={2};Password={3};Connection Timeout=5;MultipleActiveResultSets=true", mqttConnectModel.DbServer, mqttConnectModel.DbDatabase, mqttConnectModel.DbUserId, mqttConnectModel.DbPassword);
}
#endregion
#region 事件
/// <summary>
/// 断开mqtt连接
/// </summary>
/// <param name="arg"></param>
/// <returns></returns>
private async Task OnDisconnectedAsync(MqttClientDisconnectedEventArgs arg)
{
MainGridControl.Invoke((new Action(() =>
{
mqttDataRow["connectState"] = "已断开";
mqttDataRow["msg"] = "已断开连接";
})));
if (MqttSubscribeStart && !MqttSubscribeStop && !mqttClient.IsConnected)
{
MainGridControl.Invoke((new Action(() =>
{
mqttDataRow["msg"] = "连接失败,正在尝试重新连接";
})));
bool isConnect = await ConnectionMqttServerAsync();
if (isConnect)
{
await Subscribe(isConnect);
}
}
else if (MqttSubscribeStart && MqttSubscribeStop && !mqttClient.IsConnected)
{
MqttSubscribeStart = false;
}
}
/// <summary>
/// 连接mqtt服务
/// </summary>
/// <param name="arg"></param>
/// <returns></returns>
private Task OnConnectedAsync(MqttClientConnectedEventArgs arg)
{
MainGridControl.Invoke((new Action(() =>
{
mqttDataRow["msg"] = "已连接到MQTT服务器";
})));
return Task.CompletedTask;
}
/// <summary>
/// 接收订阅消息
/// </summary>
/// <param name="arg"></param>
/// <returns></returns>
private Task OnApplicationMessageReceivedAsync(MqttApplicationMessageReceivedEventArgs arg)
{
string messageReceive = Encoding.UTF8.GetString(arg.ApplicationMessage.PayloadSegment.ToArray());
//string messageReceive = jsonStr;
if (string.IsNullOrWhiteSpace(messageReceive))
{
return Task.CompletedTask;
}
using (SqlConnection sqlConnection = new SqlConnection(GetSqlConnection()))
{
sqlConnection.Open();
using (var bulkCopy = new SqlBulkCopy(sqlConnection))
{
DataTable dataTable = JsonConvert.DeserializeObject<DataTable>(messageReceive);
bulkCopy.DestinationTableName = mqttConnectModel.TempTable;
bulkCopy.BulkCopyTimeout = 0;
try
{
bulkCopy.ColumnMappings.Clear();
foreach (MqttDataFieldModel dataFieldModel in mqttDataFields)
{
string receiveFieldName = !string.IsNullOrWhiteSpace(dataFieldModel.ReceiveFieldName) ? dataFieldModel.ReceiveFieldName : dataFieldModel.FieldName;
if (dataTable.Columns.Contains(receiveFieldName))
{
bulkCopy.ColumnMappings.Add(receiveFieldName, dataFieldModel.FieldName);
}
}
bulkCopy.WriteToServer(dataTable);
}
catch (Exception ex)
{
MainGridControl.Invoke((new Action(() =>
{
string exMsg = ex != null ? ex.Message : "";
string innerMsg = ex != null && ex.InnerException != null ? ex.InnerException.Message : "";
mqttDataRow["msg"] = $"保存数据失败,原因:{exMsg},{innerMsg}";
})));
}
}
}
return Task.CompletedTask;
}
#endregion
}
}