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()