بروتوكول MQTT لتكامل الوكلاء
إعداد الوسيط، والاشتراك في الموضوعات، وتفعيل الوكيل القائم على الرسائل.
بروتوكول MQTT لتكامل الوكلاء درس مجاني في AI Agents على CoddyKit. هذا هو الدرس 1 من أصل 4. يمكنك قراءة الدرس كاملاً أدناه مجاناً — ثم تمرن عليه مباشرة في المتصفح باستخدام محرر أكواد مدمج ومدرس ذكاء اصطناعي متاح 24/7. هذا الدرس جزء من مسار التعلم في AI Agents، وتقدمك يتزامن عبر الويب وتطبيق CoddyKit. تتضمن دورة AI Agents 4 دروس في المجموع.
وكلاء MQTT وإنترنت الأشياء
إن MQTT (Message Queuing Telemetry Transport) بروتوكول خفيف للنشر والاشتراك، صُمم لبيئات إنترنت الأشياء ذات النطاق الترددي المنخفض وزمن الاستجابة المرتفع. يمكن للوكلاء الاشتراك في موضوعات 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. استخدم # كحرف بدل لمطابقة جميع الموضوعات الفرعية، أو + لمطابقة مستوى واحد فقط. أرفق دالة callback باسم 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
يشترك وكيل المستشعر في موضوعات المستشعرات الأولية، ويتحقق من صحة البيانات الواردة، ويقرر ما إذا كان ينبغي تشغيل إجراء. يستخدم منطق اتخاذ القرار نموذج 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) للوحات المعلومات والوكلاء المستندين إلى المتصفح الاتصال بوسطاء 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')رسالة الوصية الأخيرة
تتيح رسالة الوصية الأخيرة (LWT) في MQTT للوسيط نشر رسالة نيابةً عن العميل إذا انقطع اتصال العميل بشكل غير متوقع. استخدم ذلك لإبلاغ الوكلاء الآخرين بأن وكيل المستشعر أصبح خارج الخدمة، حتى يتمكنوا من التبديل إلى وضع بديل.
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 يضمن تسليم الرسالة مرة واحدة بالضبط؟
مراجعة: MQTT لتكامل الوكلاء
ممتاز! إليك ما تعلمته في هذا الدرس:
- paho-mqtt:
connect()وloop_start()وsubscribe()وpublish() - مستويات QoS: 0 = مرة واحدة كحد أقصى، و1 = مرة واحدة على الأقل، و2 = مرة واحدة بالضبط
- الرسائل المحتفَظ بها: يخزن الوسيط القيمة الأخيرة، ويتلقاها المشتركون الجدد فورًا
- LWT: ينشر الوسيط رسالة تفيد بأن العميل أصبح خارج الخدمة عند انقطاع الاتصال بشكل غير متوقع
- إعادة الاتصال: أعد الاشتراك في
on_connect، واستخدمreconnect_delay_set
التالي: معالجة تدفقات بيانات السلاسل الزمنية — النوافذ المتحركة، والمتوسطات المتحركة، واكتشاف القفزات.
الأسئلة الشائعة
هل درس «بروتوكول MQTT لتكامل الوكلاء» مجاني؟
نعم — نص درس «بروتوكول MQTT لتكامل الوكلاء» كامل متاح مجاناً هنا على الويب. لتمرينه بشكل تفاعلي (محرر أكواد مدمج ومدرس ذكاء اصطناعي متاح 24/7) وفتح باقي دورة AI Agents، انتقل إلى CoddyKit PRO. تتضمن دورة AI Agents 4 دروس في المجموع.
ماذا ستتعلم في «بروتوكول MQTT لتكامل الوكلاء»؟
إعداد الوسيط، والاشتراك في الموضوعات، وتفعيل الوكيل القائم على الرسائل. تتمرن على AI Agents مع أكواد عملية تشغلها مباشرة في المتصفح، ومدرس ذكاء اصطناعي متاح 24/7 يجيب على أسئلتك أثناء عملك.
هل أحتاج إلى خبرة سابقة لأبدأ AI Agents؟
لا تُشترط خبرة سابقة. AI Agents على CoddyKit منظم للمبتدئين حتى المتقدمين، لذا يمكنك البدء من هنا أو من البداية والتقدم بسرعتك الخاصة. هذا هو الدرس 1 من أصل 4.
كم من الوقت يستغرق درس «بروتوكول MQTT لتكامل الوكلاء»؟
معظم دروس CoddyKit تستغرق حوالي 5–10 دقائق. كل منها موجز وتفاعلي، لذا تحرز تقدماً مستمراً وتستأنف من حيث توقفت عبر الويب والتطبيق.
هل يمكنني كتابة وتشغيل أكواد في درس AI Agents هذا؟
نعم. كل درس في AI Agents يتضمن محرر أكواد مدمج، لذا تكتب وتشغل أكواداً حقيقية مباشرة في متصفحك وتحصل على تعليقات فورية من الذكاء الاصطناعي — بدون إعداد محلي.
جميع الدروس في هذه الدورة
- بروتوكول MQTT لتكامل الوكلاء
- معالجة بيانات السلاسل الزمنية لدى الوكلاء
- الاستجابة الآلية لأحداث المستشعرات
- نشر الوكلاء خفيفي الوزن على الحافة