โพรโทคอล 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 ในทันที — ไม่ต้องติดตั้งในเครื่องของคุณ
บทเรียนทั้งหมดในหลักสูตรนี้
- โพรโทคอล MQTT สำหรับการผสานเอเจนต์
- การประมวลผลข้อมูลอนุกรมเวลาในเอเจนต์
- การตอบสนองอัตโนมัติต่อเหตุการณ์จากเซนเซอร์
- การนำเอเจนต์น้ำหนักเบาไปใช้งานที่เอดจ์