用于智能体集成的 MQTT 协议
代理设置、主题订阅和消息驱动的智能体激活
用于智能体集成的 MQTT 协议 是 CoddyKit 上的免费 AI Agents 课时。 这是第 1 节课,共 4 节。 你可以在下方免费阅读本课时的完整内容 — 然后在浏览器中使用内置代码编辑器和全天候 AI 导师进行实践。 这是 AI Agents 学习路径的一部分,你的进度在网页和 CoddyKit 应用中同步。 AI Agents 课程共包含 4 节课。
MQTT 与物联网代理
MQTT(消息队列遥测传输)是一种轻量级发布-订阅协议,专为低带宽、高延迟的物联网环境设计。代理可以订阅 MQTT 主题以接收实时传感器数据,并向设备发布命令。
安装 paho-mqtt
paho-mqtt 库是标准的 Python MQTT 客户端。请使用 pip 安装它。该库支持 MQTT 3.1.1 和 5.0、TLS 加密以及全部三种 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)连接到消息代理
连接到 MQTT 消息代理(例如 Mosquitto、HiveMQ、EMQX 或云消息代理)。connect 调用不会阻塞;请使用 loop_start() 在后台线程中运行网络循环。
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订阅主题
主题是分层字符串,例如 sensors/temperature 或 factory/line1/pressure。使用 # 作为通配符来匹配所有子主题,或使用 + 作为通配符来匹配单个层级。附加一个 on_message 回调来处理传入消息。
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 级别详解
MQTT 有三个服务质量级别:QoS 0——发出后不再确认(速度最快,但消息可能丢失);QoS 1——至少一次(消息会得到确认,但可能重复);QoS 2——恰好一次(有保证,但速度最慢)。请根据您对数据丢失与延迟的容忍程度进行选择。
# 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')保留消息
保留消息是某个主题最后一次发布的值,消息代理会存储该值,并立即发送给任何新订阅者。这非常适合传感器状态主题:新的智能体实例加入网络后,无需等待下一次更新,就能立即获知当前传感器值。
# 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')构建 MQTT 传感器 Agent
传感器智能体会订阅原始传感器主题、验证传入数据,并决定是否触发操作。决策逻辑仅在复杂情况下使用 LLM;简单的阈值检查则由纯 Python 处理,以提高速度。
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 = []通过 WebSocket 使用 MQTT
通过 WebSocket 使用 MQTT(TLS 使用端口 8083 或 8084),可以让基于浏览器的仪表板和智能体连接 MQTT 消息代理,而无需原生 TCP 连接。请配置 paho-mqtt,通过 transport='websockets' 选项使用 WebSocket 传输。
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')遗嘱消息
MQTT 的遗嘱消息(LWT)允许消息代理在客户端意外断开连接时,代表该客户端发布一条消息。您可以利用这一机制通知其他智能体某个传感器智能体已离线,以便它们切换到后备模式。
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()主题设计最佳实践
良好的主题设计可以让智能体代码更易维护和扩展。请使用分层结构:location/device_type/device_id/measurement。避免使用空格和特殊字符。主题应保持简短——每条消息都会因主题而增加额外开销。
# 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干净断开与重新连接
物联网环境中的智能体必须妥善处理网络中断。请使用 paho-mqtt 内置的重新连接逻辑:设置 reconnect_on_failure=True,并在每次重新连接时重新订阅,因为订阅默认不会持久化(除非使用持久会话)。
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()知识检查
哪个 QoS 级别能够保证消息恰好传递一次?
回顾:用于 Agent 集成的 MQTT
非常好!本课中您学到的内容:
- paho-mqtt:
connect()、loop_start()、subscribe()、publish() - QoS 级别:0 = 最多一次,1 = 至少一次,2 = 恰好一次
- 保留消息:消息代理存储最后一个值;新订阅者会立即收到该值
- LWT:意外断开时,消息代理发布离线消息
- 重新连接:在
on_connect中重新订阅;使用reconnect_delay_set
下一部分:处理时间序列数据流——滚动窗口、移动平均值和尖峰检测。
常见问题解答
「用于智能体集成的 MQTT 协议」课时是免费的吗?
是的 — 「用于智能体集成的 MQTT 协议」的完整文本可在网页上免费阅读。要进行交互式练习(内置代码编辑器和全天候 AI 导师)并解锁 AI Agents 课程的其余内容,请升级到 CoddyKit PRO。 AI Agents 课程共包含 4 节课。
「用于智能体集成的 MQTT 协议」这节课中我会学到什么?
代理设置、主题订阅和消息驱动的智能体激活 你通过在浏览器中直接运行的动手代码来练习 AI Agents,全天候 AI 导师会在你学习这节课的过程中回答你的问题。
学习 AI Agents 需要有经验吗?
无需任何先前经验。CoddyKit 上的 AI Agents 课程适合初学者到高级学习者,你可以从这里开始或从头开始,按照自己的节奏学习。 这是第 1 节课,共 4 节。
「用于智能体集成的 MQTT 协议」课时需要多长时间?
大多数 CoddyKit 课程大约需要 5–10 分钟。每节课都很精短且互动,所以你能稳步进步,并在网页和应用中从离开的地方继续。
我能在这节 AI Agents 课中编写并运行代码吗?
能。每节 AI Agents 课都包含内置代码编辑器,你可以在浏览器中直接编写并运行真实代码,并获得即时 AI 反馈 — 无需本地设置。
此课程中的所有课时
- 用于智能体集成的 MQTT 协议
- 智能体中的时间序列数据处理
- 对传感器事件的自动响应
- 轻量级智能体的边缘部署