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 connectionEin 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/celsiusSauberes 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_connectabonnieren;reconnect_delay_setverwenden
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
- MQTT-Protokoll zur Agentenintegration
- Zeitreihendaten in Agenten verarbeiten
- Automatisierte Reaktion auf Sensorereignisse
- Bereitstellung leichtgewichtiger Agenten am Edge