0Pricing
AI Agents · Lektion

MQTT-Protokoll zur Agentenintegration

Broker einrichten, Topics abonnieren und ereignisgesteuerte Agentenaktivierung umsetzen.

MQTT-Protokoll zur Agentenintegration ist eine kostenlose AI Agents-Lektion auf CoddyKit. Dies ist Lektion 1 von 4. Du kannst die komplette Lektion unten kostenlos lesen – dann übst du sie direkt im Browser mit einem integrierten Code-Editor und einem KI-Tutor rund um die Uhr. Sie ist Teil des AI Agents-Lernpfads, und dein Fortschritt wird über Web und CoddyKit-App synchronisiert. Der AI Agents-Kurs umfasst insgesamt 4 Lektionen.

MQTT- und IoT-Agenten

MQTT (Message Queuing Telemetry Transport) ist ein leichtgewichtiges Publish-Subscribe-Protokoll für IoT-Umgebungen mit geringer Bandbreite und hoher Latenz. Agenten können MQTT-Topics abonnieren, um Sensordaten in Echtzeit zu empfangen, und Befehle an Geräte zurücksenden.

paho-mqtt installieren

Die Bibliothek paho-mqtt ist der standardmäßige Python-MQTT-Client. Installieren Sie sie mit pip. Sie unterstützt MQTT 3.1.1 und 5.0, TLS-Verschlüsselung und alle drei QoS-Stufen.

# 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)

Mit dem Broker verbinden

Verbinden Sie sich mit einem MQTT-Broker (z. B. Mosquitto, HiveMQ, EMQX oder einem Cloud-Broker). Der Aufruf connect ist nicht blockierend. Verwenden Sie loop_start(), um die Netzwerkschleife in einem Hintergrundthread auszuführen.

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

Ein Topic abonnieren

Topics sind hierarchisch strukturierte Zeichenketten wie sensors/temperature oder factory/line1/pressure. Verwenden Sie # als Platzhalter für alle Untertopics oder + als Platzhalter für genau eine Ebene. Verknüpfen Sie einen on_message-Callback, um eingehende Nachrichten zu verarbeiten.

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-Stufen erklärt

MQTT verfügt über drei Quality-of-Service-Stufen: QoS 0 — senden und vergessen (am schnellsten, Nachrichten können verloren gehen), QoS 1 — mindestens einmal (Nachricht bestätigt, kann doppelt empfangen werden), QoS 2 — genau einmal (garantiert, am langsamsten). Wählen Sie die Stufe abhängig davon, wie viel Datenverlust beziehungsweise Latenz Sie tolerieren.

# 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')

Retained Messages

Eine Retained Message ist der zuletzt veröffentlichte Wert eines Topics, den der Broker speichert und sofort an jeden neuen Abonnenten sendet. Das ist ideal für Topics zum Sensorstatus: Eine neue Agenteninstanz, die dem Netzwerk beitritt, kennt sofort den aktuellen Sensorwert, ohne auf die nächste Aktualisierung warten zu müssen.

# 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')

Einen MQTT-Sensoragenten erstellen

Ein Sensoragent abonniert Topics mit rohen Sensordaten, validiert eingehende Daten und entscheidet, ob eine Aktion ausgelöst werden soll. Die Entscheidungslogik verwendet das LLM nur in komplexen Fällen; einfache Schwellenwertprüfungen werden aus Geschwindigkeitsgründen in reinem Python verarbeitet.

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 über WebSocket

MQTT über WebSocket (Port 8083 oder 8084 für TLS) ermöglicht browserbasierten Dashboards und Agenten, sich ohne native TCP-Verbindung mit MQTT-Brokern zu verbinden. Konfigurieren Sie paho-mqtt so, dass der WebSocket-Transport mit der Option transport='websockets' verwendet wird.

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')

Last Will and Testament

Mit dem Last Will and Testament (LWT) von MQTT kann der Broker im Namen eines Clients eine Nachricht veröffentlichen, wenn die Verbindung des Clients unerwartet beendet wird. Verwenden Sie dies, um andere Agenten darüber zu informieren, dass ein Sensoragent offline gegangen ist, damit sie in einen Fallback-Modus wechseln können.

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()

Best Practices für das Topic-Design

Ein gutes Topic-Design macht Agentencode wartbar und skalierbar. Verwenden Sie eine Hierarchie: location/device_type/device_id/measurement. Vermeiden Sie Leerzeichen und Sonderzeichen. Halten Sie Topics kurz — sie erzeugen bei jeder Nachricht zusätzlichen Overhead.

# 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

Sauberes Trennen und Wiederverbinden

Agenten in IoT-Umgebungen müssen Netzwerkunterbrechungen zuverlässig handhaben. Verwenden Sie die integrierte Wiederverbindungslogik von paho-mqtt: Setzen Sie reconnect_on_failure=True und abonnieren Sie bei jeder Wiederverbindung die Topics erneut, da Abonnements standardmäßig nicht gespeichert werden (außer bei Verwendung persistenter Sitzungen).

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()

Wissenscheck

Welche QoS-Stufe garantiert, dass eine Nachricht genau einmal zugestellt wird?

Zusammenfassung: MQTT für die Agentenintegration

Sehr gut! Das haben Sie in dieser Lektion gelernt:

  • paho-mqtt: connect(), loop_start(), subscribe(), publish()
  • QoS-Stufen: 0 = höchstens einmal, 1 = mindestens einmal, 2 = genau einmal
  • Retained Messages: Der Broker speichert den letzten Wert; neue Abonnenten erhalten ihn sofort
  • LWT: Der Broker veröffentlicht bei einer unerwarteten Verbindungsunterbrechung eine Offline-Nachricht
  • Wiederverbindung: Erneut in on_connect abonnieren; reconnect_delay_set verwenden

Als Nächstes: die Verarbeitung von Zeitreihendatenströmen — gleitende Fenster, gleitende Mittelwerte und Spike-Erkennung.

Häufig gestellte Fragen

Ist die Lektion „MQTT-Protokoll zur Agentenintegration“ kostenlos?

Ja — der vollständige Text von „MQTT-Protokoll zur Agentenintegration“ ist hier im Web kostenlos zu lesen. Um sie interaktiv zu üben (integrierter Code-Editor und 24/7 KI-Tutor) und den Rest des AI Agents-Kurses freizuschalten, upgrade auf CoddyKit PRO. Der AI Agents-Kurs umfasst insgesamt 4 Lektionen.

Was lerne ich in „MQTT-Protokoll zur Agentenintegration“?

Broker einrichten, Topics abonnieren und ereignisgesteuerte Agentenaktivierung umsetzen. Du übst AI Agents mit praktischem Code, den du direkt im Browser ausführst, und ein 24/7 KI-Tutor beantwortet deine Fragen während du die Lektion bearbeitest.

Brauche ich Erfahrung, um AI Agents zu starten?

Keine Vorkenntnisse erforderlich. AI Agents auf CoddyKit ist für Anfänger bis fortgeschrittene Lernende strukturiert, sodass du hier starten oder von Anfang an beginnen und in deinem eigenen Tempo voranschreiten kannst. Dies ist Lektion 1 von 4.

Wie lange dauert die Lektion „MQTT-Protokoll zur Agentenintegration“?

Die meisten CoddyKit-Lektionen dauern etwa 5–10 Minuten. Jede ist kompakt und interaktiv, sodass du stetig Fortschritte machst und genau dort weitermachst, wo du aufgehört hast – im Web und in der App.

Kann ich in dieser AI Agents-Lektion Code schreiben und ausführen?

Ja. Jede AI Agents-Lektion enthält einen integrierten Code-Editor, sodass du echten Code direkt in deinem Browser schreibst und ausführst und sofort KI-Feedback erhältst — ohne lokale Einrichtung erforderlich.

Alle Lektionen in diesem Kurs

  1. MQTT-Protokoll zur Agentenintegration
  2. Zeitreihendaten in Agenten verarbeiten
  3. Automatisierte Reaktion auf Sensorereignisse
  4. Bereitstellung leichtgewichtiger Agenten am Edge
← Zurück zu AI Agents