0Pricing
AI Agents · Lekcja

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 connection

Subskrybowanie 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/celsius

Poprawne 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

  1. Protokół MQTT do integracji agentów
  2. Przetwarzanie danych szeregów czasowych przez agentów
  3. Automatyczne reagowanie na zdarzenia z czujników
  4. Wdrażanie lekkich agentów na brzegu sieci
← Powrót do AI Agents