rtu_mqtt_client.py 7.6 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202
  1. # -*- coding: utf-8 -*-
  2. # File Name:mqtt_chat_client.py
  3. # Python Version:3.5.1
  4. import os
  5. import django
  6. import sys
  7. BASE_DIR = os.path.dirname(os.path.abspath(__file__)) # 定位到你的django根目录
  8. sys.path.append(os.path.abspath(os.path.join(BASE_DIR, os.pardir)))
  9. os.environ.setdefault("DJANGO_SETTINGS_MODULE",
  10. "yfwlw_pro.settings") # project_name 项目名称
  11. django.setup()
  12. print("<-----python_mqtt_client_cbd is run----->")
  13. import paho.mqtt.client as mqtt
  14. import json
  15. from apps.AppInfoManage.models import Equip, Equip_type, TCCBdata, TCCBstatus, MyUser,Alarm_record, RTUstatus, RTUdata
  16. import re
  17. import datetime
  18. import shutil,os
  19. import threading
  20. class CJSONEncoder(json.JSONEncoder):
  21. def default(self, obj):
  22. if isinstance(obj, datetime):
  23. return obj.strftime('%Y-%m-%d %H:%M:%S')
  24. elif isinstance(obj, date):
  25. return obj.strftime('%Y-%m-%d')
  26. else:
  27. return json.JSONEncoder.default(self, obj)
  28. # 连接后的操作: 0为成功
  29. def on_connect(client, userdata, flags, rc):
  30. # print("Connected with result code "+str(rc))
  31. client.subscribe("yfkj/cbd_rtu/c2s/#")
  32. client.subscribe("yfkj/cbd_rtu/offline/#")
  33. # 发布消息完成回调函数:
  34. def on_publish(msg, rc):
  35. if rc == 0:
  36. print("publish success,msg = " + msg)
  37. def msg_thread(msg):
  38. print('\r')
  39. print('=================================================')
  40. print('\r')
  41. print("<-----topic:\n" + msg.topic + ';\n')
  42. print("Message:\n" + str(msg.payload) + "----->\n")
  43. # 从主题中获取imei
  44. # imei = msg.topic[14:len(msg.topic)]
  45. reg = re.compile(r"(?<=/)\d+")
  46. imei = reg.search(msg.topic).group(0)
  47. print("<-----imei:", imei, "----->")
  48. # 判断主题:
  49. if "c2s" in msg.topic:
  50. # 将json字符串解析:
  51. try:
  52. payload = json.loads(msg.payload.decode())
  53. except:
  54. pass
  55. if payload.get("cmd") == "data":
  56. print("<-----uploading status!----->")
  57. extdata = payload.get("ext")
  58. # print("extdata:", extdata)
  59. rtu_exist = Equip.objects.filter(equip_id=imei)
  60. # 设备存在,进一步判断状态表是否存在:
  61. if rtu_exist.exists():
  62. # 设备状态表存在、刷新状态表:
  63. print("<-----this equip is existed!----->")
  64. e_id = Equip.objects.get(equip_id=imei)
  65. RTUdata.objects.create(equip_id=e_id, rtu_data=extdata)
  66. if RTUstatus.objects.filter(equip_id=imei).exists():
  67. print("<-----this equip's status is existed!----->")
  68. try:
  69. sta = RTUstatus.objects.get(equip_id=imei)
  70. sta.rtu_status=extdata
  71. sta.is_online = '1'
  72. sta.save()
  73. print("<-----status update success!----->")
  74. except:
  75. print("<-----status update failed!----->")
  76. else:
  77. # 设备状态表不存在、创建状态表:
  78. print("<-----this equip's status is not existed!----->")
  79. try:
  80. e_id = Equip.objects.get(equip_id=imei)
  81. try:
  82. RTUstatus.objects.create(equip_id=e_id, rtu_status=extdata, is_online = '1')
  83. print("<-----this equip's status table re-create successed!----->")
  84. except:
  85. print("<-----this equip's status table re-create failed!----->")
  86. except:
  87. print("<-----this equip didn't exist!----->")
  88. else:
  89. # 设备不存在,在设备列表中创建:
  90. try:
  91. # 得到设备类型实例:
  92. equip_t = Equip_type.objects.get(type_id="10")
  93. try:
  94. e_id = Equip.objects.create(equip_id=imei, equip_type=equip_t)
  95. print("<-----this imei add successed!----->")
  96. RTUstatus.objects.create(equip_id=e_id, rtu_status=extdata, is_online = '1')
  97. RTUdata.objects.create(equip_id=e_id, rtu_data=extdata)
  98. except:
  99. print("<-----this imei add failed!----->")
  100. except:
  101. print("<-----register failed because this equip type is not existed!----->")
  102. # 设备异常预警
  103. elif payload.get("cmd") == "warn":
  104. print("<-----uploading warn!----->")
  105. print("%s is warn ! ! !" % imei)
  106. extdata = payload.get("ext")
  107. cbd_exist = Equip.objects.filter(equip_id=imei)
  108. # 设备存在,进一步判断状态表是否存在:
  109. if cbd_exist.exists():
  110. try:
  111. e_id = Equip.objects.get(equip_id=imei)
  112. # 创建预警记录:
  113. Alarm_record.objects.create(equip_id=e_id, alarm_desc=extdata, e_type="10")
  114. print("update warn ok!")
  115. except:
  116. print("update warn failed!")
  117. else:
  118. print("this imei is not exist!")
  119. # 离线消息:
  120. elif "offline" in msg.topic:
  121. try:
  122. # 将json字符串解析:
  123. payload = json.loads(msg.payload.decode())
  124. if payload.get("cmd") == "offline":
  125. print("<-----离线消息!----->")
  126. print("%s is offline!" % imei)
  127. cbd_exist = Equip.objects.filter(equip_id=imei)
  128. # 设备存在,进一步判断状态表是否存在:
  129. if cbd_exist.exists():
  130. try:
  131. e_id = Equip.objects.get(equip_id=imei)
  132. RTUstatus.objects.filter(equip_id=imei).update(is_online = '0',off_time = datetime.datetime.now())
  133. # 创建预警记录:
  134. Alarm_record.objects.create(equip_id=e_id, alarm_desc="{'status':0,'type':'offline'}", e_type="10")
  135. print("update offline ok!")
  136. except:
  137. print("update offline failed!")
  138. else:
  139. print("this imei is not exist!")
  140. except:
  141. pass
  142. # 从服务器接收消息的回调函数 :
  143. def on_message(client, userdata, msg):
  144. t = threading.Thread(target=msg_thread,args=(msg,))
  145. #打印出当前线程的名称和id
  146. # print(threading.currentThread().name)
  147. # t.setDaemon(True)
  148. t.start()
  149. return
  150. if __name__ == '__main__':
  151. client = mqtt.Client(
  152. client_id="PY_MQTT_CLIENT_RTU",
  153. clean_session=True,
  154. userdata=None,
  155. # protocol=MQTTv311,# 数据库版本
  156. )
  157. # 必须设置,否则会返回「Connected with result code 4」
  158. client.username_pw_set("admin", "password")
  159. #设置连接上服务器回调函数:
  160. client.on_connect = on_connect
  161. #设置接收到服务器消息回调函数:
  162. client.on_message = on_message
  163. # #设置与服务器断开连接回调函数
  164. # client.on_disconnect = on_disconnect
  165. HOST = "127.0.0.1"
  166. client.connect(HOST, 1883, 60)
  167. client.loop_forever()
  168. # # 输入发布的话题名称:
  169. # # user = input("请输入名称:")
  170. # topic = "/yfkj/scd/cmd/2001"
  171. # client.user_data_set(topic)
  172. # client.loop_start()
  173. # while True:
  174. # str = input()
  175. # if str:
  176. # client.publish("/yfkj/scd/cmd/2001", json.dumps({"topic": topic, "cmd": str}))