0Pricing
AI Agents · Leçon

Protocole MQTT pour l’intégration d’agents

Configuration du courtier, abonnement aux sujets et activation des agents par les messages.

Protocole MQTT pour l’intégration d’agents est une leçon AI Agents gratuite sur CoddyKit. Ceci est la leçon 1 sur 4. Tu peux lire la leçon complète ci-dessous gratuitement — puis la pratiquer en direct dans le navigateur avec un éditeur de code intégré et un tuteur IA 24/7. Elle fait partie du parcours d'apprentissage AI Agents, et ta progression se synchronise sur le web et l'application CoddyKit. Le cours AI Agents comprend 4 leçons au total.

MQTT et agents IdO

MQTT (transport de télémétrie par mise en file de messages) est un protocole léger de publication-abonnement conçu pour les environnements IdO à faible bande passante et à latence élevée. Les agents peuvent s’abonner à des sujets MQTT pour recevoir les données des capteurs en temps réel et publier des commandes à destination des appareils.

Installer paho-mqtt

La bibliothèque paho-mqtt est le client MQTT Python standard. Installez-la avec pip. Elle prend en charge MQTT 3.1.1 et 5.0, le chiffrement TLS et les trois niveaux de qualité de service.

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

Se connecter au courtier

Connectez-vous à un courtier MQTT (par exemple Mosquitto, HiveMQ, EMQX ou un courtier cloud). L’appel connect est non bloquant ; utilisez loop_start() pour exécuter la boucle réseau dans un thread d’arrière-plan.

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

S’abonner à une rubrique

Les rubriques sont des chaînes hiérarchiques comme sensors/temperature ou factory/line1/pressure. Utilisez # comme caractère générique pour toutes les sous-rubriques, ou + comme caractère générique pour un seul niveau. Associez un rappel on_message pour traiter les messages entrants.

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

Explication des niveaux de QoS

MQTT propose trois niveaux de qualité de service : QoS 0 — envoyer sans attendre (le plus rapide, mais des messages peuvent être perdus), QoS 1 — au moins une fois (message confirmé, mais potentiellement dupliqué), QoS 2 — exactement une fois (garanti, mais le plus lent). Choisissez selon votre tolérance à la perte de données et à la latence.

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

Messages conservés

Un message conservé est la dernière valeur publiée pour une rubrique, que le courtier stocke et envoie immédiatement à tout nouvel abonné. Cette fonctionnalité est idéale pour les rubriques d’état des capteurs : une nouvelle instance d’agent qui rejoint le réseau connaît immédiatement la valeur actuelle du capteur, sans attendre la prochaine mise à jour.

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

Créer un agent de capteur MQTT

Un agent de capteur s’abonne aux rubriques brutes des capteurs, valide les données entrantes et décide s’il faut déclencher une action. La logique de décision utilise le LLM uniquement pour les cas complexes ; les vérifications simples de seuil sont traitées directement en Python pour plus de rapidité.

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 sur WebSocket

MQTT sur WebSocket (port 8083 ou 8084 pour TLS) permet aux tableaux de bord et aux agents exécutés dans un navigateur de se connecter aux courtiers MQTT sans connexion TCP native. Configurez paho-mqtt pour utiliser le transport WebSocket avec l’option 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')

Testament de dernière volonté

Le testament de dernière volonté (LWT) de MQTT permet au courtier de publier un message au nom d’un Client si celui-ci se déconnecte de manière inattendue. Utilisez cette fonctionnalité pour informer les autres agents qu’un agent de capteur est hors ligne, afin qu’ils puissent passer en mode de secours.

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

Bonnes pratiques de conception des rubriques

Une bonne conception des rubriques rend le code des agents plus facile à maintenir et à faire évoluer. Utilisez une hiérarchie : location/device_type/device_id/measurement. Évitez les espaces et les caractères spéciaux. Gardez les rubriques courtes : elles ajoutent une surcharge à chaque message.

# 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

Déconnexion et reconnexion propres

Les agents des environnements IoT doivent gérer les interruptions réseau avec élégance. Utilisez la logique de reconnexion intégrée à paho-mqtt : définissez reconnect_on_failure=True et réabonnez-vous à chaque reconnexion, car les abonnements ne sont pas conservés par défaut (sauf avec des sessions persistantes).

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

Vérification des connaissances

Quel niveau de QoS garantit qu’un message est livré exactement une fois ?

Récapitulatif : MQTT pour l’intégration des agents

Excellent ! Voici ce que vous avez appris dans cette leçon :

  • paho-mqtt : connect(), loop_start(), subscribe(), publish()
  • Niveaux de QoS : 0 = au plus une fois, 1 = au moins une fois, 2 = exactement une fois
  • Messages conservés : le courtier stocke la dernière valeur ; les nouveaux abonnés la reçoivent immédiatement
  • LWT : le courtier publie un message hors ligne en cas de déconnexion inattendue
  • Reconnexion : réabonnez-vous dans on_connect ; utilisez reconnect_delay_set

Ensuite : traiter des flux de données de séries temporelles — fenêtres glissantes, moyennes mobiles et détection des pics.

Questions Fréquemment Posées

La leçon « Protocole MQTT pour l’intégration d’agents » est-elle gratuite ?

Oui — le texte complet de « Protocole MQTT pour l’intégration d’agents » est gratuit à lire ici sur le web. Pour la pratiquer de manière interactive (un éditeur de code intégré et un tuteur IA 24/7) et déverrouiller le reste du cours AI Agents, passe à CoddyKit PRO. Le cours AI Agents comprend 4 leçons au total.

Qu'est-ce que j'apprendrai dans « Protocole MQTT pour l’intégration d’agents » ?

Configuration du courtier, abonnement aux sujets et activation des agents par les messages. Tu pratiques AI Agents avec du code pratique que tu exécutes directement dans le navigateur, et un tuteur IA 24/7 répond à tes questions au fur et à mesure que tu avances dans la leçon.

Dois-je avoir de l'expérience pour commencer AI Agents ?

Aucune expérience préalable n'est requise. AI Agents sur CoddyKit est structuré pour les débutants jusqu'aux apprenants avancés, donc tu peux commencer ici ou depuis le début et avancer à ton rythme. Ceci est la leçon 1 sur 4.

Combien de temps prend la leçon « Protocole MQTT pour l’intégration d’agents » ?

La plupart des leçons CoddyKit prennent environ 5–10 minutes. Chacune est courte et interactive, tu progresses régulièrement et tu repiques exactement où tu t'es arrêté sur le web et l'app.

Peux-tu écrire et exécuter du code dans cette leçon AI Agents ?

Oui. Chaque leçon AI Agents inclut un éditeur de code intégré, tu écris et exécutes du vrai code directement dans ton navigateur et tu reçois des retours IA instantanés — aucune configuration locale requise.

Toutes les leçons de ce cours

  1. Protocole MQTT pour l’intégration d’agents
  2. Traitement des séries temporelles par les agents
  3. Réponse automatisée aux événements des capteurs
  4. Déploiement en périphérie d’agents légers
← Retour à AI Agents