1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
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()
|