Protokół MQTT do integracji agentów
Konfiguracja brokera, subskrypcja tematów i aktywowanie agentów za pomocą komunikatów.
Protokół MQTT do integracji agentów to bezpłatna lekcja AI Agents na CoddyKit. To lekcja 1 z 4. Możesz przeczytać całą lekcję poniżej za darmo — a potem ćwiczyć ją interaktywnie w przeglądarce z wbudowanym edytorem kodu i tutorem AI dostępnym 24/7. To część ścieżki edukacyjnej AI Agents, a Twój postęp synchronizuje się między webem a aplikacją CoddyKit. Kurs AI Agents zawiera 4 lekcji w sumie.
MQTT i agenci IoT
MQTT (Message Queuing Telemetry Transport) to lekki protokół publikowania i subskrybowania wiadomości, zaprojektowany z myślą o środowiskach IoT o małej przepustowości i dużych opóźnieniach. Agenci mogą subskrybować tematy MQTT, aby odbierać dane z sensorów w czasie rzeczywistym, oraz publikować polecenia kierowane do urządzeń.
Instalowanie paho-mqtt
Biblioteka paho-mqtt to standardowy klient MQTT dla języka Python. Należy zainstalować ją za pomocą pip. Obsługuje MQTT 3.1.1 i 5.0, szyfrowanie TLS oraz wszystkie trzy poziomy QoS.
# 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)Łączenie z brokerem
Należy połączyć się z brokerem MQTT (np. Mosquitto, HiveMQ, EMQX lub brokerem w chmurze). Wywołanie connect nie blokuje działania programu; należy użyć loop_start(), aby uruchomić pętlę sieciową w wątku działającym w tle.
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 connectionSubskrybowanie tematu
Tematy to hierarchiczne ciągi znaków, takie jak sensors/temperature lub factory/line1/pressure. Użyj symbolu # jako symbolu wieloznacznego dla wszystkich podtematów albo symbolu + dla pojedynczego poziomu. Dołącz funkcję zwrotną on_message, aby przetwarzać przychodzące wiadomości.
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}')Wyjaśnienie poziomów QoS
MQTT udostępnia trzy poziomy jakości usług: QoS 0 — wyślij i zapomnij (najszybszy, wiadomości mogą zostać utracone), QoS 1 — co najmniej raz (wiadomość jest potwierdzana, ale może zostać zduplikowana), QoS 2 — dokładnie raz (gwarantowane dostarczenie, najwolniejszy). Wybór zależy od Państwa tolerancji na utratę danych i opóźnienia.
# 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')Wiadomości zachowane
Wiadomość zachowana to ostatnia opublikowana wartość dla danego tematu, którą broker przechowuje i natychmiast wysyła każdemu nowemu subskrybentowi. Jest to idealne rozwiązanie dla tematów zawierających stan czujnika: nowa instancja agenta dołączająca do sieci od razu zna bieżącą wartość czujnika, bez oczekiwania na kolejną aktualizację.
# 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')Budowanie agenta czujnika MQTT
Agent czujnika subskrybuje tematy zawierające surowe dane z czujników, weryfikuje przychodzące dane i decyduje, czy wywołać działanie. Logika decyzyjna wykorzystuje LLM tylko w złożonych przypadkach; proste sprawdzanie progów jest obsługiwane bezpośrednio w języku Python, aby zapewnić szybkość działania.
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 przez WebSocket
MQTT przez WebSocket (port 8083 lub 8084 dla TLS) umożliwia pulpitom nawigacyjnym i agentom działającym w przeglądarce łączenie się z brokerami MQTT bez natywnego połączenia TCP. Należy skonfigurować paho-mqtt tak, aby korzystał z transportu WebSocket, używając opcji 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')Ostatnia wola i testament
Funkcja Last Will and Testament (LWT) MQTT pozwala brokerowi opublikować wiadomość w imieniu klienta, jeśli klient rozłączy się nieoczekiwanie. Należy jej użyć do powiadamiania innych agentów, że agent czujnika przeszedł w tryb offline, aby mogły przełączyć się na tryb awaryjny.
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()Najlepsze praktyki projektowania tematów
Dobrze zaprojektowane tematy ułatwiają utrzymanie i skalowanie kodu agentów. Należy stosować hierarchię: location/device_type/device_id/measurement. Należy unikać spacji i znaków specjalnych. Tematy powinny być krótkie — zwiększają narzut każdej wiadomości.
# 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/celsiusPoprawne rozłączanie i ponowne łączenie
Agenci działający w środowiskach IoT muszą prawidłowo obsługiwać przerwy w sieci. Należy użyć wbudowanej logiki ponownego łączenia paho-mqtt: ustawić reconnect_on_failure=True i ponownie subskrybować tematy po każdym ponownym połączeniu, ponieważ subskrypcje nie są domyślnie zachowywane (chyba że używane są sesje trwałe).
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()Sprawdzenie wiedzy
Jaki poziom QoS gwarantuje, że wiadomość zostanie dostarczona dokładnie raz?
Podsumowanie: MQTT w integracji agentów
Doskonale! W tej lekcji poznali Państwo:
- paho-mqtt:
connect(),loop_start(),subscribe(),publish() - Poziomy QoS: 0 = najwyżej raz, 1 = co najmniej raz, 2 = dokładnie raz
- Wiadomości zachowane: broker przechowuje ostatnią wartość; nowi subskrybenci otrzymują ją natychmiast
- LWT: broker publikuje wiadomość o przejściu w tryb offline po nieoczekiwanym rozłączeniu
- Ponowne łączenie: należy ponownie subskrybować tematy w
on_connect; użyćreconnect_delay_set
Następny temat: przetwarzanie strumieni danych szeregów czasowych — okna kroczące, średnie kroczące i wykrywanie skoków.
Często zadawane pytania
Czy lekcja „Protokół MQTT do integracji agentów” jest bezpłatna?
Tak — pełny tekst „Protokół MQTT do integracji agentów” jest dostępny za darmo tutaj w sieci. Aby ćwiczyć ją interaktywnie (wbudowany edytor kodu i tutor AI dostępny 24/7) i odblokować resztę kursu AI Agents, przejdź na CoddyKit PRO. Kurs AI Agents zawiera 4 lekcji w sumie.
Co nauczysz się w „Protokół MQTT do integracji agentów”?
Konfiguracja brokera, subskrypcja tematów i aktywowanie agentów za pomocą komunikatów. Ćwiczysz AI Agents z praktycznym kodem, który uruchamiasz bezpośrednio w przeglądarce, a tutor AI dostępny 24/7 odpowiada na Twoje pytania podczas pracy nad lekcją.
Czy potrzebuję doświadczenia, aby zacząć AI Agents?
Nie wymagamy żadnego doświadczenia. AI Agents w CoddyKit jest strukturyzowany dla początkujących i zaawansowanych użytkowników, więc możesz zacząć tutaj lub od początku i uczyć się w swoim tempie. To lekcja 1 z 4.
Ile czasu zajmuje lekcja „Protokół MQTT do integracji agentów”?
Większość lekcji CoddyKit trwa około 5–10 minut. Każda lekcja to mały, interaktywny krok, dzięki czemu robisz systematyczne postępy i zawsze wracasz dokładnie do tego samego miejsca — na webie i w aplikacji.
Czy mogę pisać i uruchamiać kod w tej lekcji AI Agents?
Tak. Każda lekcja AI Agents zawiera wbudowany edytor kodu, więc piszesz i uruchamiasz prawdziwy kod bezpośrednio w przeglądarce i od razu otrzymujesz sprzężenie zwrotne od AI — bez konfiguracji na komputerze.
Wszystkie lekcje w tym kursie
- Protokół MQTT do integracji agentów
- Przetwarzanie danych szeregów czasowych przez agentów
- Automatyczne reagowanie na zdarzenia z czujników
- Wdrażanie lekkich agentów na brzegu sieci