Введение в потоки данных

Анализ данных Интернета вещей на Python

Matthias Voppichler

IT Developer

Что такое поток данных

  • Непрерывный поток данных
  • Примеры
    • Сообщения в Twitter
    • Новостные статьи онлайн
    • Видеопотоки
    • Данные с датчиков (IoT)
    • Биржевые заявки (финансы)
Анализ данных Интернета вещей на Python

Что такое поток данных

  • Непрерывный поток данных
  • Примеры
    • Сообщения в Twitter
    • Новостные статьи онлайн
    • Видеопотоки
    • Данные с датчиков (IoT)
    • Биржевые заявки (финансы)
Анализ данных Интернета вещей на Python

MQTT

  • Протокол обмена сообщениями
  • Публикация / подписка
  • Малый объём

Сервер -> выступает брокером сообщений

Клиент:

  • Подключается к брокеру
  • Публикует данные
  • Подписывается на топики

Message Queuing Telemetry Transport

Поток данных MQTT — от производителя к брокеру и потребителям

Анализ данных Интернета вещей на Python

Библиотека Python

Eclipse Paho™ MQTT Python Client

# Import MQTT library
import paho.mqtt

Подробная информация и документация доступны на GitHub https://github.com/eclipse/paho.mqtt.python

Анализ данных Интернета вещей на Python

Одно сообщение

import paho.mqtt.subscribe as subscribe
msg = subscribe.simple("paho/test/simple", 
                       hostname="test.mosquitto.org")

print(f"{msg.topic}, {msg.payload}")

Вывод:

paho/test/simple, {"time": 1549481572, "humidity": 77, "temp": 21}
Анализ данных Интернета вещей на Python

Обратный вызов

def on_message(client, userdata, message):

print(f"{message.topic} : {message.payload}")

Аргументы

  • client — экземпляр клиента
  • userdata — пользовательские данные
  • message — экземпляр MQTTMessage
Анализ данных Интернета вещей на Python

Обратный вызов

import paho.mqtt.subscribe as subscribe


subscribe.callback(on_message, topics="datacamp/roomtemp", hostname="test.mosquitto.org")
Анализ данных Интернета вещей на Python

Подписка MQTT

import paho.mqtt.subscribe as subscribe

def on_message(client, userdata, message): print("{} : {}".format(message.topic, message.payload))
subscribe.callback(on_message, topics="datacamp/roomtemp", hostname="test.mosquitto.org")
datacamp/roomtemp : b'{"time": 1543344857, "hum": 34, "temp": 24}'
datacamp/roomtemp : b'{"time": 1543344858, "hum": 35, "temp": 23}'
datacamp/roomtemp : b'{"time": 1543344860, "hum": 36, "temp": 22}'
datacamp/roomtemp : b'{"time": 1543344946, "hum": 37, "temp": 22}'
datacamp/roomtemp : b'{"time": 1543345010, "hum": 36, "temp": 13}'
Анализ данных Интернета вещей на Python

Давайте потренируемся!

Анализ данных Интернета вещей на Python

Preparing Video For Download...