-
Notifications
You must be signed in to change notification settings - Fork 0
/
subscriber-async.py
61 lines (50 loc) · 1.51 KB
/
subscriber-async.py
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
import asyncio
import aiomqtt
from dotenv import load_dotenv
from colorama import Fore
import os
load_dotenv()
broker = os.getenv("BROKER")
port = int(os.getenv("PORT"))
topic = os.getenv("TOPIC")
if os.getenv("USER_NAME") is None or os.getenv("USER_NAME") == "":
user_name = None
pass_word = None
else:
user_name = os.getenv("USER_NAME")
pass_word = os.getenv("PASSWORD")
# Sample async function to be called
async def processSnapshot(message: aiomqtt.Message):
await asyncio.sleep(1)
# This gets called whenever when we get an MQTT message
async def mqtt_on_message(msg: aiomqtt.Message):
print(
Fore.WHITE
+ "TOPIC: "
+ Fore.LIGHTYELLOW_EX
+ str(msg.topic)
+ Fore.WHITE
+ "\tDATA: "
+ Fore.GREEN
+ " "
+ msg.payload.decode("utf-8", "ignore")
)
return await processSnapshot(msg.payload.decode("utf-8", "ignore"))
async def main():
...
try:
async with aiomqtt.Client(
hostname=broker,
port=port,
username=user_name,
password=pass_word,
) as client:
async with client.messages() as messages:
await client.subscribe(topic)
# await client.subscribe(topic-2)
async for message in messages:
asyncio.ensure_future(mqtt_on_message(message))
except Exception as x:
print("application has shutdown or could not start\n", x)
if __name__ == "__main__":
asyncio.run(main())