import paho.mqtt.client as mqtt import json # The callback when the client receives a CONNACK response from the broker def on_connect(client, userdata, flags, rc): print("[LOG] Connected with result code "+str(rc)) client.subscribe("read/file") if rc != 0: client.reconnect() # The callback when a PUBLISH message is received from the broker def on_message(client, userdata, msg): if msg.payload.decode() == "done": print("[LOG] The file processing is " + msg.payload.decode()) client.disconnect() # 发送文件 def send_file(initiator, topic, filename): try: with open(filename, 'r') as file: for line in file: radarData = json.loads(line) #initiator.publish(topic, payload=json.dumps(radarData), qos=2, retain=False) initiator.publish(topic, line) print(f"Sent: {radarData}") except Exception as e: print(f"发生了异常: {e}") def main(): import argparse parser = argparse.ArgumentParser() parser.add_argument("--ip", help="ip address of the broker", type=str, default="localhost") parser.add_argument("--port", help="broker port", type=int, default=1883) parser.add_argument("--file", help="radar filename", type=str, default="trackJson.txt") args = parser.parse_args() initiator = mqtt.Client(parser.prog) initiator.connect(args.ip, args.port) initiator.on_connect = on_connect initiator.on_message = on_message #initiator.publish("read/file", "start") # 开始发送文件 send_file(initiator, 'RadarRecord', args.file) initiator.loop_forever() if __name__ == "__main__": main()