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 { /// /// 主界面GridControl /// public static GridControl MainGridControl; /// /// 主界面GridView /// public static GridView MainGridView; /// /// 启动连接标识 /// public bool MqttSubscribeStart; /// /// 终止连接标识 /// public bool MqttSubscribeStop; /// /// MQTT主配置表 /// //private readonly string mainTableName = "Mqtt_SubscribeMainTab"; /// /// MQTT字段配置表 /// private readonly string fieldTableName = "Mqtt_SubscribeFieldsTab"; /// /// 数据库连接类 /// private SqlHelper sqlHelper; /// /// MQTT客户端唯一Id /// private readonly string clientId = Guid.NewGuid().ToString().Substring(0, 13); /// /// MQTT客户端对象 /// private MqttClient mqttClient = null; /// /// MQTT连接属性行数据 /// private DataRow mqttDataRow = null; /// /// MQTT连接属性 /// private MqttConnectModel mqttConnectModel = null; /// /// MQTT连接保存字段配置 /// private List mqttDataFields = new List(); /// /// MQTT保存业务表字段表 /// 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().Where(n => (n["name"] + "").Equals(mqttDataFieldModel.FieldName)).Count() > 0) { mqttDataFields.Add(mqttDataFieldModel); } } } #region 方法 /// /// 连接mqtt服务器 /// /// private async Task 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; } /// /// 订阅主题 /// 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}成功"; })); } } /// /// 停止连接 /// 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"] = "正在停止连接,请稍候"; }))); } } /// /// 获取数据库连接 /// /// 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 事件 /// /// 断开mqtt连接 /// /// /// 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; } } /// /// 连接mqtt服务 /// /// /// private Task OnConnectedAsync(MqttClientConnectedEventArgs arg) { MainGridControl.Invoke((new Action(() => { mqttDataRow["msg"] = "已连接到MQTT服务器"; }))); return Task.CompletedTask; } /// /// 接收订阅消息 /// /// /// 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(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 } }