找回密码
 立即注册

QQ登录

只需一步,快速开始

查看: 1787|回复: 5

物联网核心之MQTT移植

[复制链接]

13

主题

3

回帖

164

积分

版主

积分
164
发表于 2017-9-12 08:20:22 | 显示全部楼层 |阅读模式

高速电路PCB网,专注于嵌入式方案,信号完整性和电源完整性仿真分析,高速电路PCB设计,各种EDA工具(Cadence\Mentor\\AD\\CAM\ANSYS HFSS)交流学习。

您需要 登录 才可以下载或查看,没有账号?立即注册

×
本帖最后由 蓝凌风 于 2017-9-12 08:24 编辑

    在上一篇文章中,只是讲了MQTT的主要内容,至于怎么移植到STM32上,怎么使用才是最重要的关键。这里使用的平台是RT8711WIFI SOC,使用的LWIPFreeRTOS,移植使用跟STM32+LWIP是没什么区别的。
   先在Github上找到Eclipse的开源MQTT客户端程序https://github.com/eclipse/paho.mqtt.embedded-c.git,并把源码下载起来
   解压源码,再进入MQTTPacket文件夹,里面有三个文件夹
   把src里面的所有文件和samples下的transport.ctransport.h两个文件复制到工程目录下。这里我们主要的移植工作就在transport里面。打开transport.c文件,这个是MQTT连接,发送,接收的接口,源码是LinuxWindows平台,用的标准的Socket接口函数,我们这里的移植工作量很小,因为LWIP也是支持标准的Socket接口函数,只不过里面有些函数接口是LWIP不支持的,主要就是transport_open这个连接函数有区别。把原来的transport_open函数注释掉,重新写一个。
  1. int transport_open(char* addr, int port)
  2. {
  3.         int* sock = &mysock;
  4.         struct hostent *server;
  5.     struct sockaddr_in serv_addr;
  6.         static struct timeval tv;
  7.         int timeout = 1000;
  8.         fd_set readset;
  9.         fd_set writeset;

  10.     *sock = socket(AF_INET, SOCK_STREAM, 0);
  11.     if(*sock < 0)
  12.         DiagPrintf("[ERROR] Create socket failed\n");
  13.    
  14.     server = gethostbyname(addr);
  15.     if(server == NULL)
  16.         DiagPrintf("[ERROR] Get host ip failed\n");
  17.    
  18.     memset(&serv_addr,0,sizeof(serv_addr));
  19.     serv_addr.sin_family = AF_INET;
  20.     serv_addr.sin_port = htons(port);
  21.     memcpy(&serv_addr.sin_addr.s_addr,server->h_addr,server->h_length);
  22.    
  23.     if (connect(*sock,(struct sockaddr *)&serv_addr,sizeof(serv_addr)) < 0){
  24.         DiagPrintf("[ERROR] connect failed\n");
  25.         return -1;
  26.         }
  27.         tv.tv_sec = 10;  /* 1 second Timeout */
  28.         tv.tv_usec = 0;  
  29.         setsockopt(mysock, SOL_SOCKET, SO_RCVTIMEO, (char *)&timeout,sizeof(timeout));
  30.         return mysock;
  31. }
复制代码

   到这里其实移植工作就已经完成了,就是这么的简单,剩下就是怎么使用MQTT。直接上代码,这里用FreeRTOS新建一个MQTT的任务。
  1. #define THREAD_STACK_SIZE 512
  2. #define HOST_NAME "m2m.eclipse.org"
  3. #define HOST_PORT 1883

  4. void mqtt_thread( void *arg)
  5. {
  6.         MQTTPacket_connectData data = MQTTPacket_connectData_initializer;
  7.         MQTTString receivedTopic;
  8.         int rc = 0;
  9.         char buf[200];
  10.         int buflen = sizeof(buf);
  11.         int mysock = 0;
  12.         MQTTString topicString = MQTTString_initializer;
  13.         int payloadlen_in;
  14.         unsigned char* payload_in;
  15.         unsigned short msgid = 1;
  16.         int subcount;
  17.         int granted_qos =0;
  18.         unsigned char sessionPresent, connack_rc;
  19.         unsigned short submsgid;
  20.         int len = 0;
  21.         int req_qos = 1;
  22.         unsigned char dup;
  23.         int qos;
  24.         unsigned char retained;

  25.         char *host = "m2m.eclipse.org";
  26.         int port = 1883;
  27.         uint8_t  msgtypes = CONNECT;
  28.         uint32_t curtick = xTaskGetTickCount();
  29.         log_info("socket connect to server");
  30.         mysock = transport_open(host,port);
  31.         if(mysock < 0)
  32.                 return mysock;
  33.         log_notice("Sending to hostname %s port %d\n", host, port);
  34.         data.clientID.cstring = "me";
  35.         data.keepAliveInterval = 50;
  36.         data.cleansession = 1;
  37.         data.username.cstring = "";
  38.         data.password.cstring = "";
  39.         data.MQTTVersion = 4;
  40.         while(1)
  41.         {
  42.                 if((xTaskGetTickCount() - curtick) >(data.keepAliveInterval/2*1000))
  43.                 {
  44.                         if(msgtypes == 0)
  45.                         {
  46.                                 curtick = xTaskGetTickCount();
  47.                                 msgtypes = PINGREQ;
  48.                         }
  49.                 }

  50.                 switch(msgtypes)
  51.                 {
  52.                         case CONNECT:        len = MQTTSerialize_connect(buf, buflen, &data);       
  53.                                                         rc = transport_sendPacketBuffer(mysock, (unsigned char*)buf, len);
  54.                                                         if (rc == len)
  55.                                                                 log_info("send CONNECT Successfully");
  56.                                                         else
  57.                                                                 log_debug("send CONNECT failed");               
  58.                                                         log_info("MQTT concet to server!");
  59.                                                         msgtypes = 0;
  60.                                                         break;
  61.                         case CONNACK:         if (MQTTDeserialize_connack(&sessionPresent, &connack_rc, buf, buflen) != 1 || connack_rc != 0)
  62.                                                         {
  63.                                                                 log_info("Unable to connect, return code %d\n", connack_rc);
  64.                                                         }
  65.                                                         else log_info("MQTT is concet OK!");
  66.                                                         msgtypes = SUBSCRIBE;
  67.                                                         break;

  68.                         case SUBSCRIBE:        topicString.cstring = "ledtest";
  69.                                                         len = MQTTSerialize_subscribe(buf, buflen, 0, msgid, 1, &topicString, &req_qos);
  70.                                                         rc = transport_sendPacketBuffer(mysock, (unsigned char*)buf, len);
  71.                                                         if (rc == len)
  72.                                                                 log_info("send SUBSCRIBE Successfully\n");
  73.                                                         else
  74.                                                                 log_debug("send SUBSCRIBE failed\n");       
  75.                                                         log_info("client subscribe:[%s]",topicString.cstring);
  76.                                                         msgtypes = 0;
  77.                                                         break;
  78.                         case SUBACK:        rc = MQTTDeserialize_suback(&submsgid, 1, &subcount, &granted_qos, buf, buflen);                                                       
  79.                                                         log_info("granted qos is %d\n", granted_qos);                                                               
  80.                                                         msgtypes = 0;
  81.                                                         break;
  82.                         case PUBLISH:        rc = MQTTDeserialize_publish(&dup, &qos, &retained, &msgid, &receivedTopic,        &payload_in, &payloadlen_in, buf, buflen);
  83.                                                         log_info("message arrived : %s\n", payload_in);
  84.                                                         if(strstr(payload_in,"on"))
  85.                                                         {
  86.                                                                 log_notice("LED on!!");
  87.                                                         }
  88.                                                         else if(strstr(payload_in,"off"))
  89.                                                         {
  90.                                                                 log_notice("LED off!!");
  91.                                                         }
  92.                                                         if(qos == 1)
  93.                                                         {
  94.                                                                 log_info("publish qos is 1,send publish ack.");
  95.                                                                 memset(buf,0,buflen);
  96.                                                                 len = MQTTSerialize_ack(buf,buflen,PUBACK,dup,msgid);   //publish ack                       
  97.                                                                 rc = transport_sendPacketBuffer(mysock, (unsigned char*)buf, len);
  98.                                                                 if (rc == len)
  99.                                                                         log_info("send PUBACK Successfully");
  100.                                                                 else
  101.                                                                         log_debug("send PUBACK failed");                                       
  102.                                                         }
  103.                                                         msgtypes = 0;
  104.                                                         break;

  105.                         case PUBACK:        log_info("PUBACK!");
  106.                                                         msgtypes = 0;
  107.                                                         break;
  108.                         case PUBREC:        log_info("PUBREC!");     //just for qos2
  109.                                                         break;
  110.                         case PUBREL:        log_info("PUBREL!");        //just for qos2
  111.                                                         break;
  112.                         case PUBCOMP:        log_info("PUBCOMP!");        //just for qos2
  113.                                                         break;
  114.                         case PINGREQ:        len = MQTTSerialize_pingreq(buf, buflen);
  115.                                                         rc = transport_sendPacketBuffer(mysock, (unsigned char*)buf, len);
  116.                                                         if (rc == len)
  117.                                                                 log_info("send PINGREQ Successfully\n");
  118.                                                         else
  119.                                                                 log_debug("send PINGREQ failed\n");       
  120.                                                         log_info("time to ping mqtt server to take alive!");
  121.                                                         msgtypes = 0;
  122.                                                         break;
  123.                         case PINGRESP:        log_info("mqtt server Pong");                                                       
  124.                                                         msgtypes = 0;
  125.                                                         break;
  126.                 }
  127.                                 memset(buf,0,buflen);
  128.                                 rc=MQTTPacket_read(buf, buflen, transport_getdata);       
  129.                                 if(rc >0)
  130.                                 {
  131.                                         msgtypes = rc;
  132.                                         log_info("MQTT is get recv:");
  133.                                 }
  134.                                 gpio_write(&gpio_led, !gpio_read(&gpio_led));
  135.         }

  136. exit:
  137.         transport_close(mysock);
  138.     log_info("mqtt thread exit.");
  139.     vTaskDelete(NULL);
  140. }
复制代码


       MQTT服务器仍然是上篇文章用的m2m.eclipse.org,一开始就是调用 前面移植的transport_open(host,port)去连接MQTT服务器,并返回套接字。然后就是前面说的登录、订阅、发布、心跳等操作,这里在While里用状态来实现整个连接。
首先就是CONNECT登录,MQTTPacket_connectData data = MQTTPacket_connectData_initializer初始化登录数据结构体,然后再对data进行初始化,data.clientID.cstring = "me";data.keepAliveInterval = 50;data.cleansession = 1;
data.username.cstring = "";data.password.cstring = "";data.MQTTVersion = 4;这里表示cilentIDme,心跳时间为50,用户名跟密码都为空,结构初始化后,就是要对数据进去打包,调用MQTTSerialize_connect函数。打包后就是调用transport_sendPacketBuffer发送数据包。数据发送完后是调用MQTTPacket_read去接收服务器返回来的数据,并得到返回的数据包类型,根据数据包类型进行不同的逻辑处理。发送CONNECT后服务器对相应的返回CONNACK数据包表示登录成功。
登录成功后,开始订阅我们想要Topics,这里订阅一个”ledtest”Tipics。数据包的初始化、打包、发送跟登录是相似的,就不详述,直接看代码就可以了,服务器会相应的返回SUBACK
    接下来是心跳,通过获取系统的Tick来判断是否要发PINGREQ,如果的话,让状态机状态为PINGREQ后去发送心跳包,相应服务器会返回PINGRESP
    最后是PUBLISH,这是MQTT的最主要的通信协议,这里我只实现的客户端接收PUBLISH,其它的使用,其实都可以在源码的MQTTPacket里面sample可以找到例程。PUBLISH实现也很简单,定阅了Topics之后,只要其它客户端向这个Topics推送数据,服务器就会转发到订阅者上。MQTTPacket_read接收到服务器的推送,解析成PUBLISH数据包,然后状态机改变状态到PUBLISH再调用MQTTDeserialize_publish进行解包,得到推送的内容,最后根据推送的内容执行相应的动作,就实现远程控制。
     整个代码的效果如下:
    可以看到MQTT的登录,订阅,还有在PC上通过MQTT.fx推送信息给ledtestRT8711上收到推送的内容,并执行打开LED的动作。
    也许有人会问,如果STM32或者其它单片机是用WIFI模块或者GPRS模块,没有LWIP的怎么办。其实只要理解的MQTT的源码,就不难用GPRS或者WiFi模块去实现。MQTT的源码里都是对协议包进行打包解包,数据传输都是在tranport.c里面,我们完全不用transport,可以自己写通信接口,然后把打包的数据包通过模块发出去,写接收接口,把模块接收到服务器数据调用MQTT解包接口解析就可以了。

0

主题

87

回帖

174

积分

技术员

积分
174
发表于 2017-9-12 08:20:22 | 显示全部楼层
好贴,值得分享,

0

主题

58

回帖

116

积分

技术员

积分
116
发表于 2017-9-12 08:25:43 | 显示全部楼层

0

主题

63

回帖

134

积分

技术员

积分
134
发表于 2017-9-12 08:42:17 | 显示全部楼层
值得一读,谢谢,
您需要登录后才可以回帖 登录 | 立即注册

本版积分规则

关闭

站长推荐上一条 /4 下一条

内容正在加载中,请稍候……

QQ|Archiver|手机版|小黑屋|TechAIRD(gaosupcb Inc.)

GMT+8, 2026-8-25 01:17 AM , Processed in 0.055995 second(s), 27 queries .

Powered by Discuz! X3.5

© 2001-2026 Discuz! Team.

快速回复 返回顶部 返回列表