diff options
Diffstat (limited to 'mqtt_log.py')
| -rw-r--r-- | mqtt_log.py | 120 |
1 files changed, 120 insertions, 0 deletions
diff --git a/mqtt_log.py b/mqtt_log.py new file mode 100644 index 0000000..cac4a01 --- /dev/null +++ b/mqtt_log.py @@ -0,0 +1,120 @@ +import os +import json +import time +from pathlib import Path +import requests +import paho.mqtt.client as mqtt +from dotenv import load_dotenv + +load_dotenv() + +def file_print(data_string, data_data): + """ Обработчик данных от сервера MQTT. Передаем: топик, данные """ + try: + # Извлекаем имя файла из топика + if "/" in data_string: + file_name = data_string[1:data_string.rfind("/")] + # Получаем расширение или вторую часть + value_part = data_string[data_string.rfind("/") + 1:] + else: + file_name = data_string + value_part = "" + + # Формируем пути + csv_path = f"{file_name}.csv" + timestamp = time.strftime('%Y/%m/%d - %H:%M:%S', time.localtime()) + formatted_data = f"/{timestamp}_{data_data}" + + # Запись в файл + fileWrite(os.getenv('FOLDER_ARHIVE_NAME'), csv_path, formatted_data) + + # Обновление кэша JSON + update_json_cache(file_name, f"{data_data}_{time.strftime('%H:%M:%S_%d/%m/%Y', time.localtime())}") + + except Exception as e: + print(f"Ошибка в file_print: {e}") + +def update_json_cache(file_name, value): + """Обновление JSON кэша""" + path = os.getenv('FILE_CACHE_JSON_NAME') + data_json = {file_name: value} + + try: + if Path(path).exists(): + with open(path, 'r', encoding='utf-8') as json_file: + data = json.load(json_file) + data.update(data_json) + else: + data = data_json + + with open(path, 'w', encoding='utf-8') as f: + json.dump(data, f, ensure_ascii=False, indent=2) + except Exception as e: + print(f"Ошибка при обновлении JSON: {e}") + +def connect_mqtt(): + """ Подключение к серверу MQTT """ + try: + # Попытка для новых версий + client = mqtt.Client(mqtt.CallbackAPIVersion.VERSION2) + except AttributeError: + # Для старых версий + client = mqtt.Client() + + client.username_pw_set(os.getenv('USERNAME'), os.getenv('PASSWORD')) + client.connect( + os.getenv('SERVER_IP'), + int(os.getenv('PORT_MQTT')) + ) + return client + +def subscribe(client: mqtt): + """ Подписываемся на топики """ + def on_message(client, userdata, msg): + try: + file_print(msg.topic, msg.payload.decode()) + except Exception as error: + print(f"Ошибка в on_message: {error}") + + client.subscribe(os.getenv('TOPIC')) + client.on_message = on_message + +def fileWrite(folderArhiveName, path, myFile): + """ Запись в файл. Передаем: название папки архива, название файла (топика), данные """ + try: + # Извлекаем имя папки + space_pos = myFile.find(' ') + if space_pos == -1: + folder_name = folderArhiveName + myFile + else: + folder_name = folderArhiveName + myFile[:space_pos] + + # Создаем путь + path_temp = Path.cwd() / folder_name + path_temp.mkdir(parents=True, exist_ok=True) + + # Записываем данные + file_path = path_temp / path + with open(file_path, "a", encoding='utf-8') as f: + # Извлекаем данные после пробела + data_part = myFile[myFile.rfind(' ') + 1:] if ' ' in myFile else myFile + f.write(data_part + '\n') + + except Exception as e: + print(f"Ошибка записи файла: {e}") + +def run(): + try: + print("Подключение к MQTT...") + client = connect_mqtt() + subscribe(client) + print("Ожидание сообщений...") + client.loop_forever() + except KeyboardInterrupt: + print("\nПрограмма остановлена") + except Exception as e: + print(f"Критическая ошибка: {e}") + +if __name__ == "__main__": + run() + |
