Skip to content
GitLab
Projects
Groups
Snippets
/
Help
Help
Support
Community forum
Keyboard shortcuts
?
Submit feedback
Sign in
Toggle navigation
Menu
Open sidebar
co2ampel2
MQTT_TO_INFLUXDB
Commits
3ce51a0e
Commit
3ce51a0e
authored
Apr 17, 2025
by
Patrick Ade
Browse files
created jsonhandler (untested), refactoring
parent
57c9cb62
Changes
5
Hide whitespace changes
Inline
Side-by-side
src/mqtt_influx_backend/
I
nfluxDBWriter.py
→
src/mqtt_influx_backend/
i
nfluxDBWriter.py
View file @
3ce51a0e
File moved
src/mqtt_influx_backend/jsonhandler.py
0 → 100644
View file @
3ce51a0e
import
json
file_name
=
"mac_to_room.json"
def
load_json
():
with
open
(
file_name
)
as
f
:
mac_room_mapping
=
json
.
load
(
f
)
return
mac_room_mapping
def
write_json
(
mac_room_mapping
:
dict
):
with
open
(
file_name
,
"w"
)
as
f
:
f
.
seek
(
0
)
json
.
dump
(
mac_room_mapping
,
f
,
indent
=
4
)
f
.
truncate
()
# TODO Check if truncate is necessary?
src/mqtt_influx_backend/
M
QTTClientHandler.py
→
src/mqtt_influx_backend/
m
QTTClientHandler.py
View file @
3ce51a0e
import
json
from
datetime
import
datetime
import
paho.mqtt.client
as
mqtt
from
mqtt_influx_backend
import
I
nfluxDBWriter
from
src.
mqtt_influx_backend
import
i
nfluxDBWriter
class
MQTTClientHandler
:
class
MQTTClientHandler
:
# Konstruktor
def
__init__
(
self
,
broker_url
:
str
,
topic
:
str
,
influx_writer
:
InfluxDBWriter
):
#logger sollte hier zugeweist werden
def
__init__
(
self
,
broker_url
:
str
,
topic
:
str
,
influx_writer
:
influxDBWriter
):
# logger sollte hier zugeweist werden
self
.
broker_url
=
broker_url
self
.
topic
=
topic
self
.
influx_writer
=
influx_writer
self
.
client
=
mqtt
.
Client
()
# Events werden hier methoden
# Events werden hier methoden
self
.
client
.
on_connect
=
self
.
on_connect
self
.
client
.
on_message
=
self
.
on_message
...
...
@@ -20,25 +22,25 @@ class MQTTClientHandler:
# log
print
(
"Connected with result code "
+
str
(
rc
))
client
.
subscribe
(
self
.
topic
)
#def save_mapping(self):
#with open(self.mapping_file, "w") as f:
#
def save_mapping(self):
#
with open(self.mapping_file, "w") as f:
# json.dump(self.mac_to_room, f, indent=4)
# eventuell refactorn und die Aufgaben in Methoden aufteilen
def
on_message
(
self
,
client
,
userdata
,
msg
):
"""
Wenn das Topic eine Nachricht bekommt wird diese Methode ausgeführt
self: ist die MQTTClientHandler instanz, die wird gebraucht
self: ist die MQTTClientHandler instanz, die wird gebraucht
"""
# log
msg
=
json
.
loads
(
msg
.
payload
)
metadate
=
msg
[
"metadata"
]
# key: mac : value : room
# hier prüfen, ob die Mac-Adresse einen Raum hat,
# hier prüfen, ob die Mac-Adresse einen Raum hat,
# wenn nicht trage es in mac_to_room leer ein
# "aa:bb:cc:dd:ee:ff" : ""
mac
=
metadate
[
"mac"
]
mac
=
metadate
[
"mac"
]
# ToImplement
"""
...
...
@@ -49,23 +51,24 @@ class MQTTClientHandler:
room = self.mac_to_room[mac]
"""
try
:
self
.
influx_writer
.
write_point
(
measurement
=
"sensor_data"
,
tags
=
{
# ToDo "room": metadate["todo"],
# ToDo "room": metadate["todo"],
"mac"
:
metadate
[
"mac-address"
]
},
fields
=
{
"co2"
:
msg
[
"co2"
],
"temperature"
:
msg
[
"temp"
],
"humidity"
:
msg
[
"rh"
]
"humidity"
:
msg
[
"rh"
]
,
},
timestamp
=
metadate
[
"time"
]
timestamp
=
metadate
[
"time"
]
,
)
print
(
"Wrote to InfluxDB:"
,
msg
)
# muss später rausgeschmiessen werden
print
(
"Wrote to InfluxDB:"
,
msg
)
# muss später rausgeschmiessen werden
except
Exception
as
e
:
# log
print
(
"Error processing message:"
,
e
)
...
...
src/mqtt_influx_backend/mac_to_room
→
src/mqtt_influx_backend/mac_to_room
.json
View file @
3ce51a0e
File moved
src/mqtt_influx_backend/
M
ain.py
→
src/mqtt_influx_backend/
m
ain.py
View file @
3ce51a0e
from
dotenv
import
load_dotenv
import
os
import
logging
from
src.mqtt_influx_backend.MQTTClientHandler
import
MQTTClientHandler
from
src.mqtt_influx_backend.InfluxDBWriter
import
InfluxDBWriter
from
src.mqtt_influx_backend.mQTTClientHandler
import
MQTTClientHandler
from
src.mqtt_influx_backend.influxDBWriter
import
InfluxDBWriter
def
get_logger
(
name
:
str
=
"default_logger"
)
->
logging
.
Logger
:
logger
=
logging
.
getLogger
(
name
)
if
not
logger
.
handlers
:
# Verhindert doppelte Handler bei mehrfacher Initialisierung
if
(
not
logger
.
handlers
):
# Verhindert doppelte Handler bei mehrfacher Initialisierung
logger
.
setLevel
(
logging
.
DEBUG
)
handler
=
logging
.
StreamHandler
()
formatter
=
logging
.
Formatter
(
'
[%(asctime)s] [%(levelname)s] [%(name)s] %(message)s
'
"
[%(asctime)s] [%(levelname)s] [%(name)s] %(message)s
"
)
handler
.
setFormatter
(
formatter
)
logger
.
addHandler
(
handler
)
...
...
@@ -19,21 +22,23 @@ def get_logger(name: str = "default_logger") -> logging.Logger:
load_dotenv
()
def
main
():
influx_writer
=
InfluxDBWriter
(
url
=
os
.
getenv
(
"INFLUXDB_URL"
),
token
=
os
.
getenv
(
"INFLUXDB_TOKEN"
),
org
=
os
.
getenv
(
"INFLUXDB_ORG"
),
bucket
=
os
.
getenv
(
"INFLUXDB_BUCKET"
)
url
=
os
.
getenv
(
"INFLUXDB_URL"
),
token
=
os
.
getenv
(
"INFLUXDB_TOKEN"
),
org
=
os
.
getenv
(
"INFLUXDB_ORG"
),
bucket
=
os
.
getenv
(
"INFLUXDB_BUCKET"
)
,
)
mqtt_handler
=
MQTTClientHandler
(
broker_url
=
os
.
getenv
(
"MQTT_BROKER_URL"
),
topic
=
os
.
getenv
(
"MQTT_TOPIC"
),
influx_writer
=
influx_writer
broker_url
=
os
.
getenv
(
"MQTT_BROKER_URL"
),
topic
=
os
.
getenv
(
"MQTT_TOPIC"
),
influx_writer
=
influx_writer
,
)
mqtt_handler
.
start
()
if
__name__
==
"__main__"
:
main
()
Write
Preview
Supports
Markdown
0%
Try again
or
attach a new file
.
Cancel
You are about to add
0
people
to the discussion. Proceed with caution.
Finish editing this message first!
Cancel
Please
register
or
sign in
to comment