0Pricing
AI Agents · บทเรียน

โพรโทคอล MQTT สำหรับการผสานเอเจนต์

ตั้งค่าตัวกลาง สมัครรับหัวข้อ และเปิดใช้งานเอเจนต์ด้วยข้อความ

โพรโทคอล MQTT สำหรับการผสานเอเจนต์ เป็นบทเรียน AI Agents ฟรีบน CoddyKit นี่คือบทเรียนที่ 1 จากทั้งหมด 4 บทเรียน คุณสามารถอ่านบทเรียนทั้งหมดด้านล่างฟรี — จากนั้นลองปฏิบัติด้วยตัวคุณเองในเบราว์เซอร์พร้อมตัวแก้ไขโค้ดในตัวและติวเตอร์ AI ตลอด 24/7 บทเรียนนี้เป็นส่วนหนึ่งของเส้นทางการเรียน AI Agents และความก้าวหน้าของคุณจะซิงค์ข้ามเว็บและแอป CoddyKit คอร์ส AI Agents มีบทเรียนทั้งหมด 4 บทเรียน

MQTT และเอเจนต์ IoT

MQTT (โพรโทคอลขนส่งข้อมูลการวัดระยะไกลแบบคิวข้อความ) เป็นโพรโทคอลเผยแพร่-สมัครรับข้อมูลขนาดเล็กที่ออกแบบมาสำหรับสภาพแวดล้อม IoT ซึ่งมีแบนด์วิดท์ต่ำและความหน่วงสูง เอเจนต์สามารถสมัครรับหัวข้อ MQTT เพื่อรับข้อมูลเซนเซอร์แบบเรียลไทม์ และเผยแพร่คำสั่งกลับไปยังอุปกรณ์ได้

การติดตั้ง paho-mqtt

ไลบรารี paho-mqtt เป็นไคลเอนต์ MQTT มาตรฐานสำหรับ Python ให้ติดตั้งด้วย 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')

ข้อความที่เก็บไว้

ข้อความที่เก็บไว้ คือค่าล่าสุดที่เผยแพร่สำหรับหัวข้อ ซึ่งโบรกเกอร์จะจัดเก็บไว้และส่งให้ผู้สมัครรับรายใหม่ทันที เหมาะอย่างยิ่งสำหรับหัวข้อสถานะของเซนเซอร์ เพราะเมื่อ Agent อินสแตนซ์ใหม่เข้าร่วมเครือข่าย ก็จะทราบค่าปัจจุบันของเซนเซอร์ทันทีโดยไม่ต้องรอการอัปเดตครั้งถัดไป

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

การสร้าง Agent เซนเซอร์ด้วย 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 = []

MQTT ผ่าน WebSocket

MQTT ผ่าน WebSocket (พอร์ต 8083 หรือ 8084 สำหรับ TLS) ช่วยให้แดชบอร์ดและ Agent ที่ทำงานบนเบราว์เซอร์เชื่อมต่อกับโบรกเกอร์ MQTT ได้โดยไม่ต้องใช้การเชื่อมต่อ TCP โดยตรง กำหนดค่า paho-mqtt ให้ใช้การขนส่งผ่าน WebSocket ด้วยตัวเลือก 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')

พินัยกรรมฉบับสุดท้าย

พินัยกรรมฉบับสุดท้าย ของ MQTT ช่วยให้โบรกเกอร์เผยแพร่ข้อความในนามของ Client หาก Client ยกเลิกการเชื่อมต่อโดยไม่คาดคิด ใช้คุณสมบัตินี้เพื่อแจ้งให้ Agent อื่นทราบว่า Agent เซนเซอร์ออฟไลน์แล้ว เพื่อให้เปลี่ยนไปใช้โหมดสำรองได้

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

แนวทางปฏิบัติที่ดีในการออกแบบหัวข้อ

การออกแบบหัวข้อที่ดีช่วยให้โค้ดของ Agent ดูแลรักษาและขยายระบบได้ง่าย ใช้ลำดับชั้นดังนี้: 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

การยกเลิกการเชื่อมต่อและการเชื่อมต่อใหม่อย่างถูกต้อง

Agent ในสภาพแวดล้อม IoT ต้องรับมือกับการหยุดชะงักของเครือข่ายอย่างเหมาะสม ใช้ตรรกะการเชื่อมต่อใหม่ที่มีอยู่ใน 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 ใดรับประกันว่าข้อความจะถูกส่ง ครั้งเดียวเท่านั้น?

ทบทวน: MQTT สำหรับการผสานรวม Agent

ยอดเยี่ยม! สิ่งที่คุณได้เรียนรู้ในบทเรียนนี้:

  • paho-mqtt: connect(), loop_start(), subscribe(), publish()
  • ระดับ QoS: 0 = ส่งไม่เกินหนึ่งครั้ง, 1 = ส่งอย่างน้อยหนึ่งครั้ง, 2 = ส่งครั้งเดียวเท่านั้น
  • ข้อความที่เก็บไว้: โบรกเกอร์จัดเก็บค่าล่าสุดไว้ และผู้สมัครรับรายใหม่จะได้รับค่าดังกล่าวทันที
  • พินัยกรรมฉบับสุดท้าย: โบรกเกอร์เผยแพร่ข้อความแจ้งออฟไลน์เมื่อมีการยกเลิกการเชื่อมต่อโดยไม่คาดคิด
  • การเชื่อมต่อใหม่: สมัครรับหัวข้อใหม่ใน on_connect และใช้ reconnect_delay_set

ถัดไป: การประมวลผลสตรีมข้อมูลอนุกรมเวลา — หน้าต่างเลื่อน ค่าเฉลี่ยเคลื่อนที่ และการตรวจจับจุดพุ่งสูง

คำถามที่พบบ่อย

บทเรียน “โพรโทคอล MQTT สำหรับการผสานเอเจนต์” ฟรีหรือไม่

ใช่ — ข้อความเต็มของ “โพรโทคอล MQTT สำหรับการผสานเอเจนต์” ฟรีให้อ่านที่นี่บนเว็บ เพื่อปฏิบัติแบบโต้ตอบ (ตัวแก้ไขโค้ดในตัวและติวเตอร์ AI ตลอด 24/7) และปลดล็อคส่วนที่เหลือของคอร์ส AI Agents ให้อัปเกรดเป็น CoddyKit PRO คอร์ส AI Agents มีบทเรียนทั้งหมด 4 บทเรียน

คุณจะเรียนรู้อะไรในบทเรียน “โพรโทคอล MQTT สำหรับการผสานเอเจนต์”

ตั้งค่าตัวกลาง สมัครรับหัวข้อ และเปิดใช้งานเอเจนต์ด้วยข้อความ คุณปฏิบัติ AI Agents ด้วยโค้ดที่ใช้งานได้จริงที่คุณเรียกใช้โดยตรงในเบราว์เซอร์ และติวเตอร์ AI ตลอด 24/7 ตอบคำถามของคุณขณะที่คุณไปผ่านบทเรียน

คุณต้องมีประสบการณ์ก่อนที่จะเริ่มเรียน AI Agents หรือไม่

ไม่จำเป็นต้องมีประสบการณ์มาก่อน AI Agents บน CoddyKit ออกแบบมาสำหรับผู้เริ่มต้นไปจนถึงผู้เรียนขั้นสูง คุณสามารถเริ่มต้นที่นี่หรือเริ่มจากตัวแรกและเรียนด้วยความเร็วของคุณเอง นี่คือบทเรียนที่ 1 จากทั้งหมด 4 บทเรียน

บทเรียน “โพรโทคอล MQTT สำหรับการผสานเอเจนต์” ใช้เวลานานแค่ไหน

บทเรียน CoddyKit ส่วนใหญ่ใช้เวลาประมาณ 5–10 นาที แต่ละบทเรียนจึงสั้นและเป็นแบบโต้ตอบ คุณสามารถก้าวหน้าอย่างต่อเนื่องและกลับมาเรียนต่อจากตรงที่เพิ่งหยุดบนเว็บและแอปได้เลย

ฉันเขียนและรันโค้ดในบทเรียน AI Agents นี้ได้ไหม

ได้ บทเรียน AI Agents ทุกบทมีตัวแก้ไขโค้ดในตัว คุณจึงเขียนและรันโค้ดจริงได้เลยในเบราว์เซอร์ และได้รับข้อเสนอแนะจาก AI ในทันที — ไม่ต้องติดตั้งในเครื่องของคุณ

บทเรียนทั้งหมดในหลักสูตรนี้

  1. โพรโทคอล MQTT สำหรับการผสานเอเจนต์
  2. การประมวลผลข้อมูลอนุกรมเวลาในเอเจนต์
  3. การตอบสนองอัตโนมัติต่อเหตุการณ์จากเซนเซอร์
  4. การนำเอเจนต์น้ำหนักเบาไปใช้งานที่เอดจ์
← กลับไปที่ AI Agents