Протокол 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 — локальная установка не требуется.
Все уроки этого курса
- Протокол MQTT для интеграции агентов
- Обработка временных рядов агентами
- Автоматическая реакция на события датчиков
- Развёртывание лёгких агентов на периферии