mirror of
https://github.com/daniviga/bite.git
synced 2024-11-24 22:06:13 +01:00
Compare commits
7 Commits
681f99d2f4
...
da92935001
Author | SHA1 | Date | |
---|---|---|---|
da92935001 | |||
465870a9c5 | |||
7e689eca23 | |||
ea9f9ef705 | |||
e3785d4669 | |||
49211437d2 | |||
23dfb6837d |
13
README.md
13
README.md
|
@ -4,7 +4,7 @@ Playing with IoT
|
|||
|
||||
[![Build Status](https://travis-ci.com/daniviga/bite.svg?branch=master)](https://travis-ci.com/daniviga/bite)
|
||||
![AGPLv3](./docs/.badges/agpl3.svg)
|
||||
![Python 3.9](./docs/.badges/python.svg)
|
||||
![Python 3.11](./docs/.badges/python.svg)
|
||||
![MQTT](./docs/.badges/mqtt.svg)
|
||||
![Moby](./docs/.badges/moby.svg)
|
||||
![docker-compose 3.7+](./docs/.badges/docker-compose.svg)
|
||||
|
@ -13,13 +13,6 @@ This project is for educational purposes only. It does not implement any
|
|||
authentication and/or encryption protocol, so it is not suitable for real
|
||||
production.
|
||||
|
||||
![Application Schema](./docs/application_chart.png)
|
||||
|
||||
### Future implementations
|
||||
|
||||
- Broker HA via [VerneMQ clustering](https://docs.vernemq.com/clustering/introduction)
|
||||
- Stream analytics via [Apache Spark](https://spark.apache.org/)
|
||||
|
||||
## Installation
|
||||
|
||||
### Requirements
|
||||
|
@ -36,8 +29,10 @@ The application stack is composed by the following components:
|
|||
- [Django](https://www.djangoproject.com/) with
|
||||
[Django REST framework](https://www.django-rest-framework.org/)
|
||||
web application (running via `gunicorn` in production mode)
|
||||
- `mqtt-to-db` custom daemon to dump telemetry into the timeseries database
|
||||
- `dispatcher` custom daemon to dump telemetry into the Kafka queue
|
||||
- `handler` custom daemon to dump telemetry into the timeseries database from the Kafka queue
|
||||
- telemetry payload is stored as json object (via PostgreSQL JSON data type)
|
||||
- [Kafka](https://kafka.apache.org/) broker
|
||||
- [Timescale](https://www.timescale.com/) DB,
|
||||
a [PostgreSQL](https://www.postgresql.org/) database with a timeseries extension
|
||||
- [Mosquitto](https://mosquitto.org/) MQTT broker (see alternatives below)
|
||||
|
|
|
@ -151,6 +151,10 @@ STATIC_URL = '/static/'
|
|||
|
||||
STATIC_ROOT = '/srv/appdata/bite/static'
|
||||
|
||||
REST_FRAMEWORK = {
|
||||
'DEFAULT_AUTHENTICATION_CLASSES': []
|
||||
}
|
||||
|
||||
SKIP_WHITELIST = True
|
||||
|
||||
MQTT_BROKER = {
|
||||
|
@ -158,6 +162,11 @@ MQTT_BROKER = {
|
|||
'PORT': '1883',
|
||||
}
|
||||
|
||||
KAFKA_BROKER = {
|
||||
'HOST': 'kafka',
|
||||
'PORT': '9092',
|
||||
}
|
||||
|
||||
# If no local_settings.py is availble in the current folder let's try to
|
||||
# load it from the application root
|
||||
try:
|
||||
|
|
|
@ -22,23 +22,27 @@ import asyncio
|
|||
import json
|
||||
import time
|
||||
import paho.mqtt.client as mqtt
|
||||
from kafka import KafkaProducer
|
||||
from kafka.errors import NoBrokersAvailable
|
||||
from asgiref.sync import sync_to_async
|
||||
from asyncio_mqtt import Client
|
||||
from aiomqtt import Client
|
||||
|
||||
from django.conf import settings
|
||||
from django.core.management.base import BaseCommand
|
||||
from django.core.exceptions import ObjectDoesNotExist
|
||||
|
||||
from api.models import Device
|
||||
from telemetry.models import Telemetry
|
||||
|
||||
MQTT_HOST = settings.MQTT_BROKER['HOST']
|
||||
MQTT_PORT = int(settings.MQTT_BROKER['PORT'])
|
||||
|
||||
|
||||
class Command(BaseCommand):
|
||||
help = 'MQTT to DB deamon'
|
||||
|
||||
MQTT_HOST = settings.MQTT_BROKER['HOST']
|
||||
MQTT_PORT = int(settings.MQTT_BROKER['PORT'])
|
||||
KAFKA_HOST = settings.KAFKA_BROKER['HOST']
|
||||
KAFKA_PORT = int(settings.KAFKA_BROKER['PORT'])
|
||||
producer = None
|
||||
|
||||
@sync_to_async
|
||||
def get_device(self, serial):
|
||||
try:
|
||||
|
@ -47,24 +51,23 @@ class Command(BaseCommand):
|
|||
return None
|
||||
|
||||
@sync_to_async
|
||||
def store_telemetry(self, device, payload):
|
||||
Telemetry.objects.create(
|
||||
device=device,
|
||||
transport='mqtt',
|
||||
clock=payload['clock'],
|
||||
payload=payload['payload']
|
||||
def dispatch(self, message):
|
||||
self.producer.send(
|
||||
'telemetry', {"transport": 'mqtt',
|
||||
"body": message}
|
||||
)
|
||||
|
||||
async def mqtt_broker(self):
|
||||
async with Client(MQTT_HOST, port=MQTT_PORT) as client:
|
||||
async with Client(self.MQTT_HOST, port=self.MQTT_PORT) as client:
|
||||
# use shared subscription for HA/balancing
|
||||
await client.subscribe("$share/telemetry/#")
|
||||
async with client.unfiltered_messages() as messages:
|
||||
async with client.messages() as messages:
|
||||
async for message in messages:
|
||||
payload = json.loads(message.payload.decode('utf-8'))
|
||||
device = await self.get_device(message.topic)
|
||||
if device is not None:
|
||||
await self.store_telemetry(device, payload)
|
||||
message_body = json.loads(
|
||||
message.payload.decode('utf-8'))
|
||||
await self.dispatch(message_body)
|
||||
else:
|
||||
self.stdout.write(
|
||||
self.style.ERROR(
|
||||
|
@ -74,13 +77,28 @@ class Command(BaseCommand):
|
|||
client = mqtt.Client()
|
||||
while True:
|
||||
try:
|
||||
client.connect(MQTT_HOST, MQTT_PORT)
|
||||
client.connect(self.MQTT_HOST, self.MQTT_PORT)
|
||||
break
|
||||
except (socket.gaierror, ConnectionRefusedError):
|
||||
self.stdout.write(
|
||||
self.style.WARNING('WARNING: Broker not available'))
|
||||
self.style.WARNING('WARNING: MQTT broker not available'))
|
||||
time.sleep(5)
|
||||
|
||||
self.stdout.write(self.style.SUCCESS('INFO: Broker subscribed'))
|
||||
while True:
|
||||
try:
|
||||
self.producer = KafkaProducer(
|
||||
bootstrap_servers='{}:{}'.format(
|
||||
self.KAFKA_HOST, self.KAFKA_PORT
|
||||
),
|
||||
value_serializer=lambda v: json.dumps(v).encode('utf-8'),
|
||||
retries=5
|
||||
)
|
||||
break
|
||||
except NoBrokersAvailable:
|
||||
self.stdout.write(
|
||||
self.style.WARNING('WARNING: Kafka broker not available'))
|
||||
time.sleep(5)
|
||||
|
||||
self.stdout.write(self.style.SUCCESS('INFO: Brokers subscribed'))
|
||||
client.disconnect()
|
||||
asyncio.run(self.mqtt_broker())
|
76
bite/telemetry/management/commands/handler.py
Normal file
76
bite/telemetry/management/commands/handler.py
Normal file
|
@ -0,0 +1,76 @@
|
|||
# -*- coding: utf-8 -*-
|
||||
# vim: tabstop=4 shiftwidth=4 softtabstop=4
|
||||
#
|
||||
# BITE - A Basic/IoT/Example
|
||||
# Copyright (C) 2020-2021 Daniele Viganò <daniele@vigano.me>
|
||||
#
|
||||
# BITE is free software: you can redistribute it and/or modify
|
||||
# it under the terms of the GNU Affero General Public License as published by
|
||||
# the Free Software Foundation, either version 3 of the License, or
|
||||
# (at your option) any later version.
|
||||
#
|
||||
# BITE is distributed in the hope that it will be useful,
|
||||
# but WITHOUT ANY WARRANTY; without even the implied warranty of
|
||||
# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
|
||||
# GNU Affero General Public License for more details.
|
||||
#
|
||||
# You should have received a copy of the GNU Affero General Public License
|
||||
# along with this program. If not, see <http://www.gnu.org/licenses/>.
|
||||
|
||||
import json
|
||||
import time
|
||||
from kafka import KafkaConsumer
|
||||
from kafka.errors import NoBrokersAvailable
|
||||
|
||||
from django.conf import settings
|
||||
from django.core.management.base import BaseCommand
|
||||
from django.core.exceptions import ObjectDoesNotExist
|
||||
|
||||
from api.models import Device
|
||||
from telemetry.models import Telemetry
|
||||
|
||||
|
||||
class Command(BaseCommand):
|
||||
help = 'MQTT to DB deamon'
|
||||
|
||||
KAFKA_HOST = settings.KAFKA_BROKER['HOST']
|
||||
KAFKA_PORT = int(settings.KAFKA_BROKER['PORT'])
|
||||
|
||||
def get_device(self, serial):
|
||||
try:
|
||||
return Device.objects.get(serial=serial)
|
||||
except ObjectDoesNotExist:
|
||||
return None
|
||||
|
||||
def store_telemetry(self, transport, message):
|
||||
Telemetry.objects.create(
|
||||
transport=transport,
|
||||
device=self.get_device(message["device"]),
|
||||
clock=message["clock"],
|
||||
payload=message["payload"]
|
||||
)
|
||||
|
||||
def handle(self, *args, **options):
|
||||
while True:
|
||||
try:
|
||||
consumer = KafkaConsumer(
|
||||
"telemetry",
|
||||
bootstrap_servers='{}:{}'.format(
|
||||
self.KAFKA_HOST, self.KAFKA_PORT
|
||||
),
|
||||
group_id="handler",
|
||||
value_deserializer=lambda m: json.loads(m.decode('utf8')),
|
||||
)
|
||||
break
|
||||
except NoBrokersAvailable:
|
||||
self.stdout.write(
|
||||
self.style.WARNING('WARNING: Kafka broker not available'))
|
||||
time.sleep(5)
|
||||
|
||||
self.stdout.write(self.style.SUCCESS('INFO: Kafka broker subscribed'))
|
||||
for message in consumer:
|
||||
self.store_telemetry(
|
||||
message.value["transport"],
|
||||
message.value["body"]
|
||||
)
|
||||
consumer.unsuscribe()
|
|
@ -17,20 +17,20 @@
|
|||
# You should have received a copy of the GNU Affero General Public License
|
||||
# along with this program. If not, see <http://www.gnu.org/licenses/>.
|
||||
|
||||
FROM python:3.9-alpine AS builder
|
||||
FROM python:3.11-alpine AS builder
|
||||
RUN apk update && apk add gcc musl-dev postgresql-dev \
|
||||
&& pip install psycopg2-binary
|
||||
|
||||
# ---
|
||||
|
||||
FROM python:3.9-alpine
|
||||
FROM python:3.11-alpine
|
||||
ENV PYTHONUNBUFFERED 1
|
||||
ENV DJANGO_SETTINGS_MODULE "bite.settings"
|
||||
|
||||
RUN apk update && apk add --no-cache postgresql-libs \
|
||||
&& wget https://github.com/jwilder/dockerize/releases/download/v0.6.1/dockerize-alpine-linux-amd64-v0.6.1.tar.gz -qO- \
|
||||
&& wget https://github.com/jwilder/dockerize/releases/download/v0.7.0/dockerize-alpine-linux-amd64-v0.7.0.tar.gz -qO- \
|
||||
| tar -xz -C /usr/local/bin
|
||||
COPY --from=builder /usr/local/lib/python3.9/site-packages/ /usr/local/lib/python3.9/site-packages/
|
||||
COPY --from=builder /usr/local/lib/python3.11/site-packages/ /usr/local/lib/python3.11/site-packages/
|
||||
COPY --chown=1000:1000 requirements.txt /srv/app/bite/requirements.txt
|
||||
RUN pip3 install --no-cache-dir -r /srv/app/bite/requirements.txt
|
||||
|
||||
|
|
|
@ -36,6 +36,10 @@ services:
|
|||
ports:
|
||||
- "${CUSTOM_DOCKER_IP:-0.0.0.0}:8000:8000"
|
||||
|
||||
kafka:
|
||||
ports:
|
||||
- "${CUSTOM_DOCKER_IP:-0.0.0.0}:9092:9092"
|
||||
|
||||
data-migration:
|
||||
volumes:
|
||||
- ../bite:/srv/app/bite
|
||||
|
@ -44,6 +48,10 @@ services:
|
|||
volumes:
|
||||
- ../bite:/srv/app/bite
|
||||
|
||||
mqtt-to-db:
|
||||
dispatcher:
|
||||
volumes:
|
||||
- ../bite:/srv/app/bite
|
||||
|
||||
handler:
|
||||
volumes:
|
||||
- ../bite:/srv/app/bite
|
||||
|
|
|
@ -29,6 +29,10 @@ services:
|
|||
volumes:
|
||||
- ./django/production.py.sample:/srv/app/bite/bite/production.py
|
||||
|
||||
mqtt-to-db:
|
||||
dispatcher:
|
||||
volumes:
|
||||
- ./django/production.py.sample:/srv/app/bite/bite/production.py
|
||||
|
||||
handler:
|
||||
volumes:
|
||||
- ./django/production.py.sample:/srv/app/bite/bite/production.py
|
||||
|
|
|
@ -62,6 +62,30 @@ services:
|
|||
ports:
|
||||
- "${CUSTOM_DOCKER_IP:-0.0.0.0}:1883:1883"
|
||||
|
||||
zookeeper:
|
||||
image: confluentinc/cp-zookeeper:latest
|
||||
networks:
|
||||
- net
|
||||
environment:
|
||||
ZOOKEEPER_CLIENT_PORT: 2181
|
||||
ZOOKEEPER_TICK_TIME: 2000
|
||||
ports:
|
||||
- 22181:2181
|
||||
|
||||
kafka:
|
||||
image: confluentinc/cp-kafka:latest
|
||||
depends_on:
|
||||
- zookeeper
|
||||
networks:
|
||||
- net
|
||||
environment:
|
||||
KAFKA_BROKER_ID: 1
|
||||
KAFKA_ZOOKEEPER_CONNECT: zookeeper:2181
|
||||
KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://kafka:9092,PLAINTEXT_HOST://localhost:29092
|
||||
KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: PLAINTEXT:PLAINTEXT,PLAINTEXT_HOST:PLAINTEXT
|
||||
KAFKA_INTER_BROKER_LISTENER_NAME: PLAINTEXT
|
||||
KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1
|
||||
|
||||
ingress:
|
||||
<<: *service_default
|
||||
image: nginx:stable-alpine
|
||||
|
@ -104,13 +128,21 @@ services:
|
|||
- staticdata:/srv/appdata/bite/static
|
||||
command: ["python3", "manage.py", "collectstatic", "--noinput"]
|
||||
|
||||
mqtt-to-db:
|
||||
dispatcher:
|
||||
<<: *service_default
|
||||
image: daniviga/bite
|
||||
command: ["python3", "manage.py", "mqtt-to-db"]
|
||||
command: ["python3", "manage.py", "dispatcher"]
|
||||
networks:
|
||||
- net
|
||||
depends_on:
|
||||
- broker
|
||||
|
||||
handler:
|
||||
<<: *service_default
|
||||
image: daniviga/bite
|
||||
command: ["python3", "manage.py", "handler"]
|
||||
networks:
|
||||
- net
|
||||
depends_on:
|
||||
- data-migration
|
||||
- timescale
|
||||
- broker
|
||||
|
|
|
@ -17,9 +17,9 @@
|
|||
# You should have received a copy of the GNU Affero General Public License
|
||||
# along with this program. If not, see <http://www.gnu.org/licenses/>.
|
||||
|
||||
FROM alpine:3.15
|
||||
FROM alpine:3.18
|
||||
|
||||
RUN apk update && apk add chrony && \
|
||||
RUN apk add --no-cache chrony && \
|
||||
chown -R chrony:chrony /var/lib/chrony
|
||||
COPY ./chrony.conf /etc/chrony/chrony.conf
|
||||
|
||||
|
|
|
@ -17,7 +17,7 @@
|
|||
# You should have received a copy of the GNU Affero General Public License
|
||||
# along with this program. If not, see <http://www.gnu.org/licenses/>.
|
||||
|
||||
FROM python:3.9-alpine
|
||||
FROM python:3.11-alpine
|
||||
|
||||
RUN pip3 install urllib3 paho-mqtt
|
||||
COPY ./device_simulator.py /opt/bite/device_simulator.py
|
||||
|
|
|
@ -94,7 +94,7 @@ def main():
|
|||
parser.add_argument('-s', '--serial',
|
||||
default=os.environ.get('IOT_SERIAL'),
|
||||
help='IoT device serial number')
|
||||
parser.add_argument('-d', '--delay', metavar='s', type=int,
|
||||
parser.add_argument('-d', '--delay', metavar='s', type=float,
|
||||
default=os.environ.get('IOT_DELAY', 10),
|
||||
help='Delay between requests')
|
||||
args = parser.parse_args()
|
||||
|
|
|
@ -1 +1 @@
|
|||
<svg xmlns="http://www.w3.org/2000/svg" xmlns:xlink="http://www.w3.org/1999/xlink" width="93.0" height="20"><linearGradient id="smooth" x2="0" y2="100%"><stop offset="0" stop-color="#bbb" stop-opacity=".1"/><stop offset="1" stop-opacity=".1"/></linearGradient><clipPath id="round"><rect width="93.0" height="20" rx="3" fill="#fff"/></clipPath><g clip-path="url(#round)"><rect width="65.5" height="20" fill="#555"/><rect x="65.5" width="27.5" height="20" fill="#007ec6"/><rect width="93.0" height="20" fill="url(#smooth)"/></g><g fill="#fff" text-anchor="middle" font-family="DejaVu Sans,Verdana,Geneva,sans-serif" font-size="110"><image x="5" y="3" width="14" height="14" xlink:href=""/><text x="422.5" y="150" fill="#010101" fill-opacity=".3" transform="scale(0.1)" textLength="385.0" lengthAdjust="spacing">python</text><text x="422.5" y="140" transform="scale(0.1)" textLength="385.0" lengthAdjust="spacing">python</text><text x="782.5" y="150" fill="#010101" fill-opacity=".3" transform="scale(0.1)" textLength="175.0" lengthAdjust="spacing">3.9</text><text x="782.5" y="140" transform="scale(0.1)" textLength="175.0" lengthAdjust="spacing">3.9</text></g></svg>
|
||||
<svg xmlns="http://www.w3.org/2000/svg" xmlns:xlink="http://www.w3.org/1999/xlink" width="93.0" height="20"><linearGradient id="smooth" x2="0" y2="100%"><stop offset="0" stop-color="#bbb" stop-opacity=".1"/><stop offset="1" stop-opacity=".1"/></linearGradient><clipPath id="round"><rect width="93.0" height="20" rx="3" fill="#fff"/></clipPath><g clip-path="url(#round)"><rect width="65.5" height="20" fill="#555"/><rect x="65.5" width="27.5" height="20" fill="#007ec6"/><rect width="93.0" height="20" fill="url(#smooth)"/></g><g fill="#fff" text-anchor="middle" font-family="DejaVu Sans,Verdana,Geneva,sans-serif" font-size="110"><image x="5" y="3" width="14" height="14" xlink:href=""/><text x="422.5" y="150" fill="#010101" fill-opacity=".3" transform="scale(0.1)" textLength="385.0" lengthAdjust="spacing">python</text><text x="422.5" y="140" transform="scale(0.1)" textLength="385.0" lengthAdjust="spacing">python</text><text x="782.5" y="150" fill="#010101" fill-opacity=".3" transform="scale(0.1)" textLength="175.0" lengthAdjust="spacing">3.11</text><text x="782.5" y="140" transform="scale(0.1)" textLength="175.0" lengthAdjust="spacing">3.11</text></g></svg>
|
||||
|
|
Before Width: | Height: | Size: 2.4 KiB After Width: | Height: | Size: 2.4 KiB |
|
@ -4,3 +4,4 @@ ipython
|
|||
flake8
|
||||
pyinstrument
|
||||
django-debug-toolbar
|
||||
urllib3
|
||||
|
|
|
@ -4,7 +4,8 @@ djangorestframework
|
|||
django-health-check
|
||||
psycopg2-binary
|
||||
paho-mqtt
|
||||
asyncio-mqtt
|
||||
kafka-python
|
||||
aiomqtt
|
||||
PyYAML
|
||||
uritemplate
|
||||
pygments
|
||||
|
|
Loading…
Reference in New Issue
Block a user