0Pricing
AI Agents · 课时

用于智能体集成的 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 反馈 — 无需本地设置。

此课程中的所有课时

  1. 用于智能体集成的 MQTT 协议
  2. 智能体中的时间序列数据处理
  3. 对传感器事件的自动响应
  4. 轻量级智能体的边缘部署
← 返回 AI Agents