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
35827b83
Commit
35827b83
authored
Apr 21, 2025
by
Gezer
Browse files
cleaned project
parent
a168b48f
Changes
6
Hide whitespace changes
Inline
Side-by-side
build/lib/mqtt_influx_backend/InfluxDBWriter.py
deleted
100644 → 0
View file @
a168b48f
from
influxdb_client
import
InfluxDBClient
,
Point
,
WritePrecision
from
influxdb_client.client.write_api
import
SYNCHRONOUS
class
InfluxDBWriter
:
def
__init__
(
self
,
url
:
str
,
token
:
str
,
org
:
str
,
bucket
:
str
):
self
.
client
=
InfluxDBClient
(
url
=
url
,
token
=
token
,
org
=
org
)
self
.
write_api
=
self
.
client
.
write_api
(
write_options
=
SYNCHRONOUS
)
self
.
bucket
=
bucket
self
.
org
=
org
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
)
build/lib/mqtt_influx_backend/MQTTClientHandler.py
deleted
100644 → 0
View file @
a168b48f
import
json
from
datetime
import
datetime
import
paho.mqtt.client
as
mqtt
from
mqtt_influx_backend
import
InfluxDBWriter
class
MQTTClientHandler
:
def
__init__
(
self
,
broker_url
:
str
,
topic
:
str
,
influx_writer
:
InfluxDBWriter
):
self
.
broker_url
=
broker_url
self
.
topic
=
topic
self
.
influx_writer
=
influx_writer
self
.
client
=
mqtt
.
Client
()
self
.
client
.
on_connect
=
self
.
on_connect
self
.
client
.
on_message
=
self
.
on_message
def
on_connect
(
self
,
client
,
userdata
,
flags
,
rc
):
print
(
"Connected with result code "
+
str
(
rc
))
client
.
subscribe
(
self
.
topic
)
def
on_message
(
self
,
client
,
userdata
,
msg
):
utfMgs
=
msg
.
payload
.
decode
(
"utf-8"
)
#dictionary = json.loads(utfMgs)
print
(
utfMgs
)
print
(
utfMgs
)
print
(
utfMgs
)
data
=
utfMgs
try
:
value
=
float
(
data
.
get
(
"co2"
))
self
.
influx_writer
.
write_point
(
measurement
=
"temperature"
,
tags
=
{
"topic"
:
msg
.
topic
},
fields
=
{
"value"
:
value
},
timestamp
=
datetime
.
utcnow
()
)
print
(
"Wrote to InfluxDB:"
,
value
)
except
Exception
as
e
:
print
(
"Error processing message:"
,
e
)
def
start
(
self
):
self
.
client
.
connect
(
self
.
broker_url
)
self
.
client
.
loop_forever
()
build/lib/mqtt_influx_backend/Main.py
deleted
100644 → 0
View file @
a168b48f
from
dotenv
import
load_dotenv
import
os
from
mqtt_influx_backend.MQTTClientHandler
import
MQTTClientHandler
from
mqtt_influx_backend.InfluxDBWriter
import
InfluxDBWriter
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"
)
)
mqtt_handler
=
MQTTClientHandler
(
broker_url
=
os
.
getenv
(
"MQTT_BROKER_URL"
),
topic
=
os
.
getenv
(
"MQTT_TOPIC"
),
influx_writer
=
influx_writer
)
mqtt_handler
.
start
()
if
__name__
==
"__main__"
:
main
()
build/lib/mqtt_to_influx.py
deleted
100644 → 0
View file @
a168b48f
import
json
from
datetime
import
datetime
import
influxdb_client
,
os
,
time
from
influxdb_client
import
InfluxDBClient
,
Point
,
WritePrecision
from
influxdb_client.client.write_api
import
SYNCHRONOUS
# InfluxDB config
INFLUXDB_URL
=
os
.
environ
.
get
(
"INFLUXDB_URL"
)
INFLUXDB_TOKEN
=
os
.
environ
.
get
(
"INFLUXDB_TOKEN"
)
INFLUXDB_ORG
=
os
.
environ
.
get
(
"INFLUXDB_ORG"
)
INFLUXDB_BUCKET
=
os
.
environ
.
get
(
"INFLUXDB_BUCKET"
)
write_client
=
influxdb_client
.
InfluxDBClient
(
url
=
INFLUXDB_URL
,
token
=
INFLUXDB_TOKEN
,
org
=
INFLUXDB_ORG
)
bucket
=
"co2-test"
write_api
=
write_client
.
write_api
(
write_options
=
SYNCHRONOUS
)
# MQTT config
MQTT_BROKER_URL
=
os
.
environ
.
get
(
"MQTT_BROKER_URL"
)
MQTT_PUBLISH_TOPIC
=
os
.
environ
.
get
(
"MQTT_TOPIC"
)
for
value
in
range
(
5
):
point
=
(
Point
(
"measurement1"
)
.
tag
(
"mac-adress"
,
"22:de:aa:21::a2"
)
.
field
(
"ppm"
,
value
)
)
write_api
.
write
(
bucket
=
bucket
,
org
=
INFLUXDB_ORG
,
record
=
point
)
time
.
sleep
(
1
)
# separate points by 1 second
'''
def on_connect(client, userdata, flags, rc):
print("Connected with result code " + str(rc))
client.subscribe(MQTT_PUBLISH_TOPIC)
def on_message(client, userdata, msg):
print(msg.topic + " " + str(msg.payload))
try:
# Decode JSON payload (z. B. {"value": 22.5, "unit": "C"})
data = json.loads(msg.payload.decode())
value = float(data.get("value"))
# Create data point
point = Point("temperature")
\
.tag("topic", msg.topic)
\
.field("value", value)
\
.time(datetime.utcnow())
# Write to InfluxDB
write_api.write(bucket=INFLUXDB_BUCKET, org=INFLUXDB_ORG, record=point)
print("Wrote to InfluxDB:", point)
except Exception as e:
print("Error processing message:", e)
'''
# Start MQTT client
'''
mqttc = mqtt.Client()
mqttc.on_connect = on_connect
mqttc.on_message = on_message
mqttc.connect(MQTT_BROKER_URL)
mqttc.loop_forever()
'''
src/mqtt_influx_backend.egg-info/SOURCES.txt
View file @
35827b83
pyproject.toml
pyproject.toml
src/mqtt_influx_backend/InfluxDBWriter.py
src/mqtt_influx_backend/MQTTClientHandler.py
src/mqtt_influx_backend/Main.py
src/mqtt_influx_backend/__init__.py
src/mqtt_influx_backend/__init__.py
src/mqtt_influx_backend/influxDBWriter.py
src/mqtt_influx_backend/jsonhandler.py
src/mqtt_influx_backend/loggingFactory.py
src/mqtt_influx_backend/mQTTClientHandler.py
src/mqtt_influx_backend/main.py
src/mqtt_influx_backend.egg-info/PKG-INFO
src/mqtt_influx_backend.egg-info/PKG-INFO
src/mqtt_influx_backend.egg-info/SOURCES.txt
src/mqtt_influx_backend.egg-info/SOURCES.txt
src/mqtt_influx_backend.egg-info/dependency_links.txt
src/mqtt_influx_backend.egg-info/dependency_links.txt
...
...
src/mqtt_influx_backend/mac_to_room.json
View file @
35827b83
...
@@ -3,5 +3,9 @@
...
@@ -3,5 +3,9 @@
"11:22:33:44:55:66"
:
"K
\u
00fcche"
,
"11:22:33:44:55:66"
:
"K
\u
00fcche"
,
"77:88:99:AA:BB:CC"
:
"Schlafzimmer"
,
"77:88:99:AA:BB:CC"
:
"Schlafzimmer"
,
"DE:AD:BE:EF:12:34"
:
""
,
"DE:AD:BE:EF:12:34"
:
""
,
"DK:AD:BE:EF:12:34"
:
""
"DK:AD:BE:EF:12:34"
:
""
,
"EK:AD:BE:EF:12:34"
:
""
,
"lK:AD:BE:EF:12:34"
:
""
,
"\u00f6K:AD:BE:EF:12:34"
:
""
,
"MK:AD:BE:EF:12:34"
:
""
}
}
\ No newline at end of file
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