C#代码实现阿里云消息服务MNS消息监听

十年河东,十年河西,莫欺少年穷

学无止境,精益求精

近几天一直都在看阿里云的IOT云服务及消息队列MNS,一头雾水好几天了,直到今天,总算有点收获了,记录下来,方便以后查阅。

首先借用阿里云的一张图来说明:设备是如何通过云服务平台和企业服务器‘通话的’

针对此图,作如下说明:

1、物联网平台作为中间组件,主要是通过消息队MNS列来实现设备和企业服务器对话的,具体可描述为:

1.1、设备发送指令至物联网平台的MNS队列,MNS队列将设备指令收录,需要说明的是:设备发送指令是通过嵌入式开发人员开发的,例如C语言

1.2、企业通过C#、JAVA、PHP等高级语言开发人员开发监听程序,当监听到MNS队列中的设备指令时,获取指令,做相关业务处理,并发送新的设备指令至MNS队列。【例如发送快递柜关门的指令】

1.3、企业发送的指令被MNS收录,设备同样通过监听程序获取企业服务器发送的关门指令,收到关门指令的设备执行相关指令,完成自动关门操作。

以上便是设备与企业服务器之间的对话过程

下面列出C#的监听MNS代码【需要MNS C# JDK 的支持】注意:消息是经过EncodeBase64编码,接受消息要解码,发送消息要编码

异步监听:

1using System; 2using System.Threading; 3using System.Threading.Tasks; 4using Aliyun.MNS; 5using Aliyun.MNS.Model; 6using IotCommon; 7using IotDtos.MongodbDtos; 8using IotService.Device; 9using IotService.MongoDb; 10 11namespace IotListener 12{ 13 class Program 14 { 15 private static MongoLogService _logService; 16 public static string _receiptHandle; 17 public static DeviceResponseService service = new DeviceResponseService(); 18 public static Queue nativeQueue; 19 static void Main(string[] args) 20 { 21 LogstoreDatabaseSettings st = new LogstoreDatabaseSettings() { LogsCollectionName = "LogsForDg_" + DateTime.Now.ToString("yyyyMMdd") }; 22 _logService = new MongoLogService(st); 23 while (true) 24 { 25 try 26 { 27 IMNS client = new Aliyun.MNS.MNSClient(IotParm._accessKeyId, IotParm._secretAccessKey, IotParm._endpoint, IotParm._stsToken); 28 nativeQueue = client.GetNativeQueue(IotParm._queueName); 29 for (int i = 0; i < IotParm._receiveTimes; i++) 30 { 31 ReceiveMessageRequest request = new ReceiveMessageRequest(1); 32 nativeQueue.BeginReceiveMessage(request, ListenerCallback, null); 33 Thread.Sleep(1); 34 } 35 } 36 catch (Exception ex) 37 { 38 Console.WriteLine("Receive message failed, exception info: " + ex.Message); 39 } 40 } 41 } 42 43 /// <summary> 44 /// 回调函数 45 /// </summary> 46 /// <param name="ar"></param> 47 public static void ListenerCallback(IAsyncResult ar) 48 { 49 try 50 { 51 Message message = nativeQueue.EndReceiveMessage(ar).Message; 52 string Json = Base64Helper.DecodeBase64(message.Body); 53 Console.WriteLine("Message: {0}", Json); 54 Console.WriteLine("----------------------------------------------------\n"); 55 var methodValue = JsonKeyHelper.GetJsonValue(Json, "method"); 56 DeviceResponse(methodValue, Json); 57 if (!string.IsNullOrEmpty(methodValue)) 58 { 59 _logService.Create(new LogsForDgModel { CreateTime = DateTime.Now, data = Json, methodNo = methodValue }); 60 } 61 _receiptHandle = message.ReceiptHandle; 62 nativeQueue.DeleteMessage(_receiptHandle); 63 } 64 catch (Exception ex) 65 { 66 Console.WriteLine("Receive message failed, exception info: " + ex.Message); 67 } 68 69 } 70 71 /// <summary> 72 /// 响应设备上传接口 73 /// </summary> 74 /// <param name="method"></param> 75 /// <param name="message"></param> 76 public static void DeviceResponse(string method, string message) 77 { 78 switch (method) 79 { 80 case "doorClosedReport": service.doorClosedReportResponse(message); break; 81 case "doorOpenReport": service.doorOpenReportResponse(message); break; 82 case "deviceStartReportToCloud": service.deviceStartReportToCloudResponse(message); break; 83 case "qryDeviceConfig": service.qryDeviceConfigResponse(message); break; 84 case "devicePingToCloud": service.devicePingToCloudResponse(message); break; 85 case "deviceFatalReport": service.deviceFatalReportResponse(message); break; 86 case "deviceVersionReport": service.deviceVersionReportResponse(message); break; 87 case "deviceFirmwareData": service.deviceFirmwareDataResponse(message); break; 88 case "deviceLocationReport": service.deviceLocationReportResponse(message); break; 89 } 90 } 91 } 92}

View Code

同步监听:

1using Aliyun.MNS; 2using Aliyun.MNS.Model; 3using System; 4using System.Collections.Generic; 5using System.Linq; 6using System.Text; 7using System.Threading; 8using System.Threading.Tasks; 9 10namespace MnsListener 11{ 12 class Program 13 { 14 #region Private Properties 15 private const string _accessKeyId = ""; 16 private const string _secretAccessKey = ""; 17 private const string _endpoint = "http://.mns.cn-shanghai.aliyuncs.com/"; 18 private const string _stsToken = null; 19 20 private const string _queueName = "Sub"; 21 private const string _queueNamePrefix = "my"; 22 private const int _receiveTimes = 1; 23 private const int _receiveInterval = 2; 24 private const int batchSize = 6; 25 private static string _receiptHandle; 26 27 #endregion 28 29 30 static void Main(string[] args) 31 { 32 while (true) 33 { 34 try 35 { 36 IMNS client = new Aliyun.MNS.MNSClient(_accessKeyId, _secretAccessKey, _endpoint, _stsToken); 37 var nativeQueue = client.GetNativeQueue(_queueName); 38 for (int i = 0; i < _receiveTimes; i++) 39 { 40 var receiveMessageResponse = nativeQueue.ReceiveMessage(3); 41 Console.WriteLine("Receive message successfully, status code: {0}", receiveMessageResponse.HttpStatusCode); 42 Console.WriteLine("----------------------------------------------------"); 43 Message message = receiveMessageResponse.Message; 44 string s = DecodeBase64(message.Body); 45 Console.WriteLine("MessageId: {0}", message.Id); 46 Console.WriteLine("ReceiptHandle: {0}", message.ReceiptHandle); 47 Console.WriteLine("MessageBody: {0}", message.Body); 48 Console.WriteLine("MessageBodyMD5: {0}", message.BodyMD5); 49 Console.WriteLine("EnqueueTime: {0}", message.EnqueueTime); 50 Console.WriteLine("NextVisibleTime: {0}", message.NextVisibleTime); 51 Console.WriteLine("FirstDequeueTime: {0}", message.FirstDequeueTime); 52 Console.WriteLine("DequeueCount: {0}", message.DequeueCount); 53 Console.WriteLine("Priority: {0}", message.Priority); 54 Console.WriteLine("----------------------------------------------------\n"); 55 56 _receiptHandle = message.ReceiptHandle; 57 nativeQueue.DeleteMessage(_receiptHandle); 58 59 Thread.Sleep(_receiveInterval); 60 } 61 } 62 catch (Exception ex) 63 { 64 Console.WriteLine("Receive message failed, exception info: " + ex.Message); 65 } 66 } 67 68 } 69 70 ///编码 71 public static string EncodeBase64(string code, string code_type= "utf-8") 72 { 73 string encode = ""; 74 byte[] bytes = Encoding.GetEncoding(code_type).GetBytes(code); 75 try 76 { 77 encode = Convert.ToBase64String(bytes); 78 } 79 catch 80 { 81 encode = code; 82 } 83 return encode; 84 } 85 ///解码 86 public static string DecodeBase64(string code, string code_type = "utf-8") 87 { 88 string decode = ""; 89 byte[] bytes = Convert.FromBase64String(code); 90 try 91 { 92 decode = Encoding.GetEncoding(code_type).GetString(bytes); 93 } 94 catch 95 { 96 decode = code; 97 } 98 return decode; 99 } 100 } 101}

View Code

发送消息:

1using Aliyun.MNS; 2using Aliyun.MNS.Model; 3using System; 4using System.Collections.Generic; 5using System.Linq; 6using System.Text; 7using System.Threading; 8using System.Threading.Tasks; 9 10namespace MnsSendMsg 11{ 12 class Program 13 { 14 #region Private Properties 15 private const string _accessKeyId = ""; 16 private const string _secretAccessKey = ""; 17 private const string _endpoint = "http://.mns.cn-shanghai.aliyuncs.com/"; 18 private const string _stsToken = null; 19 20 private const string _queueName = "Sub"; 21 private const string _queueNamePrefix = "my"; 22 private const int _receiveTimes = 1; 23 private const int _receiveInterval = 2; 24 private const int batchSize = 6; 25 private static string _receiptHandle; 26 #endregion 27 static void Main(string[] args) 28 { 29 try 30 { 31 IMNS client = new Aliyun.MNS.MNSClient(_accessKeyId, _secretAccessKey, _endpoint, _stsToken); 32 // 1. 获取Queue的实例 33 var nativeQueue = client.GetNativeQueue(_queueName); 34 var sendMessageRequest = new SendMessageRequest(EncodeBase64("阿里云<MessageBody>计算")); 35 sendMessageRequest.DelaySeconds = 2; 36 var sendMessageResponse = nativeQueue.SendMessage(sendMessageRequest); 37 Console.WriteLine("Send message successfully,{0}", 38 sendMessageResponse.ToString()); 39 Thread.Sleep(2000); 40 } 41 catch (Exception ex) 42 { 43 Console.WriteLine("Send message failed, exception info: " + ex.Message); 44 } 45 } 46 47 ///编码 48 public static string EncodeBase64(string code, string code_type = "utf-8") 49 { 50 string encode = ""; 51 byte[] bytes = Encoding.GetEncoding(code_type).GetBytes(code); 52 try 53 { 54 encode = Convert.ToBase64String(bytes); 55 } 56 catch 57 { 58 encode = code; 59 } 60 return encode; 61 } 62 ///解码 63 public static string DecodeBase64(string code, string code_type = "utf-8") 64 { 65 string decode = ""; 66 byte[] bytes = Convert.FromBase64String(code); 67 try 68 { 69 decode = Encoding.GetEncoding(code_type).GetString(bytes); 70 } 71 catch 72 { 73 decode = code; 74 } 75 return decode; 76 } 77 } 78}

View Code

关于MNS C# JDK下载,可以去阿里云:https://help.aliyun.com/document_detail/32447.html?spm=a2c4g.11186623.6.633.61395f64IfHTRo

关于MNS队列,主题,主题订阅相关知识:https://help.aliyun.com/document_detail/34445.html?spm=a2c4g.11186623.6.542.699f38c6RO3nDS

关于阿里云AMQP队列接入,可以查询:https://help.aliyun.com/document_detail/149716.html?spm=a2c4g.11186623.6.621.2cda31b4kS1zXR

关于阿里云物联网平台,请查阅:https://help.aliyun.com/document_detail/125800.html?spm=a2c4g.11186623.6.542.7b0241c8o5r6PT

最后:阿里云物联网平台

@天才卧龙的博客

点赞
收藏

评论区

加载中...

相关推荐

MySQL:[Err] 1292 - Incorrect datetime value: ‘0000-00-00 00:00:00‘ for column ‘CREATE_TIME‘ at row 1

文章目录问题用navicat导入数据时,报错:原因这是因为当前的MySQL不支持datetime为0的情况。解决修改sql\mode:sql\mode:SQLMode定义了MySQL应支持的SQL语法、数据校验等,这样可以更容易地在不同的环境中使用MySQL。全局s

Oracle 分组与拼接字符串同时使用

SELECTT.,ROWNUMIDFROM(SELECTT.EMPLID,T.NAME,T.BU,T.REALDEPART,T.FORMATDATE,SUM(T.S0)S0,MAX(UPDATETIME)CREATETIME,LISTAGG(TOCHAR(

MySQL部分从库上面因为大量的临时表tmp_table造成慢查询

背景描述Time:20190124T00:08:14.70572408:00User@Host:@Id:Schema:sentrymetaLast_errno:0Killed:0Query_time:0.315758Lock_

皕杰报表之UUID

​在我们用皕杰报表工具设计填报报表时,如何在新增行里自动增加id呢?能新增整数排序id吗?目前可以在新增行里自动增加id,但只能用uuid函数增加UUID编码,不能新增整数排序id。uuid函数说明:获取一个UUID,可以在填报表中用来创建数据ID语法:uuid()或uuid(sep)参数说明:sep布尔值,生成的uuid中是否包含分隔符'',缺省为

手写Java HashMap源码

HashMap的使用教程HashMap的使用教程HashMap的使用教程HashMap的使用教程HashMap的使用教程22

2020年前端实用代码段,为你的工作保驾护航

有空的时候,自己总结了几个代码段,在开发中也经常使用,谢谢。1、使用解构获取json数据let jsonData  id: 1,status: "OK",data: 'a', 'b';let  id, status, data: number   jsonData;console.log(id, status, number )