device.py 21 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466
  1. from rest_framework.views import APIView
  2. from rest_framework.response import Response
  3. from django.conf import settings
  4. from django.core.paginator import Paginator
  5. from django.forms.models import model_to_dict
  6. import time
  7. import os
  8. import datetime
  9. import uuid
  10. import json
  11. import logging
  12. import requests
  13. from smartfarming.utils import get_addr_by_lag_lng
  14. from smartfarming.serializers.device_serializers import DeviceSerializers
  15. from smartfarming.models.device import MongoDevice, MongoCBDData, MongoSCDData, MongoXYCBData
  16. from smartfarming.models.worm_forecast import MongoCBDphoto
  17. from smartfarming.models.weather import MongoQXZ_Base_Info, MongoQXZ_Alarm_Log_New, MongoQXZ_Conf, QXZdata_New, MongoQXZ_Alarm
  18. from django.db.models import Q
  19. from smartfarming.qxz import data_deal
  20. from collections import Counter, defaultdict
  21. from smartfarming.models.device import MongoDevice, DevicePestWarning, MongoDeviceType
  22. from kedong.decoration import kedong_deco, PortError
  23. from django.core.paginator import Paginator
  24. logger = logging.getLogger("data_ingestion")
  25. config_dict = settings.CONFIG
  26. device_type_en = config_dict.get("device_type_en")
  27. device_type_zh = config_dict.get("device_type_zh")
  28. class CbdScdXyDeviceSaveAPIView(APIView):
  29. permission_classes = []
  30. authentication_classes = []
  31. def post(self, request):
  32. # 测报灯 杀虫灯 性诱 设备及数据入库
  33. try:
  34. request_data = request.body
  35. data = json.loads(request_data)
  36. logger.warning(f"测报灯 杀虫灯 性诱 设备及数据入库原数据: {data}")
  37. topic = data.get("topic")
  38. payload = data.get("payload")
  39. cmd = payload.get("cmd")
  40. topic_msg = topic.split("/")
  41. now = int(time.time())
  42. if topic_msg and len(topic_msg) == 5 and cmd:
  43. device_id = topic_msg[-1]
  44. try:
  45. device_type = topic_msg[2]
  46. device_type_id = device_type_en.get(device_type)
  47. if device_type_id == 2:
  48. model = MongoSCDData
  49. elif device_type_id == 3:
  50. model = MongoCBDData
  51. elif device_type_id == 4:
  52. model = MongoXYCBData
  53. # 在设备信息表中查找是否有数据,如果没有数据则增加
  54. device_name = device_type_zh.get(str(device_type_id))
  55. device, is_created = MongoDevice.objects.get_or_create(
  56. device_id = device_id,
  57. defaults={
  58. "device_id": device_id,
  59. "device_type_id": device_type_id,
  60. "addtime": now
  61. }
  62. )
  63. if is_created:
  64. device.device_name = device_name
  65. device.save()
  66. logger.warning(f"新设备:{device_type} {device_id} 入库成功")
  67. # 获取数据并更新设备
  68. if cmd == "data":
  69. ext = payload.get("ext")
  70. if ext:
  71. # 获取设备上报的时间,同步设备
  72. try:
  73. stamp = ext.get("stamp")
  74. uptime = int((datetime.datetime.strptime(stamp, "%Y%m%d%H%M%S")).timestamp())
  75. except Exception as e:
  76. logger.error(f"同步设备时间失败:{device_id} {e}")
  77. # 增加设备数据
  78. model.objects.create(
  79. device_id = device.id,
  80. device_data = str(ext),
  81. addtime = now
  82. )
  83. lng = ext.get("lng")
  84. lat = ext.get("lat")
  85. dver_num = ext.get("dver")
  86. device.device_status = 1
  87. device.uptime = uptime
  88. if dver_num:
  89. device.dver_num = dver_num
  90. if lng and lat and dver_num and device.gps != 0:
  91. device.lng = lng
  92. device.lat = lat
  93. # 根据经纬度获取省市级
  94. is_success, province, city, district = get_addr_by_lag_lng(lat, lng)
  95. if is_success:
  96. # 更新地理位置坐标
  97. device.province = province
  98. device.city = city
  99. device.district = district
  100. device.save()
  101. elif cmd == "offline":
  102. ext = data.get("ext")
  103. if ext:
  104. # 增加设备数据
  105. model.objects.create(
  106. device_id=device_id,
  107. device_data = str(ext),
  108. addtime=now
  109. )
  110. # 更新设备状态
  111. device = MongoDevice.objects.filter(device_id=device_id).first()
  112. device.device_status = 0
  113. device.save()
  114. return Response({"code": 0, "msg": "success"})
  115. except Exception as e:
  116. logger.error(f"测报灯设备 {device_id} 处理上报数据或增加设备失败,错误原因:{e.args}")
  117. return Response({"code": 2, "msg": f"处理测报灯上报数据失败 {device_id}"})
  118. else:
  119. return Response({"code": 2, "msg": "请核对数据结构"})
  120. except Exception as e:
  121. logger.error(f"测报灯、杀虫灯、性诱设备 {e.args}")
  122. return Response({"code": 2, "msg": "failer"})
  123. class CbdPhotoAPIView(APIView):
  124. permission_classes = []
  125. authentication_classes = []
  126. def post(self, request):
  127. try:
  128. request_data = request.data
  129. logger.warning(f"测报灯图片数据入库原数据: {request_data}")
  130. device_id = request_data.get("imei")
  131. device = MongoDevice.objects.filter(device_id=device_id)
  132. media = config_dict.get("media")
  133. img_content_url = f'{config_dict.get("image_url").get("image")}/media'
  134. if device:
  135. device = device.first()
  136. d_id = device.id
  137. # 把原图下载到本地
  138. addr = config_dict.get("oss") + request_data.get("Result_image")
  139. inden_path = os.path.join(media, f"/cbd/{device.device_id}")
  140. os.makedirs(inden_path) if not os.path.exists(inden_path) else None
  141. stamp = int(datetime.now().timestamp())
  142. unique_id = uuid.uuid4()
  143. combined_id = str(stamp) + "-" + str(unique_id)
  144. addr_org = os.path.join(inden_path, f"{combined_id}.jpg")
  145. remote_content = requests.get(addr).content
  146. with open(addr_org, "wb") as f:
  147. f.write(remote_content)
  148. # 把识别后的图片下载到本地
  149. indentify = config_dict.get("oss") + request_data.get("Result")
  150. inden_path = os.path.join(media, f"/result/cbd/{device.device_id}")
  151. os.makedirs(inden_path) if not os.path.exists(inden_path) else None
  152. stamp = int(datetime.now().timestamp())
  153. unique_id = uuid.uuid4()
  154. combined_id = str(stamp) + "-" + str(unique_id)
  155. indentify_photo = os.path.join(inden_path, f"{combined_id}.jpg")
  156. indentify_remote_content = requests.get(indentify).content
  157. with open(indentify_photo, "wb") as f:
  158. f.write(indentify_remote_content)
  159. data = {
  160. "device_id": d_id,
  161. "addr": addr_org.replace(media, img_content_url),
  162. "indentify_photo":indentify_photo.replace(media, img_content_url),
  163. "indentify_result": request_data.get("Result"),
  164. "label": request_data.get("Result_code"),
  165. "photo_status": 1,
  166. "uptime": int(time.time()),
  167. "addtime": int(time.time())
  168. }
  169. photo = MongoCBDphoto(**data)
  170. photo.save()
  171. return Response({"code": 0, "msg": "success"})
  172. else:
  173. return Response({"code": 2, "msg": "该项目不存在此设备"})
  174. except Exception as e:
  175. logger.error(f"测报灯图片 {e.args}")
  176. return Response({"code": 2, "msg": "failer"})
  177. class QxzDeviceAddAPIViw(APIView):
  178. permission_classes = []
  179. authentication_classes = []
  180. def post(self, request):
  181. # 气象站上传数据
  182. try:
  183. request_data = request.body
  184. request_data = json.loads(request_data)
  185. device_id = request_data.get("StationID")
  186. uptime = request_data.get("MonitorTime")
  187. data = request_data.get("data")
  188. terminalStatus = request_data.get("terminalStatus")
  189. cmd = request_data.get("cmd")
  190. up = uptime.replace(" ", "")
  191. uptime_tp = int((datetime.datetime.strptime(up, "%Y-%m-%d%H:%M:%S")).timestamp())
  192. # 获取该设备的预警配置数据
  193. alarm = MongoQXZ_Alarm.objects.filter(device_id=device_id)
  194. qxz_e_conf = MongoQXZ_Conf.objects.filter(device_id=device_id)
  195. qxz_e_conf = qxz_e_conf.first() if qxz_e_conf else None
  196. if not qxz_e_conf:
  197. return Response({"code": 2, "msg": "failer"})
  198. if data:
  199. qxz_e_conf = model_to_dict(qxz_e_conf)
  200. qx_ek = {}
  201. result_tp_fin = ""
  202. for i in data:
  203. tp_value = i.get("eValue")
  204. if tp_value:
  205. ek = i.get("eKey")
  206. qx_ek[ek] = tp_value
  207. if alarm:
  208. # 存在预警配置文件, 先查看预警配置文件中是否有配置 -- "0#6" 表示大于6则报警 "1#5" 表示小于5报警 "0#" 表示不配置
  209. alarms = alarm.first()
  210. alarm_config = eval(alarms.conf)
  211. dat = alarm_config.get("dat")
  212. for m, n in dat.items():
  213. n_sp = n.split("#")
  214. if n_sp[1]:
  215. if ek == m:
  216. # 查询具体含义
  217. zh = qxz_e_conf.get(m)
  218. zh_k = zh.split("#")
  219. result = ""
  220. if n_sp[0] == "0":
  221. if float(tp_value) > float(n_sp[1]):
  222. # 组织预警信息
  223. result = f"为{tp_value},大于{n_sp[1]}"
  224. elif n_sp[0] == "1":
  225. if float(tp_value) < float(n_sp[1]):
  226. result = f"为{tp_value},小于{n_sp[1]}"
  227. if result:
  228. result_tp = f"{zh_k[0]}{result}{zh_k[1]},"
  229. result_tp_fin += result_tp
  230. if result_tp_fin:
  231. alarm_new = MongoQXZ_Alarm_Log_New()
  232. alarm_new.warning_content = result_tp_fin
  233. alarm_new.upl_time = uptime_tp
  234. alarm_new.save()
  235. logger.warning(f"{device_id} 产生预警")
  236. # 30分钟上报一次的数据
  237. qx_ek["device_id"] = device_id
  238. qx_ek["uptime"] = uptime_tp
  239. qxz_data = QXZdata_New(**qx_ek)
  240. qxz_data.save()
  241. MongoDevice.objects.filter(device_id=device_id).update(uptime=uptime_tp, device_status=1)
  242. return Response({"code": 0, "msg": "success"})
  243. if terminalStatus:
  244. base_info_obj, is_created = MongoQXZ_Base_Info.objects.update_or_create(
  245. device_id=device_id,
  246. defaults={
  247. "volt": terminalStatus.get("VOLT"),
  248. "rssi": terminalStatus.get("RSSI"),
  249. "uptime": uptime_tp
  250. }
  251. )
  252. iccid = terminalStatus.get("ICCID")
  253. lng = terminalStatus.get("longitude")
  254. lat = terminalStatus.get("latitude")
  255. led = terminalStatus.get("Dotled")
  256. dver = terminalStatus.get("Version")
  257. device, is_created = MongoDevice.objects.get_or_create(device_id=device_id)
  258. if iccid:
  259. base_info_obj.iccid = iccid
  260. if lng:
  261. base_info_obj.lng = lng
  262. device.lng = lng
  263. if lat:
  264. base_info_obj.lat = lat
  265. device.lat = lat
  266. if led:
  267. base_info_obj.led = led
  268. if dver:
  269. base_info_obj.dver = dver
  270. # 如果经纬度均存在
  271. if lat and lng:
  272. is_success, province, city, district = get_addr_by_lag_lng(lat, lng)
  273. if is_success:
  274. device.province = province
  275. device.city = city
  276. device.district = district
  277. base_info_obj.save()
  278. device.uptime = uptime_tp
  279. device.save()
  280. return Response({"code": 0, "msg": "success"})
  281. if cmd:
  282. ext = request_data.get("ext")
  283. imei = ext.get("imei")
  284. device_info = MongoDevice.objects.get(device_id=imei)
  285. if cmd == "online":
  286. device_info.device_status = 1
  287. if cmd == "offline":
  288. device_info.device_status = 0
  289. device_info.save()
  290. except Exception as e:
  291. logger.error(f"气象站设备 {device_id} 处理上报数据或增加设备失败,错误原因:{e.args}")
  292. return Response({"code": 2, "msg": "failer"})
  293. class DeviceListAPIView(APIView):
  294. def post(self, request):
  295. # 设备列表
  296. request_data = request.data
  297. device_id = request_data.get("device_id")
  298. device_status = request_data.get("device_status")
  299. search = request_data.get("search")
  300. page_num = int(request_data.get("pagenum")) if request_data.get("pagenum") else 1
  301. page_size = int(request_data.get("pagesize")) if request_data.get("pagesize") else 10
  302. if device_id:
  303. queryset = MongoDevice.objects.filter(device_id=device_id).order_by("-uptime")
  304. elif device_status:
  305. queryset = MongoDevice.objects.filter(device_status=device_status).order_by("-uptime")
  306. elif search:
  307. queryset = MongoDevice.objects.filter(Q(device_name__icontains=search) | Q(device_id__icontains=search))
  308. else:
  309. queryset = MongoDevice.objects.all().order_by("-uptime")
  310. total_obj = queryset.count()
  311. paginator = Paginator(queryset, page_size)
  312. page_obj = paginator.get_page(page_num)
  313. serializers = DeviceSerializers(page_obj, many=True)
  314. return Response({"code": 0, "msg": "success", "data": serializers.data, "count": total_obj})
  315. class DeviceChangeAPIView(APIView):
  316. def post(self, request):
  317. # 修改设备信息
  318. request_data = request.data
  319. device_name = request_data.get("device_name")
  320. device_id = request_data.get("device_id")
  321. lng = request_data.get("lng")
  322. lat = request_data.get("lat")
  323. device = MongoDevice.objects.get(device_id=device_id)
  324. if device_name:
  325. device.device_name = device_name
  326. if lng and lat:
  327. device.lng = lng
  328. device.lat = lat
  329. is_success, province, city, district = get_addr_by_lag_lng(lat, lng)
  330. if is_success:
  331. # 更新地理位置坐标
  332. device.province = province
  333. device.city = city
  334. device.district = district
  335. device.save()
  336. return Response({"code": 0, "msg": "success"})
  337. class DeviceListInfoAPIView(APIView):
  338. def post(self, request):
  339. # 设备列表
  340. request_data = request.data
  341. device_id = request_data.get("device_id")
  342. device_status = request_data.get("device_status")
  343. search = request_data.get("search")
  344. page_num = int(request_data.get("pagenum")) if request_data.get("pagenum") else 1
  345. page_size = int(request_data.get("pagesize")) if request_data.get("pagesize") else 10
  346. if device_id:
  347. queryset = MongoDevice.objects.filter(device_id=device_id).order_by("-uptime")
  348. elif device_status:
  349. queryset = MongoDevice.objects.filter(device_status=device_status).order_by("-uptime")
  350. elif search:
  351. queryset = MongoDevice.objects.filter(Q(device_name__icontains=search) | Q(device_id__icontains=search))
  352. else:
  353. queryset = MongoDevice.objects.all().order_by("-uptime")
  354. total_obj = queryset.count()
  355. paginator = Paginator(queryset, page_size)
  356. page_obj = paginator.get_page(page_num)
  357. serializers = DeviceSerializers(page_obj, many=True)
  358. return Response({"code": 0, "msg": "success", "data": serializers.data, "count": total_obj})
  359. class DeviceListAPIView(APIView):
  360. def post(self, request):
  361. queryset = MongoDevice.objects.order_by('-id')
  362. type_dict = {d.id: d.type_name for d in MongoDeviceType.objects.all()}
  363. result = []
  364. offline_list = []
  365. type_counter = Counter()
  366. device_dict = defaultdict(list)
  367. for item in queryset:
  368. device_info = item
  369. device_type_id = device_info.device_type_id
  370. device_status = device_info.device_status
  371. is_offline = False if device_status == 1 else True
  372. device_id = device_info.device_id
  373. if is_offline:
  374. offline_list.append(device_id)
  375. type_counter[device_type_id] += 1
  376. device_dict[device_type_id].append(device_id)
  377. coordinates = ""
  378. lng = device_info.lng
  379. lat = device_info.lat
  380. if not (lng and lat):
  381. continue
  382. coordinates = f"[{str(float(lng))},{str(float(lat))}]"
  383. device_name = device_info.device_name
  384. result.append({
  385. 'ld_id': item.id,
  386. 'device_id': device_id,
  387. 'device_name': device_name or device_id,
  388. 'tpye_name': type_dict[device_info.device_type_id],
  389. 'device_type_id': device_type_id,
  390. 'coordinates': coordinates,
  391. 'offline': False if device_status == 1 else True,
  392. 'is_warning': False
  393. })
  394. statistic_list = []
  395. for k, v in type_counter.items():
  396. statistic_list.append({
  397. 'type_id': k,
  398. 'type_count': v,
  399. 'type_name': type_dict[k]
  400. })
  401. warning_model_dict = {
  402. 3: DevicePestWarning
  403. }
  404. warning_list = []
  405. for k, v in device_dict.items():
  406. try:
  407. model_obj = warning_model_dict[k]
  408. except KeyError as e:
  409. continue
  410. war_list = [d.device_id for d in model_obj.objects.filter(device_id__in=v, status=0)]
  411. if war_list:
  412. warning_list.extend(war_list)
  413. for item in result:
  414. device_id = item['device_id']
  415. if device_id in warning_list:
  416. item['is_warning'] = True
  417. warn_and_off = set(warning_list) & set(offline_list)
  418. data = {
  419. "statistic": statistic_list,
  420. 'offline': {
  421. 'count': len(offline_list),
  422. 'result': offline_list
  423. },
  424. "waring": {
  425. 'count': len(warning_list),
  426. 'result': warning_list
  427. },
  428. "warn_and_off": {
  429. 'count': len(warn_and_off),
  430. 'result': warning_list
  431. },
  432. "data": result
  433. }
  434. return Response({"code": 0, "message": "success", "data": data})