Commit da0dc0bb authored by Gezer's avatar Gezer
Browse files

Refactor backend and stream processing structure

- Updated Dockerfile in backend to set PYTHONPATH and modify CMD for Django server.
- Changed permissions for multiple Python files in backend to make them executable.
- Adjusted import paths in influxdb_service.py and views.py to reflect new structure.
- Removed unused utils and loggingFactory files from stream_processing.
- Added new utils for logging and flux query building.
- Enhanced error handling and logging in mQTTClientHandler.py.
- Updated docker-compose.yaml to include utils directory for both backend and stream processing services.
- Improved JSON handling in jsonhandler.py with better exception management.
parent 7ba8a99b
File mode changed from 100644 to 100755
File mode changed from 100644 to 100755
File mode changed from 100644 to 100755
File mode changed from 100644 to 100755
File mode changed from 100644 to 100755
File mode changed from 100644 to 100755
...@@ -9,6 +9,7 @@ services: ...@@ -9,6 +9,7 @@ services:
container_name: stream-processing container_name: stream-processing
volumes: volumes:
- ./stream_processing/mac_to_room.json:/app/mac_to_room.json - ./stream_processing/mac_to_room.json:/app/mac_to_room.json
- ./utils:/app/utils
env_file: env_file:
- ./backend/.env.docker - ./backend/.env.docker
restart: unless-stopped restart: unless-stopped
...@@ -28,6 +29,7 @@ services: ...@@ -28,6 +29,7 @@ services:
- "8000:8000" - "8000:8000"
volumes: volumes:
- ./backend:/app - ./backend:/app
- ./utils:/app/utils
- /app/.venv - /app/.venv
frontend: frontend:
......
#Stream-Processing
FROM python:latest FROM python:latest
# Setze das Arbeitsverzeichnis im Container # Setze das Arbeitsverzeichnis im Container
WORKDIR /app WORKDIR /app
ENV PYTHONPATH="${PYTHONPATH}:/app:/app/utils"
# Kopiere die pyproject.toml und uv.lock aus dem Backend-Hauptverzeichnis # Kopiere die pyproject.toml und uv.lock aus dem Backend-Hauptverzeichnis
COPY . . COPY . .
......
import json import json
import os from typing import Dict
def load_json(file_name: str) -> dict: def load_json(file_name: str) -> dict:
""" """
ladet eine JSON Datei, wenn diese existiert, Lädt eine JSON-Datei und gibt deren Inhalt als Dictionary zurück.
und gibt diese als dictionary zurück
key : value Falls die Datei nicht existiert oder nicht lesbar ist, wird ein leeres Dictionary zurückgegeben.
:param file_name: Pfad zur JSON-Datei.
:return: Dictionary mit dem Inhalt der JSON-Datei.
""" """
print("JSONHANDLER") try:
if not os.path.exists(file_name): with open(file_name, "r", encoding="utf-8") as f:
print("FILE IST OFFEN") return json.load(f)
except FileNotFoundError:
return {} return {}
with open(file_name) as f: except json.JSONDecodeError as e:
mac_room_mapping = json.load(f) raise ValueError(f"Fehler beim Parsen von JSON: {e}") from e
return mac_room_mapping
def write_json(mac_room_mapping: dict, file_name: str): def write_json(mac_room_mapping: Dict[str, str], file_name: str) -> None:
print("ES SCHREIBT IN JASON") """
Schreibt ein Dictionary als JSON-Datei an den angegebenen Pfad.
:param mac_room_mapping: Das zu speichernde Dictionary.
:param file_name: Pfad zur Zieldatei.
"""
try: try:
print(f"→ Schreibe Datei: {file_name}") with open(file_name, "w", encoding="utf-8") as f:
with open(file_name, "w") as f:
print("FILE IST OPEN")
json.dump(mac_room_mapping, f, indent=4) json.dump(mac_room_mapping, f, indent=4)
print("✅ Datei erfolgreich geschrieben")
except Exception as e: except Exception as e:
print("❌ Fehler beim Schreiben:", e) raise IOError(f"Fehler beim Schreiben der Datei '{file_name}': {e}") from e
import json import json
import os
import paho.mqtt.client as mqtt
from utils.loggingFactory import LoggerFactory from utils.loggingFactory import LoggerFactory
from datetime import datetime
import paho.mqtt.client as mqtt
import jsonhandler
from utils.influx import InfluxDBHelper from utils.influx import InfluxDBHelper
import os import jsonhandler
class MQTTClientHandler: class MQTTClientHandler:
MAPPING_FILE_NAME = os.path.join( MAPPING_FILE_NAME = os.path.join(
os.path.dirname(os.path.abspath(__file__)), os.path.dirname(os.path.abspath(__file__)),
"mac_to_room.json" "mac_to_room.json"
) )
MEASUREMENT_NAME = "sensor_data" def __init__(self, broker_url: str, topic: str, influx_writer: InfluxDBHelper):
TAG_ROOM = "room"
TAG_MAC = "mac"
FIELD_CO2 = "co2"
FIELD_TEMP = "temperature"
FIELD_HUMIDITY = "humidity"
# Konstruktor
def __init__(
self, broker_url: str, topic: str, influx_writer: InfluxDBHelper
):
print("DAS IST EIN TEST")
self.logger = LoggerFactory.get_logger(__name__) self.logger = LoggerFactory.get_logger(__name__)
# key: mac : value : room
self.mac_to_room = jsonhandler.load_json(self.MAPPING_FILE_NAME) self.mac_to_room = jsonhandler.load_json(self.MAPPING_FILE_NAME)
self.broker_url = broker_url self.broker_url = broker_url
self.topic = topic self.topic = topic
self.influx_writer = influx_writer self.influx_writer = influx_writer
self.client = mqtt.Client() self.client = mqtt.Client()
# Methoden werden hier Events zugeteilt
self.client.on_connect = self.on_connect self.client.on_connect = self.on_connect
self.client.on_message = self.on_message self.client.on_message = self.on_message
def on_connect(self, client, userdata, flags, rc): def on_connect(self, client, userdata, flags, rc):
self.logger.info("Connected with result code " + str(rc)) if rc == 0:
print("Connected") self.logger.info("MQTT-Verbindung erfolgreich.")
client.subscribe(self.topic) client.subscribe(self.topic)
self.logger.info("Subscribed to " + self.topic) self.logger.info(f"Abonniert: {self.topic}")
else:
self.logger.error(f"MQTT-Verbindung fehlgeschlagen mit Code {rc}")
# eventuell refactorn und die Aufgaben in Methoden aufteilen
def on_message(self, client, userdata, msg): def on_message(self, client, userdata, msg):
""" """
Wenn das Topic eine Nachricht bekommt wird diese Methode ausgeführt Wird aufgerufen, wenn eine Nachricht empfangen wird.
self: ist die MQTTClientHandler instanz, die wird gebraucht um die Einträge in
die InfluxDB zu schreiben
""" """
try:
payload = json.loads(msg.payload)
metadata = payload.get("metadata", {})
mac = metadata.get("mac-address")
print("Message") if not mac:
self.logger.warning("MAC-Adresse fehlt in empfangener Nachricht.")
msg = json.loads(msg.payload) return
metadate = msg["metadata"]
# hier prüfen, ob die Mac-Adresse einen Raum hat, if mac not in self.mac_to_room:
# wenn nicht trage es in mac_to_room leer ein self._handle_unknown_mac(mac)
# "aa:bb:cc:dd:ee:ff" : "" return
mac = metadate["mac-address"]
if mac not in self.mac_to_room: self._write_to_influxdb(payload, metadata)
self.logger.warning(
f"Neue MAC-Adresse gefunden: {mac}. Mapping wird ergänzt."
)
print("MAC war nicht in File")
self.mac_to_room[mac] = "" # leerer Platzhalter
jsonhandler.write_json(self.mac_to_room, self.MAPPING_FILE_NAME)
self.mac_to_room = jsonhandler.load_json(self.MAPPING_FILE_NAME)
print("keine Ahnung")
return
self.write_to_influxDB(msg, metadate) except json.JSONDecodeError as e:
self.logger.error(f"Fehler beim Parsen der Nachricht: {e}")
except Exception as e:
self.logger.exception(f"Unerwarteter Fehler beim Verarbeiten der Nachricht: {e}")
def _handle_unknown_mac(self, mac: str):
"""
Fügt eine neue MAC-Adresse zum Mapping hinzu und speichert das JSON neu.
"""
self.logger.warning(f"Neue MAC-Adresse entdeckt: {mac} – Mapping wird ergänzt.")
self.mac_to_room[mac] = ""
jsonhandler.write_json(self.mac_to_room, self.MAPPING_FILE_NAME)
self.mac_to_room = jsonhandler.load_json(self.MAPPING_FILE_NAME)
def write_to_influxDB(self, msg: dict, metadate: dict): def _write_to_influxdb(self, payload: dict, metadata: dict):
"""
Schreibt die Sensorwerte in die InfluxDB.
"""
try: try:
print(msg)
self.influx_writer.write_point( self.influx_writer.write_point(
measurement=self.MEASUREMENT_NAME, measurement=self.influx_writer.MEASUREMENT_NAME,
tags={ tags={
self.TAG_ROOM: self.mac_to_room[metadate["mac-address"]], self.influx_writer.TAG_ROOM: self.mac_to_room[metadata["mac-address"]],
self.TAG_MAC: metadate["mac-address"], self.influx_writer.TAG_MAC: metadata["mac-address"],
}, },
fields={ fields={
self.FIELD_CO2: msg["co2"], self.influx_writer.FIELD_CO2: payload["co2"],
self.FIELD_TEMP: msg["temp"], self.influx_writer.FIELD_TEMP: payload["temp"],
self.FIELD_HUMIDITY: msg["rh"], self.influx_writer.FIELD_HUMIDITY: payload["rh"],
}, },
timestamp=metadate["time"], # fix timestamp=metadata["time"],
)
self.logger.info(
f"Wrote to InfluxDB: {msg}"
) )
self.logger.info(f"Daten geschrieben für MAC: {metadata['mac-address']}")
except Exception as e: except Exception as e:
self.logger.error(f"Failed writing to InfluxDb: {e}") self.logger.error(f"Fehler beim Schreiben in InfluxDB: {e}")
def start(self): def start(self):
self.client.connect(self.broker_url) """
self.client.loop_forever() Startet die Verbindung zum MQTT-Broker.
"""
try:
self.client.connect(self.broker_url)
self.client.loop_forever()
except Exception as e:
self.logger.critical(f"Fehler beim Starten des MQTT-Clients: {e}")
from dotenv import load_dotenv
import os import os
from dotenv import load_dotenv
from mQTTClientHandler import MQTTClientHandler from mQTTClientHandler import MQTTClientHandler
from utils.influx import InfluxDBHelper from utils.influx import InfluxDBHelper
from utils.loggingFactory import LoggerFactory
load_dotenv() load_dotenv()
......
from typing import List, Optional
class FluxQueryBuilder:
"""
A builder class for constructing Flux queries (InfluxDB 2.x).
Supports fluent method chaining for readable query construction.
Example:
query = (
FluxQueryBuilder()
.bucket("sensor_data")
.time_range("-1h", "now()")
.filter_measurement("sensor_data")
.filter_fields("co2", "temperature", "humidity")
.filter_field("room", "1/210")
.pivot()
.mean()
.build()
)
"""
def __init__(self) -> None:
self._bucket: Optional[str] = None
self._start: Optional[str] = None
self._stop: Optional[str] = None
self._measurement: Optional[str] = None
self._fields: List[str] = []
self._field_filters: dict = {} # e.g. {"room": "1/210"}
self._use_mean: bool = False
self._use_pivot: bool = False
def bucket(self, name: str) -> "FluxQueryBuilder":
"""Set the InfluxDB bucket name."""
self._bucket = name
return self
def time_range(self, start: str, stop: str) -> "FluxQueryBuilder":
"""Set time range for the query."""
self._start = start
self._stop = stop
return self
def filter_measurement(self, measurement: str) -> "FluxQueryBuilder":
"""Filter for a specific _measurement value."""
self._measurement = measurement
return self
def filter_field(self, field: str, value: str) -> "FluxQueryBuilder":
"""Add a tag filter like r["room"] == "1/210"."""
self._field_filters[field] = value
return self
def filter_fields(self, *fields: str) -> "FluxQueryBuilder":
"""Filter for multiple _field values using OR."""
self._fields.extend(fields)
return self
def mean(self) -> "FluxQueryBuilder":
"""Apply mean() aggregation."""
self._use_mean = True
return self
def pivot(self) -> "FluxQueryBuilder":
"""Apply pivot() to restructure results with multiple fields per timestamp."""
self._use_pivot = True
return self
def build(self) -> str:
"""Construct and return the final Flux query string."""
if not self._bucket:
raise ValueError("Bucket name is required.")
if not (self._start and self._stop):
raise ValueError("Start and stop times are required.")
if not self._measurement:
raise ValueError("Measurement is required.")
lines: List[str] = []
lines.append(f'from(bucket: "{self._bucket}")')
lines.append(f' |> range(start: {self._start}, stop: {self._stop})')
lines.append(f' |> filter(fn: (r) => r["_measurement"] == "{self._measurement}")')
# Optional tag filters (e.g., room or mac)
for key, value in self._field_filters.items():
lines.append(f' |> filter(fn: (r) => r["{key}"] == "{value}")')
# Optional field filters (_field == ...)
if self._fields:
or_expr = " or ".join(f'r["_field"] == "{f}"' for f in self._fields)
lines.append(f' |> filter(fn: (r) => {or_expr})')
if self._use_mean:
lines.append(' |> mean()')
if self._use_pivot:
lines.append(' |> pivot(rowKey:["_time"], columnKey: ["_field"], valueColumn: "_value")')
return "\n".join(lines)
def reset(self) -> "FluxQueryBuilder":
"""Reset the builder so a new query can be built from scratch."""
self.__init__()
return self
import os
from influxdb_client import InfluxDBClient, Point, WritePrecision
from influxdb_client.client.write_api import WriteOptions
class InfluxDBHelper:
def __init__(self, url: str, token: str, org: str, bucket: str):
self.client = InfluxDBClient(url=url, token=token, org=org)
self.bucket = bucket
self.org = org
self.query_api = self.client.query_api()
# self.write_api = self.client.write_api(write_options=SYNCHRONOUS) good for debug
self.write_api = self.client.write_api(
write_options=WriteOptions(batch_size=1000, flush_interval=10000)
)
def write_point(
self, measurement: str, tags: dict, fields: dict, timestamp=None
):
""" """
point = Point(measurement)
for k, v in tags.items():
point.tag(k, v)
for k, v in fields.items():
point.field(k, v)
if timestamp:
point.time(timestamp, WritePrecision.NS)
self.write_api.write(bucket=self.bucket, org=self.org, record=point)
def ping(self) -> bool:
return self.client.ping()
def get_all_data(self):
""" """
query = f'''
from(bucket: "{self.bucket}")
|> range(start: -20d)
|> filter(fn: (r) => r["_measurement"] == "sensor_data")
'''
return self.query_api.query(org=self.org, query=query)
def get_latest_room_data(self, room_id: str):
""" """
query = f'''
from(bucket: "{self.bucket}")
|> range(start: -5m)
|> filter(fn: (r) => r["_measurement"] == "co2")
|> filter(fn: (r) => r["room"] == "{room_id}")
|> last()
'''
return self.query_api.query(org=self.org, query=query)
import logging
import os
from logging.handlers import RotatingFileHandler
LOG_DIR = "logs"
LOG_FILE = "app.log"
LOG_PATH = os.path.join(LOG_DIR, LOG_FILE)
class LoggerFactory:
# logger.info("Connected with result code %s", str(rc))
# logger.warning("Neue MAC-Adresse gefunden: %s", mac)
# logger.error("Failed writing to InfluxDb: %s", e)
@staticmethod
def get_logger(name: str, level=logging.DEBUG) -> logging.Logger:
if not os.path.exists(LOG_DIR):
os.makedirs(LOG_DIR)
logger = logging.getLogger(name)
if logger.hasHandlers():
return logger # vermeidet doppelte Handler
logger.setLevel(level)
formatter = logging.Formatter(
"[%(asctime)s] %(levelname)s in %(name)s: %(message)s",
datefmt="%Y-%m-%d %H:%M:%S",
)
file_handler = RotatingFileHandler(
LOG_PATH, maxBytes=5_000_000, backupCount=5
)
file_handler.setFormatter(formatter)
logger.addHandler(file_handler)
return logger
Supports Markdown
0% or .
You are about to add 0 people to the discussion. Proceed with caution.
Finish editing this message first!
Please register or to comment