エージェント連携のための MQTT プロトコル
ブローカーの設定、トピックの購読、メッセージ駆動によるエージェントの起動を学びます。
「エージェント連携のための MQTT プロトコル」はCoddyKit上の無料AI Agentsレッスンです。 これはレッスン1/4です。 下記で完全なレッスンを無料で読むことができます。その後、ブラウザ内の組み込みコードエディタと24時間対応のAIチューターでハンズオン演習できます。 これはAI Agents学習パスの一部であり、ウェブとCoddyKitアプリ全体で進捗が同期されます。 AI Agentsコースには全4レッスンが含まれています。
MQTTとIoTエージェント
MQTT(Message Queuing Telemetry Transport)は、低帯域幅・高遅延のIoT環境向けに設計された軽量なパブリッシュ/サブスクライブプロトコルです。エージェントはMQTTトピックをサブスクライブしてセンサーデータをリアルタイムで受信し、デバイスにコマンドをパブリッシュできます。
paho-mqttのインストール
paho-mqttライブラリは、標準的なPython用MQTTクライアントです。pipでインストールできます。MQTT 3.1.1と5.0、TLS暗号化、3つすべての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のような階層構造を持つ文字列です。すべてのサブトピックには#をワイルドカードとして使用し、1階層分のワイルドカードには+を使用します。受信メッセージを処理するには、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には、Quality of Serviceのレベルが3つあります。QoS 0 — 送信後の確認なし(最速ですが、メッセージが失われる可能性があります)、QoS 1 — 少なくとも1回(メッセージは確認されますが、重複する可能性があります)、QoS 2 — 正確に1回(保証されますが、最も遅くなります)。データ損失と遅延のどちらをどの程度許容できるかに基づいて選択します。
# 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')Retainメッセージ
Retainメッセージとは、ブローカーが保存し、新しいサブスクライバーにすぐ送信する、あるトピックで最後にパブリッシュされた値です。センサーの状態を表すトピックに最適です。ネットワークに参加したエージェントの新しいインスタンスは、次の更新を待たずに現在のセンサー値をすぐ把握できます。
# 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 = []WebSocket経由のMQTT
WebSocket経由のMQTT(TLSの場合はポート8083または8084)を使用すると、ブラウザベースのダッシュボードやエージェントが、ネイティブなTCP接続なしでMQTTブローカーに接続できます。transport='websockets'オプションを使って、paho-mqttが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')Last Will and Testament
MQTTのLast Will and Testament (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正常な切断と再接続
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()知識チェック
メッセージが正確に1回配信されることを保証するQoSレベルはどれですか?
まとめ:エージェント統合のためのMQTT
すばらしいです!このレッスンで学んだこと:
- paho-mqtt:
connect()、loop_start()、subscribe()、publish() - QoSレベル:0 = 最大1回、1 = 少なくとも1回、2 = 正確に1回
- Retainメッセージ:ブローカーが最後の値を保存し、新しいサブスクライバーにすぐ配信
- LWT:予期しない切断時にブローカーがオフラインメッセージをパブリッシュ
- 再接続:
on_connectで再サブスクライブし、reconnect_delay_setを使用
次は、時系列データストリームの処理です。ローリングウィンドウ、移動平均、スパイク検出について学びます。
よくある質問
「エージェント連携のための MQTT プロトコル」レッスンは無料ですか?
はい。「エージェント連携のための MQTT プロトコル」の完全なテキストはこのウェブで無料で読めます。インタラクティブに演習し(組み込みコードエディタと24時間対応のAIチューター)、AI Agentsコースの残りをアンロックするには、CoddyKit PROにアップグレードしてください。 AI Agentsコースには全4レッスンが含まれています。
「エージェント連携のための MQTT プロトコル」で何を学びますか?
ブローカーの設定、トピックの購読、メッセージ駆動によるエージェントの起動を学びます。 ブラウザで直接実行するハンズオンコードでAI Agentsを演習し、24時間対応のAIチューターがレッスンを進める中での質問に答えます。
AI Agentsを始めるのに経験は必要ですか?
事前経験は必要ありません。CoddyKitのAI Agentsは初級者から上級者向けに構成されているため、ここから始めるか最初から始めて、自分のペースで進むことができます。 これはレッスン1/4です。
「エージェント連携のための MQTT プロトコル」レッスンにはどのくらい時間がかかりますか?
ほとんどのCoddyKitレッスンは約5~10分かかります。各レッスンはコンパクトでインタラクティブなので、着実に進歩し、ウェブとアプリ全体で正確に前回の場所から再開できます。
このAI Agentsレッスンでコードを書いて実行できますか?
はい。すべてのAI Agentsレッスンに組み込みコードエディタが含まれているため、ブラウザでリアルコードを書いて実行し、即座のAIフィードバックを取得できます。ローカル設定は不要です。
このコースのすべてのレッスン
- エージェント連携のための MQTT プロトコル
- エージェントにおける時系列データ処理
- センサーイベントへの自動応答
- 軽量エージェントのエッジデプロイ