Protocolo MQTT para la integración de agentes
Configuración del broker, suscripción a temas y activación de agentes basada en mensajes.
Protocolo MQTT para la integración de agentes es una lección gratuita de AI Agents en CoddyKit. Esta es la lección 1 de 4. Puedes leer la lección completa abajo gratuitamente — luego la practicas en el navegador con un editor de código integrado y un tutor de IA 24/7. Forma parte de la ruta de aprendizaje de AI Agents, y tu progreso se sincroniza en la web y la app de CoddyKit. El curso de AI Agents incluye 4 lecciones en total.
MQTT y agentes de IoT
MQTT (Message Queuing Telemetry Transport) es un protocolo ligero de publicación y suscripción diseñado para entornos de IoT con poco ancho de banda y alta latencia. Los agentes pueden suscribirse a temas MQTT para recibir datos de sensores en tiempo real y publicar comandos de vuelta a los dispositivos.
Instalación de paho-mqtt
La biblioteca paho-mqtt es el cliente MQTT estándar para Python. Instálela con pip. Es compatible con MQTT 3.1.1 y 5.0, con cifrado TLS y con los tres niveles de 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)Conexión al broker
Conéctese a un broker MQTT (por ejemplo, Mosquitto, HiveMQ, EMQX o un broker en la nube). La llamada a connect no bloquea la ejecución; utilice loop_start() para ejecutar el bucle de red en un hilo en segundo plano.
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 connectionSuscribirse a un topic
Los topics son cadenas jerárquicas como sensors/temperature o factory/line1/pressure. Use # como comodín para todos los subtopics, o + como comodín de un solo nivel. Asocie un callback on_message para procesar los mensajes entrantes.
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}')Niveles de QoS explicados
MQTT tiene tres niveles de calidad de servicio: QoS 0 — enviar y olvidar (el más rápido, pero puede perder mensajes), QoS 1 — al menos una vez (el mensaje se confirma, pero puede duplicarse), QoS 2 — exactamente una vez (garantizado, pero el más lento). Elija el nivel según su tolerancia a la pérdida de datos frente a la latencia.
# 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')Mensajes retenidos
Un mensaje retenido es el último valor publicado para un topic que el broker almacena y envía inmediatamente a cualquier suscriptor nuevo. Es ideal para los topics que representan el estado de sensores: una instancia nueva del agente que se incorpora a la red conoce de inmediato el valor actual del sensor sin tener que esperar a la siguiente actualización.
# 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')Construcción de un agente sensor MQTT
Un agente sensor se suscribe a topics de sensores sin procesar, valida los datos entrantes y decide si debe activar una acción. La lógica de decisión usa el LLM solo para los casos complejos; las comprobaciones sencillas de umbrales se ejecutan en Python puro para obtener mayor velocidad.
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 mediante WebSocket
MQTT mediante WebSocket (puerto 8083 u 8084 para TLS) permite que los dashboards y agentes basados en navegador se conecten a brokers MQTT sin una conexión TCP nativa. Configure paho-mqtt para usar el transporte WebSocket con la opción 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')Última voluntad y testamento
La última voluntad y testamento (LWT) de MQTT permite que el broker publique un mensaje en nombre de un cliente si este se desconecta inesperadamente. Úselo para notificar a otros agentes que un agente sensor se desconectó, de modo que puedan cambiar a un modo alternativo.
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()Buenas prácticas para diseñar topics
Un buen diseño de topics facilita el mantenimiento y la escalabilidad del código de los agentes. Use una jerarquía: location/device_type/device_id/measurement. Evite los espacios y los caracteres especiales. Mantenga los topics cortos: añaden sobrecarga a cada mensaje.
# 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/celsiusDesconexión y reconexión ordenadas
Los agentes en entornos IoT deben gestionar las interrupciones de red correctamente. Use la lógica integrada de reconexión de paho-mqtt: establezca reconnect_on_failure=True y vuelva a suscribirse en cada reconexión, ya que las suscripciones no se conservan de forma predeterminada (salvo cuando se usan sesiones persistentes).
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()Comprobación de conocimientos
¿Qué nivel de QoS garantiza que un mensaje se entregue exactamente una vez?
Repaso: MQTT para la integración de agentes
¡Excelente! Esto es lo que aprendió en esta lección:
- paho-mqtt:
connect(),loop_start(),subscribe(),publish() - Niveles de QoS: 0 = como máximo una vez, 1 = al menos una vez, 2 = exactamente una vez
- Mensajes retenidos: el broker almacena el último valor; los nuevos suscriptores lo reciben de inmediato
- LWT: el broker publica un mensaje de desconexión cuando se produce una desconexión inesperada
- Reconexión: vuelva a suscribirse en
on_connect; usereconnect_delay_set
A continuación: procesamiento de flujos de datos de series temporales — ventanas deslizantes, promedios móviles y detección de picos.
Preguntas frecuentes
¿La lección «Protocolo MQTT para la integración de agentes» es gratis?
Sí — el texto completo de «Protocolo MQTT para la integración de agentes» es gratis para leer aquí en la web. Para practicarla de forma interactiva (editor de código integrado y tutor de IA 24/7) y desbloquear el resto del curso de AI Agents, actualiza a CoddyKit PRO. El curso de AI Agents incluye 4 lecciones en total.
¿Qué aprenderé en «Protocolo MQTT para la integración de agentes»?
Configuración del broker, suscripción a temas y activación de agentes basada en mensajes. Practicas AI Agents con código real que ejecutas directamente en el navegador, y un tutor de IA 24/7 responde tus preguntas mientras trabajas en la lección.
¿Necesito experiencia previa para empezar AI Agents?
No se requiere experiencia previa. AI Agents en CoddyKit está estructurado para principiantes hasta estudiantes avanzados, así que puedes empezar aquí o desde el inicio y avanzar a tu ritmo. Esta es la lección 1 de 4.
¿Cuánto tiempo toma la lección «Protocolo MQTT para la integración de agentes»?
La mayoría de las lecciones de CoddyKit toman alrededor de 5–10 minutos. Cada una es compacta e interactiva, así que avanzas constantemente y retomas exactamente por donde dejaste en la web y la app.
¿Puedo escribir y ejecutar código en esta lección de AI Agents?
Sí. Cada lección de AI Agents incluye un editor de código integrado, así que escribes y ejecutas código real directamente en tu navegador y obtienes retroalimentación instantánea de IA — sin configuración local necesaria.
Todas las lecciones de este curso
- Protocolo MQTT para la integración de agentes
- Procesamiento de datos de series temporales en agentes
- Respuesta automatizada a eventos de sensores
- Despliegue en el edge de agentes ligeros