0Pricing
AI Agents · Урок

Протокол MQTT для интеграции агентов

Настройка брокера, подписка на темы и активация агентов по сообщениям.

«Протокол MQTT для интеграции агентов» — бесплатный урок AI Agents на CoddyKit. Это урок 1 из 4. Ты можешь прочитать весь урок бесплатно ниже — а потом практиковать его прямо в браузере с встроенным редактором кода и ИИ-репетитором 24/7. Это часть пути обучения AI Agents, и твой прогресс синхронизируется между веб-версией и приложением CoddyKit. Курс AI Agents содержит 4 уроков всего.

MQTT и агенты Интернета вещей

MQTT (протокол передачи телеметрии с очередями сообщений) — это лёгкий протокол публикации и подписки, предназначенный для сред Интернета вещей с низкой пропускной способностью и высокой задержкой. Агенты могут подписываться на темы MQTT, чтобы получать данные датчиков в реальном времени, и публиковать команды для устройств.

Установка paho-mqtt

Библиотека paho-mqtt — стандартный MQTT-клиент для Python. Установите её с помощью pip. Она поддерживает MQTT 3.1.1 и 5.0, шифрование TLS и все три уровня качества обслуживания.

# pip install paho-mqtt

import paho.mqtt.client as mqtt

# Create a client instance
client = mqtt.Client(client_id='agent_001', protocol=mqtt.MQTTv311)

# Optional: set credentials if broker requires authentication
client.username_pw_set(username='agent_user', password='YOUR_PASSWORD')

# Optional: enable TLS for secure connections
# client.tls_set('/path/to/ca.crt')

print('MQTT client created:', client._client_id)

Подключение к брокеру

Подключитесь к брокеру MQTT (например, Mosquitto, HiveMQ, EMQX или облачному брокеру). Вызов connect не блокирует выполнение; используйте loop_start(), чтобы запустить сетевой цикл в фоновом потоке.

import paho.mqtt.client as mqtt
import time

BROKER_HOST = 'broker.hivemq.com'  # public test broker
BROKER_PORT = 1883
KEEP_ALIVE_SECONDS = 60

def on_connect(client, userdata, flags, rc):
    status = {
        0: 'Connected successfully',
        1: 'Refused: wrong protocol',
        2: 'Refused: client ID rejected',
        3: 'Refused: server unavailable',
        4: 'Refused: bad credentials',
        5: 'Refused: not authorised'
    }
    print(f'Connect result: {status.get(rc, f"Unknown code {rc}")}')

client = mqtt.Client(client_id='agent_001')
client.on_connect = on_connect
client.connect(BROKER_HOST, BROKER_PORT, keepalive=KEEP_ALIVE_SECONDS)
client.loop_start()  # background thread
time.sleep(1)  # wait for connection

Подписка на топик

Топики — это иерархические строки, например sensors/temperature или factory/line1/pressure. Используйте # как подстановочный символ для всех дочерних топиков или + как подстановочный символ для одного уровня. Добавьте функцию обратного вызова on_message для обработки входящих сообщений.

import json

class FakeMsg:
    def __init__(self, topic, payload):
        self.topic = topic
        self.payload = payload

class FakeMQTTClient:
    def __init__(self):
        self.on_message = None
        self.userdata = {}
    def user_data_set(self, userdata):
        self.userdata = userdata
    def subscribe(self, topic, qos=0):
        print(f'Subscribed to {topic} (qos={qos})')
    def simulate_message(self, topic, payload_bytes):
        self.on_message(self, self.userdata, FakeMsg(topic, payload_bytes))

class Agent:
    def process_sensor_data(self, topic, payload):
        print(f'Agent processing {topic} -> {payload}')

def on_message(client, userdata, msg):
    topic = msg.topic
    try:
        payload = json.loads(msg.payload.decode('utf-8'))
    except (json.JSONDecodeError, UnicodeDecodeError):
        payload = msg.payload.decode('utf-8', errors='replace')

    print(f'Received on {topic}: {payload}')
    userdata['agent'].process_sensor_data(topic, payload)

agent_state = {'agent': Agent()}
client = FakeMQTTClient()
client.user_data_set(agent_state)
client.on_message = on_message

client.subscribe('sensors/temperature', qos=1)
client.subscribe('sensors/+/humidity', qos=1)
client.subscribe('factory/#', qos=1)

client.simulate_message('sensors/temperature', b'{"value": 21.5}')

Объяснение уровней QoS

В MQTT предусмотрено три уровня качества обслуживания: QoS 0 — отправить и забыть (самый быстрый вариант, сообщения могут потеряться), QoS 1 — как минимум один раз (сообщение подтверждается, но могут появиться дубликаты), QoS 2 — ровно один раз (гарантированная доставка, самый медленный вариант). Выбирайте уровень с учётом допустимого риска потери данных и задержки.

# QoS level guidelines for IoT agents:

# QoS 0 — temperature readings updated every second
# (loss of one reading is acceptable)
client.subscribe('sensors/temperature', qos=0)

# QoS 1 — alert notifications (must arrive, duplicates are OK)
client.subscribe('sensors/alerts', qos=1)

# QoS 2 — billing/counting events (each event must be processed exactly once)
client.subscribe('meters/energy_consumed', qos=2)

# When publishing, specify QoS:
client.publish(
    topic='agents/response',
    payload='{"action": "turn_on_cooling"}',
    qos=1,
    retain=False
)
print('Published command with QoS 1')

Сохранённые сообщения

Сохранённое сообщение — это последнее опубликованное значение для топика, которое брокер сохраняет и сразу отправляет каждому новому подписчику. Это особенно удобно для топиков состояния датчиков: новый экземпляр агента, подключившись к сети, сразу узнаёт текущее значение датчика, не дожидаясь следующего обновления.

# Publishing with retain=True persists the last value on the broker
client.publish(
    topic='sensors/thermostat/current_temp',
    payload='{"value": 22.5, "unit": "celsius"}',
    qos=1,
    retain=True  # broker stores this message
)

# Any new subscriber will receive this immediately on subscribe,
# even if it was published hours ago.

# To clear a retained message, publish empty payload:
client.publish(
    topic='sensors/thermostat/current_temp',
    payload='',  # empty payload clears retention
    retain=True
)
print('Retained message cleared')

Создание сенсорного агента MQTT

Агент датчика подписывается на топики с исходными показаниями датчиков, проверяет входящие данные и решает, нужно ли запустить действие. Для сложных случаев логика принятия решений использует LLM, а простые проверки пороговых значений выполняются на чистом Python ради скорости.

import anthropic

class SensorAgent:
    def __init__(self, mqtt_client, llm_api_key: str):
        self.client = mqtt_client
        self.llm = anthropic.Anthropic(api_key=llm_api_key)
        self.readings = []

    def process_sensor_data(self, topic: str, payload: dict):
        value = payload.get('value')
        if value is None:
            return
        self.readings.append({'topic': topic, 'value': value})
        # Fast path: simple threshold
        if topic == 'sensors/temperature' and value > 35:
            self._trigger_action('COOLING_ON', f'Temperature {value}C exceeds threshold')
        # Slow path: complex reasoning via LLM
        elif len(self.readings) >= 10:
            self._llm_analyze()

    def _trigger_action(self, action: str, reason: str):
        payload = '{"action": "' + action + '", "reason": "' + reason + '"}'
        self.client.publish('agents/actions', payload, qos=1)
        print(f'Action triggered: {action} — {reason}')

    def _llm_analyze(self):
        summary = str(self.readings[-10:])
        result = self.llm.messages.create(
            model='claude-opus-4-5', max_tokens=128,
            messages=[{'role': 'user', 'content':
                f'Sensor readings: {summary}. Any anomalies?'}]
        )
        print('LLM analysis:', result.content[0].text)
        self.readings = []

MQTT через WebSocket

MQTT через WebSocket (порт 8083 или 8084 для TLS) позволяет панелям мониторинга и агентам, работающим в браузере, подключаться к брокерам MQTT без собственного TCP-соединения. Настройте paho-mqtt на использование транспорта WebSocket с параметром transport='websockets'.

import paho.mqtt.client as mqtt

# MQTT over WebSocket configuration
ws_client = mqtt.Client(
    client_id='dashboard_agent',
    transport='websockets',  # use WS instead of TCP
    protocol=mqtt.MQTTv311
)

# WebSocket path (broker-specific)
ws_client.ws_set_options(path='/mqtt', headers=None)

# TLS over WebSocket (WSS, port 8084)
# ws_client.tls_set()  # uses system CA store

ws_client.connect(
    host='broker.hivemq.com',
    port=8884,  # WSS port
    keepalive=60
)
ws_client.loop_start()
print('Connected via WebSocket')

Последнее завещание

Последнее завещание и распоряжение (LWT) в MQTT позволяет брокеру опубликовать сообщение от имени клиента, если клиент неожиданно отключится. Используйте эту возможность, чтобы уведомить другие агенты о переходе агента датчика в автономный режим и позволить им переключиться в резервный режим.

import paho.mqtt.client as mqtt

client = mqtt.Client(client_id='agent_001')

# Set LWT before connecting
client.will_set(
    topic='agents/status/agent_001',
    payload='{"status": "offline", "reason": "unexpected_disconnect"}',
    qos=1,
    retain=True  # retain so new subscribers see the last known status
)

def on_connect(client, userdata, flags, rc):
    if rc == 0:
        # Publish 'online' status on connect (overrides LWT retain)
        client.publish(
            'agents/status/agent_001',
            '{"status": "online"}',
            qos=1, retain=True
        )

client.on_connect = on_connect
client.connect('broker.hivemq.com', 1883, keepalive=60)
client.loop_start()

Рекомендации по проектированию топиков

Хорошо спроектированные топики упрощают сопровождение и масштабирование кода агентов. Используйте иерархию: location/device_type/device_id/measurement. Избегайте пробелов и специальных символов. Делайте топики короткими: они увеличивают служебные затраты каждого сообщения.

# Good topic hierarchy examples:
# factory/line1/sensor_042/temperature
# home/living_room/thermostat_01/setpoint
# agents/agent_001/commands/turn_on
# agents/agent_001/status

TOPIC_SCHEMA = {
    'sensor_data': '{location}/{device_type}/{device_id}/{measurement}',
    'agent_command': 'agents/{agent_id}/commands/{action}',
    'agent_status': 'agents/{agent_id}/status',
    'alert': 'alerts/{severity}/{location}'
}

def build_topic(schema_key: str, **kwargs) -> str:
    template = TOPIC_SCHEMA[schema_key]
    return template.format(**kwargs)

# Usage:
topic = build_topic('sensor_data',
                    location='factory', device_type='temp_sensor',
                    device_id='042', measurement='celsius')
print(topic)  # factory/temp_sensor/042/celsius

Корректное отключение и повторное подключение

Агенты в средах IoT должны корректно обрабатывать перебои в сети. Используйте встроенную логику повторного подключения paho-mqtt: установите reconnect_on_failure=True и повторно подписывайтесь при каждом подключении, поскольку подписки по умолчанию не сохраняются (если не используются постоянные сеансы).

import paho.mqtt.client as mqtt
import time

SUBSCRIPTIONS = [
    ('sensors/temperature', 1),
    ('sensors/humidity', 1),
    ('agents/commands', 2)
]

def on_connect(client, userdata, flags, rc):
    if rc == 0:
        print('Connected — re-subscribing to topics')
        for topic, qos in SUBSCRIPTIONS:
            client.subscribe(topic, qos=qos)
    else:
        print(f'Connection failed with code {rc}')

client = mqtt.Client(client_id='agent_reliable', clean_session=True)
client.on_connect = on_connect
client.reconnect_delay_set(min_delay=1, max_delay=30)
client.connect_async('broker.hivemq.com', 1883, keepalive=60)
client.loop_start()

# Graceful shutdown:
# client.disconnect()
# client.loop_stop()

Проверка знаний

Какой уровень QoS гарантирует доставку сообщения ровно один раз?

Итоги: MQTT для интеграции агентов

Отлично! В этом уроке Вы узнали:

  • paho-mqtt: connect(), loop_start(), subscribe(), publish()
  • Уровни QoS: 0 = не более одного раза, 1 = как минимум один раз, 2 = ровно один раз
  • Сохранённые сообщения: брокер хранит последнее значение, а новые подписчики получают его сразу
  • LWT: брокер публикует сообщение о переходе в автономный режим при неожиданном отключении
  • Повторное подключение: повторно подписывайтесь в on_connect; используйте reconnect_delay_set

Далее: обработка потоков данных временных рядов — скользящие окна, скользящие средние и обнаружение выбросов.

Часто задаваемые вопросы

Урок «Протокол MQTT для интеграции агентов» бесплатный?

Да — полный текст урока «Протокол MQTT для интеграции агентов» бесплатно доступен здесь в веб-версии. Чтобы практиковать его интерактивно (встроенный редактор кода и ИИ-репетитор 24/7) и разблокировать остальной курс AI Agents, подпишись на CoddyKit PRO. Курс AI Agents содержит 4 уроков всего.

Чему я научусь в уроке «Протокол MQTT для интеграции агентов»?

Настройка брокера, подписка на темы и активация агентов по сообщениям. Ты практикуешь AI Agents с помощью реального кода, который запускаешь прямо в браузере, и ИИ-репетитор 24/7 отвечает на твои вопросы во время урока.

Нужен ли мне опыт, чтобы начать AI Agents?

Предыдущий опыт не требуется. AI Agents на CoddyKit структурирован для всех уровней — от новичков до продвинутых, поэтому ты можешь начать отсюда или с самого начала и учиться в своем темпе. Это урок 1 из 4.

Сколько времени занимает урок «Протокол MQTT для интеграции агентов»?

Большинство уроков CoddyKit занимают около 5–10 минут. Каждый из них компактный и интерактивный, поэтому ты постоянно делаешь прогресс и продолжаешь с того же места в веб-версии и приложении.

Можно ли писать и запускать код в этом уроке AI Agents?

Да. Каждый урок AI Agents включает встроенный редактор кода, поэтому ты пишешь и запускаешь реальный код прямо в браузере и получаешь моментальную обратную связь от AI — локальная установка не требуется.

Все уроки этого курса

  1. Протокол MQTT для интеграции агентов
  2. Обработка временных рядов агентами
  3. Автоматическая реакция на события датчиков
  4. Развёртывание лёгких агентов на периферии
← Назад к AI Agents