# -*- coding: utf-8 -*- #!/usr/bin/python # -*- coding: utf-8 -* import paho.mqtt.client as mqtt import json import time import threading import db_operate from linkListUtil import Node, LinkedList from shapelyUtil import ShapelyUtil from queue import Queue def gettime(): time1=time.strftime("%Y-%m-%d %H:%M:%S",time.localtime()) return time1 # 服务器地址 host = '7.65.0.207' #'10.5.226.22 7.65.0.207' # 通信端口 默认端口1883 port = 1883 username = 'eadpoint' password = '321123' objs = LinkedList() # 创建数据队列 data_queue = Queue() #[[12146938.804545141,3106558.9023983314],[12147505.698887747,3106620.0104321464],[12147188.527927307,3108004.3420571982],[12146593.188072586,3107932.8914925163]] #经度:121.342248, 纬度:30.814109 经度:121.335855, 纬度:30.843758 经度:121.374132, 纬度:30.849655, 经度:121.38213, 纬度:30.819001 shapeUtil = ShapelyUtil([[12134224.8,3081410.9],[12133585.5,3084375.8],[12137413.2,3084965.5],[12138213,3081900.1]]) arr_topic = ['uav_zte_radar'] send_topic = ['uavAddObj','uavAddTrack','uavRemoveObj'] # 连接后事件 def on_connect(client, userdata, flags, respons_code): if respons_code == 0: print('Connection Succeed!') else: print('Connect Error status {0}'.format(respons_code)) for x in arr_topic: client.subscribe(x) # 发送文件 def send_file(client, topic, filename): try: with open(filename, 'r') as file: for line in file: radarData = json.loads(line) #client.publish(topic, payload=json.dumps(radarData), qos=2, retain=False) client.publish(topic, line) print(f"Sent: {radarData}") except Exception as e: print(f"发生了异常: {e}") # 接收到数据后事件 def on_message(client, userdata, msg): global dddd # 打印订阅消息主题 #print(msg.payload) try: jsondata=json.loads(msg.payload) #print(jsondata) if jsondata['object'] is None: print('has no object data') else : #if jsondata['cellid'] == "11" : for x in jsondata['object']: x['time'] = time.strftime("%Y-%m-%d %H:%M:%S",time.localtime(x['ts'])) x['recvtime'] = int(time.time()) x['minOfYear'] = jsondata['minOfYear'] x['second'] = jsondata['second'] x['direction'] = 0 x['fusDevID'] = jsondata['fusDevID'] x['fps'] = jsondata['fps'] x['cellid'] = jsondata['cellid'] x['msgCnt'] = jsondata['msgCnt'] x['nofobjects'] = jsondata['nofobjects'] x['dataSource'] = jsondata['dataSource'] x['refPos'] = jsondata['refPosList'][0] x['isInside'] = 1 if shapeUtil.isInPolygon_(x['lng'],x['lat']) else 0 node = objs.find_by_value(x['object_id']) if node is None: objs.insert_value_to_head(x['object_id'], [x]) node = objs.find_by_value(x['object_id']) if x['isInside'] == 1: client.publish(send_topic[0], json.dumps(x)) data_queue.put(node.track) else : obj0 = node.track[0] x['direction'] = shapeUtil.angle(obj0['lat'], obj0['lng'], x['lat'], x['lng']) if x['isInside'] == 1: client.publish(send_topic[1], json.dumps(x)) else : client.publish(send_topic[2],x['object_id']) objs.delete_by_value(x['object_id'], node.track) data_queue.put(node.track) node.track.append(x) except Exception as e: print(f"发生了异常: {e}") def monitor_list(linked_list: LinkedList, client: mqtt.Client): p = linked_list.head countBoat = 0 while True: if p is None: #print(f"当前区域内船舶数量:{countBoat}: ") p = linked_list.head countBoat = 0 continue try: item = p.item track = p.track p = p.next if abs(track[-1]['recvtime'] - int(time.time())) > 5 : client.publish(send_topic[2],item) track[-1]['isInside'] = -1 linked_list.delete_by_value(item, track) data_queue.put(track) #else: #if track[-1]['isInside'] == 1 and len(track)%20 == 1: #print(f"更新状态:{item}") #data_queue.put(track) countBoat = countBoat+1 if track[-1]['isInside'] == 1 else countBoat except Exception as e: print(f"监控异常: {e}") def queueListen(data_queue): while True: data = data_queue.get() if data is None: continue #print(f"无人机: {data[-1]['object_id']} 状态: {data[-1]['isInside']}") db_operate.uav_zte_dt_save(data) def main(): db_operate.uav_zte_old_dt_clean() #队列监听 db_thread = threading.Thread(target=queueListen, args=(data_queue,)) db_thread.start() client = mqtt.Client() # 注册事件 client.on_connect = on_connect client.on_message = on_message # 设置账号密码(如果需要的话) client.username_pw_set(username, password=password) # 连接到服务器 client.connect(host, port=port, keepalive=60) list_thread = threading.Thread(target=monitor_list, args=(objs, client)) list_thread.start() # 守护连接状态 client.loop_forever() if __name__ == '__main__': #print(time.strftime("%Y-%m-%d %H:%M:%S",time.localtime(1712900732))) main()