summaryrefslogtreecommitdiff
path: root/mqtt_log.py
diff options
context:
space:
mode:
authorvlapa <vlapa@ya.ru>2026-06-17 03:07:13 +0300
committervlapa <vlapa@ya.ru>2026-06-17 03:07:13 +0300
commitc28d7285227cb4a5fafff0b8aa4744e9a26ef865 (patch)
tree7b0e68cd727aed22366b16542987e621a48a1dd4 /mqtt_log.py
Firstmain
Diffstat (limited to 'mqtt_log.py')
-rw-r--r--mqtt_log.py120
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()
+