에이전트 통합을 위한 MQTT 프로토콜
브로커 설정, 주제 구독, 메시지 기반 에이전트 활성화를 학습합니다.
에이전트 통합을 위한 MQTT 프로토콜은(는) CoddyKit의 무료 AI Agents 강의입니다. 이것은 4개 중 1번째 강의입니다. 아래에서 전체 강의를 무료로 읽을 수 있으며, 내장 코드 에디터와 24/7 AI 튜터와 함께 브라우저에서 직접 실습할 수 있습니다. 이 강의는 AI Agents 학습 경로의 일부이며, 진행 상황이 웹과 CoddyKit 앱에 동기화됩니다. AI Agents 강의에는 총 4개의 강의가 포함되어 있습니다.
MQTT 및 사물 인터넷 에이전트
MQTT(메시지 큐 원격 측정 전송)는 낮은 대역폭과 높은 지연 시간이 특징인 사물 인터넷 환경을 위해 설계된 경량 발행-구독 프로토콜입니다. 에이전트는 MQTT 주제를 구독하여 실시간 센서 데이터를 받고, 장치로 명령을 다시 발행할 수 있습니다.
paho-mqtt 설치
paho-mqtt 라이브러리는 파이썬의 표준 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')보존 메시지
보존 메시지는 브로커가 저장해 두었다가 새로운 구독자에게 즉시 보내는 해당 토픽의 마지막 발행 값입니다. 센서 상태 토픽에 이상적입니다. 네트워크에 새로 참여한 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')MQTT 센서 Agent 구축
센서 Agent는 원시 센서 토픽을 구독하고, 수신 데이터를 검증하며, 동작을 실행할지 결정합니다. 결정 로직은 복잡한 경우에만 LLM을 사용하고, 간단한 임계값 확인은 속도를 위해 순수한 파이썬으로 처리합니다.
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)를 사용하면 브라우저 기반 대시보드와 Agent가 네이티브 TCP 연결 없이 MQTT 브로커에 연결할 수 있습니다. transport='websockets' 옵션으로 WebSocket 전송을 사용하도록 paho-mqtt를 구성하십시오.
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의 최종 유언 및 유언장 기능을 사용하면 클라이언트가 예기치 않게 연결을 끊었을 때 브로커가 클라이언트를 대신하여 메시지를 발행할 수 있습니다. 이를 사용해 센서 Agent가 오프라인 상태가 되었음을 다른 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정상적인 연결 해제 및 재연결
IoT 환경의 Agent는 네트워크 중단을 적절하게 처리해야 합니다. 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 = 정확히 한 번
- 보존 메시지: 브로커가 마지막 값을 저장하며, 새 구독자가 즉시 수신함
- 최종 유언 및 유언장: 예기치 않은 연결 해제 시 브로커가 오프라인 메시지를 발행함
- 재연결:
on_connect에서 다시 구독하고reconnect_delay_set을 사용함
다음 주제: 시계열 데이터 스트림 처리 — 롤링 윈도, 이동 평균 및 급증 감지
자주 묻는 질문
“에이전트 통합을 위한 MQTT 프로토콜” 강의는 무료인가요?
네 — “에이전트 통합을 위한 MQTT 프로토콜” 전체 내용을 이 웹사이트에서 무료로 읽을 수 있습니다. 인터랙티브하게 실습하려면(내장 코드 에디터와 24/7 AI 튜터), CoddyKit PRO로 업그레이드하면 AI Agents 강의 전체를 잠금 해제할 수 있습니다. AI Agents 강의에는 총 4개의 강의가 포함되어 있습니다.
“에이전트 통합을 위한 MQTT 프로토콜”에서 뭘 배우나요?
브로커 설정, 주제 구독, 메시지 기반 에이전트 활성화를 학습합니다. 브라우저에서 직접 실행하는 실습 코드로 AI Agents을(를) 배우며, 24/7 AI 튜터가 강의를 진행하면서 질문에 답변해줍니다.
AI Agents을(를) 시작하는 데 경험이 필요한가요?
사전 경험은 필요하지 않습니다. CoddyKit의 AI Agents은(는) 초급자부터 고급 학습자까지를 위해 구성되어 있으므로, 여기서 시작하거나 처음부터 시작할 수 있으며 자신의 속도대로 진행할 수 있습니다. 이것은 4개 중 1번째 강의입니다.
“에이전트 통합을 위한 MQTT 프로토콜” 강의는 얼마나 걸리나요?
대부분의 CoddyKit 강의는 약 5~10분이 소요됩니다. 각 강의는 간결하고 인터랙티브하여 꾸준한 진행이 가능하며, 웹과 앱에서 중단한 부분부터 바로 시작할 수 있습니다.
이 AI Agents 강의에서 코드를 작성하고 실행할 수 있나요?
네. 모든 AI Agents 강의에는 내장 코드 에디터가 포함되어 있으므로, 브라우저에서 바로 실제 코드를 작성하고 실행한 후 즉시 AI 피드백을 받을 수 있습니다 — 로컬 설정이 필요 없습니다.
이 강의의 모든 강의
- 에이전트 통합을 위한 MQTT 프로토콜
- 에이전트의 시계열 데이터 처리
- 센서 이벤트에 대한 자동 응답
- 경량 에이전트의 엣지 배포