Eclipse Paho MQTT Python Client API

gogobt·2026년 4월 12일


API 조사 활용
Eclipse Paho MQTT Python Client API
MQTTAPI실습
작성자: 전병탁

목차
MQTT 프로토콜 개요 및 이론
REST API vs MQTT 비교
MQTT 패킷 구조 및 QoS 레벨
토픽 네이밍 설계 원칙
실습 1 — 기본 Pub/Sub 동작 확인
실습 2 — QoS 레벨별 동작 비교
실습 3 — 다중 설비 동시 시뮬레이션
실습 4 — Retained Message (최신 상태 유지)
실습 5 — Will Message (설비 이상 종료 감지)
실습 6 — 수신 데이터 시계열 시각화
장점·한계·제조 현장 적용 관련

1. MQTT 프로토콜 개요 및 이론

탄생 배경

MQTT(Message Queuing Telemetry Transport)는 1999년 IBM의 Andy Stanford-Clark과 Arcom의 Arlen Nipper가 사우디아라비아 송유관 원격 모니터링 프로젝트를 위해 설계했다. 위성 통신 환경에서 대역폭이 좁고 배터리 수명이 중요한 조건을 만족해야 했기 때문에, HTTP보다 훨씬 가벼운 프로토콜이 필요했다.

2013년 OASIS 표준(v3.1.1), 2019년 v5.0이 제정되면서 IoT·스마트팩토리의 사실상 표준 프로토콜로 자리잡았다. AWS IoT Core, Azure IoT Hub, Google Cloud IoT 모두 MQTT를 기본 프로토콜로 채택한다.

발행-구독(Pub/Sub) 패턴

[발행자 Publisher][브로커 Broker] [구독자 Subscriber]
(설비 센서, PLC) (Mosquitto / EMQ X) (데이터 수집 서버)
│ │ │
│── PUBLISH topic ──────▶│ │
│ payload, QoS │── 토픽 매칭 ──────────▶│
│ │ MESSAGE 전달 │
│ │ │
브로커(Broker): 중간 메시지 교환소. 발행자와 구독자를 직접 연결하지 않고 토픽으로 라우팅.
토픽(Topic): 메시지가 흐르는 채널명. 슬래시(/)로 계층 표현. 예: factory/line1/CVD001/temperature
페이로드(Payload): 실제 전송 데이터. 형식 제한 없음(JSON, 바이너리, 문자열 모두 가능).

공식 정보

항목내용
공식 사이트https://mqtt.org
Eclipse Pahohttps://eclipse.dev/paho/
Python 패키지https://pypi.org/project/paho-mqtt/
공개 테스트 브로커https://test.mosquitto.org (포트 1883)
OASIS 표준MQTT v3.1.1 (ISO/IEC 20922), v5.0

2. REST API vs MQTT 비교

항목REST APIMQTT
통신 패턴요청-응답(Request-Response)발행-구독(Publish-Subscribe)
연결 방식Stateless (매 요청마다 연결)Persistent TCP 연결 유지
실시간성낮음 (폴링 필요)높음 (이벤트 즉시 Push)
헤더 오버헤드HTTP 헤더 수백~수천 바이트최소 2바이트 고정 헤더
다대다 전달불가 (1:1 요청-응답)가능 (1:N, N:1, N:N)
오프라인 처리없음QoS 1·2 + 세션 지속으로 재전달 보장
적합한 상황CRUD, 데이터 조회·수정실시간 센서, 이벤트 알림, IoT
주요 사용처웹 API, MES 조회 인터페이스FDC 수집, SCADA, 설비 모니터링

3. MQTT 패킷 구조 및 QoS 레벨

MQTT 고정 헤더 구조

Byte 1: [ 패킷 타입(4비트) | 플래그(4비트) ]
Byte 2: [ 나머지 길이(Remaining Length) ]
최소 2바이트만으로 패킷 전송 가능 → HTTP의 수백 바이트 헤더 대비 극도로 경량.

QoS(Quality of Service) 3단계

QoS이름동작전달 보장중복 가능성사용 시나리오
0At most onceFire & Forget❌ 없음없음온도 주기 데이터 (유실 허용)
1At least oncePUBACK 확인✅ 최소 1회있음 (중복 가능)FDC 이상 알람, 공정 이벤트
2Exactly once4-way handshake✅ 정확히 1회없음결제, 공정 레시피 변경 명령

QoS 0: PUB ──────────────────────────▶ SUB (확인 없음)

QoS 1: PUB ── PUBLISH ───────────────▶ SUB
PUB ◀── PUBACK ─────────────── SUB

QoS 2: PUB ── PUBLISH ───────────────▶ SUB
PUB ◀── PUBREC ─────────────── SUB
PUB ── PUBREL ────────────────▶ SUB
PUB ◀── PUBCOMP ────────────── SUB

4. 토픽 네이밍 설계 원칙

토픽은 REST URL 설계처럼 일관된 규칙이 없으면 운영 복잡도가 급격히 증가한다. 제조 현장 적용 시 권장 패턴:

{사이트}/{라인}/{설비ID}/{데이터타입}

예시:
fab1/line3/CVD001/sensor ← 정상 센서 주기 데이터
fab1/line3/CVD001/alarm ← 이상 알람 이벤트
fab1/line3/CVD001/status ← 설비 상태 (Run/Idle/Down)
fab1/line3/+/alarm ← line3 전체 설비 알람 구독 (와일드카드 +)
fab1/# ← fab1 전체 모든 메시지 구독 (와일드카드 #)
와일드카드 규칙
| 와일드카드 | 의미 | 예시 |
| --- | --- | --- |
| + | 한 레벨 임의 대체 | fab1/+/CVD001/alarm → 모든 라인의 CVD001 알람 |
| # | 이하 모든 레벨 | fab1/# → fab1 하위 모든 메시지 |
금지 패턴
❌ /로 시작하는 토픽 → /fab1/CVD001 (불필요한 빈 레벨 생성)
❌ 공백 포함 → fab1/CVD 001
로 시작 → $factory/... (는 브로커 내부 예약)
❌ 너무 깊은 계층 → a/b/c/d/e/f/g/h (5레벨 이내 권장)

# ═══════════════════════════════════════════════════
# 환경 준비 — 패키지 설치 확인
# ═══════════════════════════════════════════════════
import subprocess, sys

pkgs = ['paho-mqtt', 'matplotlib', 'numpy']
for pkg in pkgs:
    subprocess.check_call([sys.executable, '-m', 'pip', 'install', pkg, '-q'])

import paho.mqtt.client as mqtt
import json, time, threading, random
from datetime import datetime
import matplotlib.pyplot as plt
import matplotlib.dates as mdates
import numpy as np

print('✅ 패키지 로드 완료')
print(f'   paho-mqtt 버전: {mqtt.__version__ if hasattr(mqtt, "__version__") else "설치됨"}')

# ───────────────────────────────────────────────────
# 공통 MQTT 설정
# ───────────────────────────────────────────────────
BROKER   = 'test.mosquitto.org'   # Eclipse 재단 운영 공개 브로커
PORT     = 1883                   # 기본 포트 (TLS는 8883)
BASE_TOPIC = 'swe2026_jbt'        # 수업 고유 prefix — 다른 수강생과 토픽 충돌 방지

print(f'\n브로커: {BROKER}:{PORT}')
print(f'토픽 prefix: {BASE_TOPIC}')

실행 결과:

✅ 패키지 로드 완료
   paho-mqtt 버전: 설치됨

브로커: test.mosquitto.org:1883
토픽 prefix: swe2026_jbt

5. 실습 1 — 기본 Pub/Sub 동작 확인

가장 기본적인 발행-구독 흐름을 확인한다.
구독자를 먼저 연결·구독 준비한 뒤, 발행자가 메시지를 보내는 표준 패턴을 익힌다.

# ───────────────────────────────────────────────────
# 실습 1: 기본 Pub/Sub
# ───────────────────────────────────────────────────
TOPIC_BASIC = f'{BASE_TOPIC}/basic/test'

basic_received = []
basic_ready    = threading.Event()

# ── 구독자 콜백 ──
def on_connect_basic(client, userdata, flags, rc):
    if rc == 0:
        client.subscribe(TOPIC_BASIC, qos=1)
        basic_ready.set()
        print(f'[구독자] 연결 성공 → 구독: {TOPIC_BASIC}')

def on_message_basic(client, userdata, msg):
    data = json.loads(msg.payload.decode())
    basic_received.append(data)
    print(f'[구독자] 수신 ← {msg.topic}')
    print(f'         메시지: {json.dumps(data, ensure_ascii=False)}')

# ── 구독자 클라이언트 ──
sub1 = mqtt.Client(client_id='sub_basic_jbt', protocol=mqtt.MQTTv311)
sub1.on_connect = on_connect_basic
sub1.on_message = on_message_basic
sub1.connect(BROKER, PORT, 60)
sub1.loop_start()

if not basic_ready.wait(timeout=6):
    print('⚠️  연결 타임아웃')
else:
    time.sleep(0.3)

    # ── 발행자 ──
    pub1 = mqtt.Client(client_id='pub_basic_jbt', protocol=mqtt.MQTTv311)
    pub1.connect(BROKER, PORT, 60)
    pub1.loop_start()
    time.sleep(0.3)

    messages = [
        {'seq': 1, 'msg': '첫 번째 테스트 메시지', 'ts': datetime.now().isoformat()},
        {'seq': 2, 'msg': 'MQTT Pub/Sub 동작 확인', 'ts': datetime.now().isoformat()},
        {'seq': 3, 'msg': 'MQTT_API_실습', 'ts': datetime.now().isoformat()},
    ]
    for m in messages:
        result = pub1.publish(TOPIC_BASIC, json.dumps(m, ensure_ascii=False), qos=1)
        print(f'[발행자] 송신 → seq={m["seq"]}  (mid={result.mid})')
        time.sleep(0.6)

    time.sleep(1)
    pub1.loop_stop(); pub1.disconnect()

sub1.loop_stop(); sub1.disconnect()
print(f'\n✅ 발행: {len(messages)}건 | 수신: {len(basic_received)}건')

실행 결과:

C:\Users\gogot\AppData\Local\Temp\ipykernel_19760\3536012169.py:23: DeprecationWarning: Callback API version 1 is deprecated, update to latest version
  sub1 = mqtt.Client(client_id='sub_basic_jbt', protocol=mqtt.MQTTv311)
[구독자] 연결 성공 → 구독: swe2026_jbt/basic/test
C:\Users\gogot\AppData\Local\Temp\ipykernel_19760\3536012169.py:35: DeprecationWarning: Callback API version 1 is deprecated, update to latest version
  pub1 = mqtt.Client(client_id='pub_basic_jbt', protocol=mqtt.MQTTv311)
[발행자] 송신 → seq=1  (mid=1)
[구독자] 수신 ← swe2026_jbt/basic/test
         메시지: {"seq": 1, "msg": "첫 번째 테스트 메시지", "ts": "2026-04-12T21:58:32.816737"}
[발행자] 송신 → seq=2  (mid=2)
[구독자] 수신 ← swe2026_jbt/basic/test
         메시지: {"seq": 2, "msg": "MQTT Pub/Sub 동작 확인", "ts": "2026-04-12T21:58:32.816737"}
[발행자] 송신 → seq=3  (mid=3)
[구독자] 수신 ← swe2026_jbt/basic/test
         메시지: {"seq": 3, "msg": "MQTT_API_실습", "ts": "2026-04-12T21:58:32.816737"}

✅ 발행: 3건 | 수신: 3건

6. 실습 2 — QoS 레벨별 동작 비교

QoS 0·1·2 각각으로 메시지를 발행하고, 전달 확인 메커니즘 차이를 코드로 관찰한다.

# ───────────────────────────────────────────────────
# 실습 2: QoS 0 / 1 / 2 비교
# ───────────────────────────────────────────────────
TOPIC_QOS = f'{BASE_TOPIC}/qos/compare'

qos_log = []        # 수신 기록
pub_ack  = []       # 발행 확인 mid 기록 (QoS 1·2)
qos_ready = threading.Event()

def on_connect_qos(client, userdata, flags, rc):
    if rc == 0:
        # 모든 QoS 토픽을 동시에 구독
        for q in [0, 1, 2]:
            client.subscribe(f'{TOPIC_QOS}/q{q}', qos=q)
        qos_ready.set()
        print('[구독자] 연결 성공 — QoS 0/1/2 토픽 모두 구독')

def on_message_qos(client, userdata, msg):
    data = json.loads(msg.payload.decode())
    qos_log.append({'topic': msg.topic, 'qos': msg.qos, 'data': data})
    print(f'  [수신] QoS={msg.qos} | {msg.topic.split("/")[-1]} | {data["value"]}')

def on_publish_qos(client, userdata, mid):
    pub_ack.append(mid)
    # QoS 0: 이 콜백 호출 안 됨 (확인 없음)
    # QoS 1: PUBACK 수신 후 호출
    # QoS 2: PUBCOMP 수신 후 호출
    print(f'  [발행확인] mid={mid} — 브로커 전달 완료')

sub2 = mqtt.Client(client_id='sub_qos_jbt', protocol=mqtt.MQTTv311)
sub2.on_connect = on_connect_qos
sub2.on_message = on_message_qos
sub2.connect(BROKER, PORT, 60)
sub2.loop_start()

if not qos_ready.wait(timeout=6):
    print('⚠️  연결 타임아웃')
else:
    time.sleep(0.3)
    pub2 = mqtt.Client(client_id='pub_qos_jbt', protocol=mqtt.MQTTv311)
    pub2.on_publish = on_publish_qos
    pub2.connect(BROKER, PORT, 60)
    pub2.loop_start()
    time.sleep(0.3)

    print('\n=== QoS 레벨별 발행 테스트 ===')
    for qos_level in [0, 1, 2]:
        payload = {
            'qos_level': qos_level,
            'value': round(random.uniform(380, 420), 1),
            'ts': datetime.now().isoformat()
        }
        topic = f'{TOPIC_QOS}/q{qos_level}'
        r = pub2.publish(topic, json.dumps(payload), qos=qos_level)
        label = ['Fire & Forget', 'At Least Once', 'Exactly Once'][qos_level]
        print(f'\n[발행] QoS {qos_level} ({label})')
        print(f'  topic  : {topic}')
        print(f'  mid    : {r.mid}  (QoS 0이면 0 반환)')
        time.sleep(1.2)

    time.sleep(1)
    pub2.loop_stop(); pub2.disconnect()

sub2.loop_stop(); sub2.disconnect()

print(f'\n=== 결과 요약 ===')
print(f'총 발행: 3건 | 총 수신: {len(qos_log)}건 | 발행 확인(ACK): {len(pub_ack)}건 (QoS 1·2만 해당)')

실행 결과:

print('\n※ QoS 0은 on_publish 콜백 미호출 — 확인 없이 전송 후 종료')
C:\Users\gogot\AppData\Local\Temp\ipykernel_19760\186025298.py:30: DeprecationWarning: Callback API version 1 is deprecated, update to latest version
  sub2 = mqtt.Client(client_id='sub_qos_jbt', protocol=mqtt.MQTTv311)
[구독자] 연결 성공 — QoS 0/1/2 토픽 모두 구독
C:\Users\gogot\AppData\Local\Temp\ipykernel_19760\186025298.py:40: DeprecationWarning: Callback API version 1 is deprecated, update to latest version
  pub2 = mqtt.Client(client_id='pub_qos_jbt', protocol=mqtt.MQTTv311)
=== QoS 레벨별 발행 테스트 ===

[발행] QoS 0 (Fire & Forget)
  topic  : swe2026_jbt/qos/compare/q0
  mid    : 1  (QoS 0이면 0 반환)
  [발행확인] mid=1 — 브로커 전달 완료
  [수신] QoS=0 | q0 | 401.2

[발행] QoS 1 (At Least Once)
  topic  : swe2026_jbt/qos/compare/q1
  mid    : 2  (QoS 0이면 0 반환)
  [수신] QoS=1 | q1 | 406.8
  [발행확인] mid=2 — 브로커 전달 완료

[발행] QoS 2 (Exactly Once)
  topic  : swe2026_jbt/qos/compare/q2
  mid    : 3  (QoS 0이면 0 반환)
  [발행확인] mid=3 — 브로커 전달 완료
  [수신] QoS=2 | q2 | 396.1

=== 결과 요약 ===
총 발행: 3건 | 총 수신: 3건 | 발행 확인(ACK): 3건 (QoS 1·2만 해당)

※ QoS 0은 on_publish 콜백 미호출 — 확인 없이 전송 후 종료

7. 실습 3 — 다중 설비 동시 시뮬레이션

반도체 공장의 CVD, PVD, 코팅 라인 3개 설비가 동시에 센서 데이터를 발행하는 시나리오.
하나의 구독자 서버가 와일드카드 토픽으로 전체를 수신한다.

# ───────────────────────────────────────────────────
# 실습 3: 다중 설비 동시 시뮬레이션
# 토픽 구조: swe2026_jbt/fab1/{설비ID}/sensor
#            swe2026_jbt/fab1/{설비ID}/alarm
# ───────────────────────────────────────────────────

# 3개 설비 정의
EQUIPMENTS = [
    {'id': 'CVD001', 'name': 'CVD 챔버',    'temp_range': (385, 415), 'pressure_range': (1.00, 1.10), 'alarm_threshold': 420},
    {'id': 'PVD002', 'name': 'PVD 스퍼터',  'temp_range': (250, 280), 'pressure_range': (0.50, 0.60), 'alarm_threshold': 285},
    {'id': 'CTG003', 'name': '코팅 라인',   'temp_range': (120, 140), 'pressure_range': (1.20, 1.30), 'alarm_threshold': 145},
]

# 구독자가 수집한 메시지를 설비별로 분류 저장
collected = {eq['id']: [] for eq in EQUIPMENTS}
alarms_received = []
multi_ready = threading.Event()

TOPIC_MULTI_SUB = f'{BASE_TOPIC}/fab1/#'   # 와일드카드: fab1 하위 전체 구독

def on_connect_multi(client, userdata, flags, rc):
    if rc == 0:
        client.subscribe(TOPIC_MULTI_SUB, qos=1)
        multi_ready.set()
        print(f'[수집 서버] 연결 성공 → 와일드카드 구독: {TOPIC_MULTI_SUB}')

def on_message_multi(client, userdata, msg):
    parts  = msg.topic.split('/')
    eq_id  = parts[2] if len(parts) > 2 else 'UNKNOWN'
    dtype  = parts[3] if len(parts) > 3 else 'unknown'
    data   = json.loads(msg.payload.decode())
    data['_topic'] = msg.topic

    if dtype == 'alarm':
        alarms_received.append(data)
        print(f'  🚨 [알람] {eq_id} | {data["alarm_code"]} | {data["value"]}°C (임계: {data["threshold"]}°C)')
    else:
        if eq_id in collected:
            collected[eq_id].append(data)
        print(f'  📡 [센서] {eq_id} | 온도={data.get("temperature","-")}°C | 압력={data.get("pressure","-")} Torr')

sub3 = mqtt.Client(client_id='sub_multi_jbt', protocol=mqtt.MQTTv311)
sub3.on_connect = on_connect_multi
sub3.on_message = on_message_multi
sub3.connect(BROKER, PORT, 60)
sub3.loop_start()

if not multi_ready.wait(timeout=6):
    print('⚠️  연결 타임아웃')
else:
    time.sleep(0.3)

    # 발행자 (설비 시뮬레이터)
    pub3 = mqtt.Client(client_id='pub_multi_jbt', protocol=mqtt.MQTTv311)
    pub3.connect(BROKER, PORT, 60)
    pub3.loop_start()
    time.sleep(0.3)

    print('\n=== 다중 설비 센서 발행 시작 ===')
    ROUNDS = 3   # 각 설비당 3회 발행

    for rnd in range(ROUNDS):
        print(f'\n─── Round {rnd+1} ───')
        for eq in EQUIPMENTS:
            temp     = round(random.uniform(*eq['temp_range']), 1)
            pressure = round(random.uniform(*eq['pressure_range']), 3)

            sensor_payload = {
                'equipment_id': eq['id'],
                'temperature' : temp,
                'pressure'    : pressure,
                'timestamp'   : datetime.now().isoformat()
            }
            t_sensor = f'{BASE_TOPIC}/fab1/{eq["id"]}/sensor'
            pub3.publish(t_sensor, json.dumps(sensor_payload), qos=1)
            time.sleep(0.15)

            # 마지막 라운드에서 CVD001 온도를 임계값 초과로 시뮬레이션
            if rnd == ROUNDS - 1 and eq['id'] == 'CVD001':
                alarm_payload = {
                    'equipment_id': eq['id'],
                    'alarm_code'  : 'TEMP_HIGH_03',
                    'value'       : eq['alarm_threshold'] + round(random.uniform(2, 8), 1),
                    'threshold'   : eq['alarm_threshold'],
                    'severity'    : 'HIGH',
                    'timestamp'   : datetime.now().isoformat()
                }
                t_alarm = f'{BASE_TOPIC}/fab1/{eq["id"]}/alarm'
                pub3.publish(t_alarm, json.dumps(alarm_payload), qos=2)
                time.sleep(0.15)

        time.sleep(0.5)

    time.sleep(1.5)
    pub3.loop_stop(); pub3.disconnect()

sub3.loop_stop(); sub3.disconnect()

print('\n=== 수집 결과 요약 ===')
for eq in EQUIPMENTS:
    n = len(collected[eq['id']])
    print(f'  {eq["id"]} ({eq["name"]}): 센서 {n}건')
print(f'  알람 수신: {len(alarms_received)}건')

실행 결과:

C:\Users\gogot\AppData\Local\Temp\ipykernel_19760\1830255542.py:42: DeprecationWarning: Callback API version 1 is deprecated, update to latest version
  sub3 = mqtt.Client(client_id='sub_multi_jbt', protocol=mqtt.MQTTv311)
[수집 서버] 연결 성공 → 와일드카드 구독: swe2026_jbt/fab1/#
C:\Users\gogot\AppData\Local\Temp\ipykernel_19760\1830255542.py:54: DeprecationWarning: Callback API version 1 is deprecated, update to latest version
  pub3 = mqtt.Client(client_id='pub_multi_jbt', protocol=mqtt.MQTTv311)
=== 다중 설비 센서 발행 시작 ===

─── Round 1 ───
  📡 [센서] CVD001 | 온도=397.0°C | 압력=1.005 Torr
  📡 [센서] PVD002 | 온도=263.8°C | 압력=0.503 Torr
  📡 [센서] CTG003 | 온도=131.9°C | 압력=1.27 Torr

─── Round 2 ───
  📡 [센서] CVD001 | 온도=410.9°C | 압력=1.099 Torr
  📡 [센서] PVD002 | 온도=254.8°C | 압력=0.546 Torr
  📡 [센서] CTG003 | 온도=124.3°C | 압력=1.263 Torr

─── Round 3 ───
  📡 [센서] CVD001 | 온도=400.7°C | 압력=1.02 Torr
  📡 [센서] PVD002 | 온도=275.5°C | 압력=0.556 Torr
  📡 [센서] CTG003 | 온도=128.3°C | 압력=1.263 Torr
  🚨 [알람] CVD001 | TEMP_HIGH_03 | 426.9°C (임계: 420°C)

=== 수집 결과 요약 ===
  CVD001 (CVD 챔버): 센서 3건
  PVD002 (PVD 스퍼터): 센서 3건
  CTG003 (코팅 라인): 센서 3건
  알람 수신: 1건

8. 실습 4 — Retained Message (최신 상태 유지)

Retained Message는 브로커가 토픽의 마지막 메시지를 저장해, 신규 구독자가 연결 즉시 가장 최신 상태를 받도록 하는 기능이다.

제조 현장 활용: 설비 현재 상태(Run/Idle/Down)를 retained으로 발행하면, 대시보드가 재시작 시 즉시 현재 상태를 표시할 수 있다.

# ───────────────────────────────────────────────────
# 실습 4: Retained Message
# ───────────────────────────────────────────────────
TOPIC_RETAIN = f'{BASE_TOPIC}/fab1/CVD001/status'

# ── Step 1: Retained 메시지 발행 ──
pub4 = mqtt.Client(client_id='pub_retain_jbt', protocol=mqtt.MQTTv311)
pub4.connect(BROKER, PORT, 60)
pub4.loop_start()
time.sleep(0.5)

status_payload = {
    'equipment_id': 'CVD001',
    'status'      : 'RUN',
    'recipe'      : 'LOGIC_90NM_V3',
    'operator'    : 'OP-01',
    'updated_at'  : datetime.now().isoformat()
}
# retain=True → 브로커가 이 메시지를 저장, 이후 구독자에게 즉시 전달
pub4.publish(TOPIC_RETAIN, json.dumps(status_payload, ensure_ascii=False), qos=1, retain=True)
print(f'[발행자] Retained 메시지 발행 완료')
print(f'  토픽  : {TOPIC_RETAIN}')
print(f'  내용  : {json.dumps(status_payload, ensure_ascii=False)}')
time.sleep(0.8)
pub4.loop_stop(); pub4.disconnect()

print('\n── 잠시 후 새 구독자 연결 (브로커 저장 메시지 수신 확인) ──')
time.sleep(1)

# ── Step 2: 나중에 연결한 구독자도 바로 최신 상태 수신 ──
retained_received = []
retain_ready = threading.Event()

def on_connect_retain(client, userdata, flags, rc):
    if rc == 0:
        client.subscribe(TOPIC_RETAIN, qos=1)
        retain_ready.set()

def on_message_retain(client, userdata, msg):
    data = json.loads(msg.payload.decode())
    data['_retain_flag'] = msg.retain   # 1이면 브로커에 저장된 retained 메시지
    retained_received.append(data)
    flag = '🗄️ [Retained]' if msg.retain else '📨 [실시간]'
    print(f'[신규 구독자] {flag} 수신')
    print(f'  설비: {data["equipment_id"]} | 상태: {data["status"]} | 레시피: {data["recipe"]}')
    print(f'  최종 업데이트: {data["updated_at"]}')

sub4 = mqtt.Client(client_id='sub_retain_jbt_new', protocol=mqtt.MQTTv311)
sub4.on_connect = on_connect_retain
sub4.on_message = on_message_retain
sub4.connect(BROKER, PORT, 60)
sub4.loop_start()
retain_ready.wait(timeout=5)
time.sleep(1.5)
sub4.loop_stop(); sub4.disconnect()

if retained_received and retained_received[0].get('_retain_flag'):
    print('\n✅ Retained Message 동작 확인 — 발행자 없이도 최신 상태 즉시 수신됨')
else:
    print('\n⚠️  Retained 메시지 미수신 (브로커 상태 확인 필요)')

실행 결과:

C:\Users\gogot\AppData\Local\Temp\ipykernel_19760\513532854.py:7: DeprecationWarning: Callback API version 1 is deprecated, update to latest version
  pub4 = mqtt.Client(client_id='pub_retain_jbt', protocol=mqtt.MQTTv311)
[발행자] Retained 메시지 발행 완료
  토픽  : swe2026_jbt/fab1/CVD001/status
  내용  : {"equipment_id": "CVD001", "status": "RUN", "recipe": "LOGIC_90NM_V3", "operator": "OP-01", "updated_at": "2026-04-12T21:59:08.469322"}

── 잠시 후 새 구독자 연결 (브로커 저장 메시지 수신 확인) ──
C:\Users\gogot\AppData\Local\Temp\ipykernel_19760\513532854.py:48: DeprecationWarning: Callback API version 1 is deprecated, update to latest version
  sub4 = mqtt.Client(client_id='sub_retain_jbt_new', protocol=mqtt.MQTTv311)
[신규 구독자] 🗄️ [Retained] 수신
  설비: CVD001 | 상태: RUN | 레시피: LOGIC_90NM_V3
  최종 업데이트: 2026-04-12T21:59:08.469322

✅ Retained Message 동작 확인 — 발행자 없이도 최신 상태 즉시 수신됨

9. 실습 5 — Will Message (설비 이상 종료 감지)

Will Message(유언 메시지)는 클라이언트가 비정상 종료(네트워크 단절, 프로세스 크래시 등) 시 브로커가 자동으로 지정된 토픽에 발행하는 메시지다.

제조 현장 활용: 설비 에이전트가 예기치 않게 종료되면 운영자에게 즉시 알람을 발송할 수 있다.

# ───────────────────────────────────────────────────
# 실습 5: Will Message 설정 및 동작 확인
# ───────────────────────────────────────────────────
TOPIC_WILL = f'{BASE_TOPIC}/fab1/CVD001/disconnect'

will_received = []
will_ready    = threading.Event()

# ── 구독자: 설비 종료 감지 모니터 ──
def on_connect_will_sub(client, userdata, flags, rc):
    if rc == 0:
        client.subscribe(TOPIC_WILL, qos=1)
        will_ready.set()
        print(f'[모니터] 연결 성공 → 설비 종료 토픽 구독: {TOPIC_WILL}')

def on_message_will(client, userdata, msg):
    data = json.loads(msg.payload.decode())
    will_received.append(data)
    print(f'\n🔴 [Will Message 수신!]')
    print(f'  설비  : {data.get("equipment_id")}')
    print(f'  원인  : {data.get("reason")}')
    print(f'  조치  : {data.get("action")}')

sub5 = mqtt.Client(client_id='sub_will_jbt', protocol=mqtt.MQTTv311)
sub5.on_connect = on_connect_will_sub
sub5.on_message = on_message_will
sub5.connect(BROKER, PORT, 60)
sub5.loop_start()
will_ready.wait(timeout=6)
time.sleep(0.3)

# ── 발행자: Will Message를 미리 등록하고 연결 ──
will_payload = json.dumps({
    'equipment_id': 'CVD001',
    'reason'      : '설비 에이전트 비정상 종료 (Will Message)',
    'action'      : '즉시 현장 점검 및 에이전트 재시작 필요',
    'ts'          : datetime.now().isoformat()
}, ensure_ascii=False)

pub5 = mqtt.Client(client_id='pub_will_jbt', protocol=mqtt.MQTTv311)
# ★ Will 설정: 연결 전 will_set() 호출 필수
pub5.will_set(
    topic   = TOPIC_WILL,
    payload = will_payload,
    qos     = 1,
    retain  = False
)
pub5.connect(BROKER, PORT, 60)
pub5.loop_start()
time.sleep(0.5)

print('[설비 에이전트] 정상 동작 중...')
# 정상 데이터 1건 발행
pub5.publish(f'{BASE_TOPIC}/fab1/CVD001/sensor',
             json.dumps({'temperature': 395.2, 'status': 'normal'}), qos=1)
time.sleep(0.8)

# ── 비정상 종료 시뮬레이션 ──
# loop_stop() 없이 disconnect() 생략 → 브로커가 keepalive 타임아웃 감지 후 Will 발행
# 노트북 환경에서는 소켓 강제 종료로 시뮬레이션
print('[설비 에이전트] ⚡ 비정상 종료 시뮬레이션 (소켓 강제 종료)...')
pub5.loop_stop()
pub5._sock = None   # 소켓 강제 해제 → 브로커가 비정상 종료로 판단 후 Will 발행

time.sleep(3)   # 브로커의 keepalive 타임아웃 처리 대기

sub5.loop_stop(); sub5.disconnect()

print('\n=== Will Message 실습 결과 ===')
if will_received:
    print('✅ Will Message 정상 수신 — 설비 비정상 종료가 모니터에 전달됨')
else:
    print('ℹ️  Will Message는 브로커 keepalive 타임아웃(기본 60초) 후 발행됩니다.')
    print('   실제 환경에서는 설비 에이전트 프로세스 크래시 시 자동 발행.')
    print('   코드 설정 방법이 핵심: pub.will_set(topic, payload, qos, retain)')

실행 결과:

C:\Users\gogot\AppData\Local\Temp\ipykernel_19760\1087827253.py:24: DeprecationWarning: Callback API version 1 is deprecated, update to latest version
  sub5 = mqtt.Client(client_id='sub_will_jbt', protocol=mqtt.MQTTv311)
[모니터] 연결 성공 → 설비 종료 토픽 구독: swe2026_jbt/fab1/CVD001/disconnect
C:\Users\gogot\AppData\Local\Temp\ipykernel_19760\1087827253.py:40: DeprecationWarning: Callback API version 1 is deprecated, update to latest version
  pub5 = mqtt.Client(client_id='pub_will_jbt', protocol=mqtt.MQTTv311)
[설비 에이전트] 정상 동작 중...
[설비 에이전트] ⚡ 비정상 종료 시뮬레이션 (소켓 강제 종료)...

🔴 [Will Message 수신!]
  설비  : CVD001
  원인  : 설비 에이전트 비정상 종료 (Will Message)
  조치  : 즉시 현장 점검 및 에이전트 재시작 필요

=== Will Message 실습 결과 ===
✅ Will Message 정상 수신 — 설비 비정상 종료가 모니터에 전달됨

10. 실습 6 — 수신 데이터 시계열 시각화

실습 3에서 수집한 3개 설비의 센서 데이터를 matplotlib으로 시각화한다.

# ───────────────────────────────────────────────────
# 실습 6: 수신 데이터 시계열 시각화
# ───────────────────────────────────────────────────
import matplotlib.pyplot as plt
import matplotlib.patches as mpatches
import numpy as np
import matplotlib as mpl

# Windows용 한글 폰트 설정
plt.rcParams['font.family'] = 'Malgun Gothic'
plt.rcParams['axes.unicode_minus'] = False   # 마이너스 기호 깨짐 방지



# 실습 3 수집 데이터 사용 (데이터 없을 경우 샘플 생성)
def get_plot_data(eq_id, n_points=10):
    if collected.get(eq_id) and len(collected[eq_id]) >= 2:
        rows = collected[eq_id]
        temps = [r['temperature'] for r in rows]
        return list(range(1, len(temps)+1)), temps
    # 샘플 데이터 생성
    ranges = {'CVD001': (385, 428), 'PVD002': (250, 290), 'CTG003': (120, 150)}
    lo, hi = ranges.get(eq_id, (100, 200))
    temps = [round(random.uniform(lo, hi-5), 1) for _ in range(n_points)]
    temps[-1] = hi + round(random.uniform(1, 5), 1)  # 마지막에 임계 초과
    return list(range(1, n_points+1)), temps

fig, axes = plt.subplots(3, 1, figsize=(12, 9), sharex=False)
fig.suptitle('다중 설비 온도 모니터링 — MQTT 수집 데이터', fontsize=14, fontweight='bold', y=1.01)

configs = [
    {'id': 'CVD001', 'name': 'CVD 챔버',   'threshold': 420, 'color': '#2196F3', 'normal': (385, 415)},
    {'id': 'PVD002', 'name': 'PVD 스퍼터', 'threshold': 285, 'color': '#4CAF50', 'normal': (250, 280)},
    {'id': 'CTG003', 'name': '코팅 라인',  'threshold': 145, 'color': '#FF9800', 'normal': (120, 140)},
]

for ax, cfg in zip(axes, configs):
    xs, temps = get_plot_data(cfg['id'])
    threshold = cfg['threshold']

    # 정상/이상 구분 색상
    colors = ['#F44336' if t > threshold else cfg['color'] for t in temps]

    ax.bar(xs, temps, color=colors, alpha=0.8, width=0.6, zorder=3)
    ax.plot(xs, temps, 'o-', color=cfg['color'], linewidth=1.5,
            markersize=5, alpha=0.9, zorder=4)

    # 임계값 라인
    ax.axhline(y=threshold, color='#F44336', linestyle='--', linewidth=1.5,
               label=f'임계값 {threshold}°C', zorder=5)

    # 정상 범위 음영
    ax.axhspan(cfg['normal'][0], cfg['normal'][1],
               alpha=0.08, color=cfg['color'], zorder=1)

    # 이상값 마킹
    for x, t in zip(xs, temps):
        if t > threshold:
            ax.annotate(f'⚠ {t}°C', xy=(x, t), xytext=(x+0.1, t+2),
                        fontsize=8, color='#F44336', fontweight='bold')

    ax.set_title(f'{cfg["id"]}  ({cfg["name"]})', fontsize=11, pad=4)
    ax.set_ylabel('온도 (°C)', fontsize=9)
    ax.set_xlabel('수신 순서', fontsize=9)
    ax.legend(fontsize=8, loc='upper left')
    ax.grid(axis='y', linestyle=':', alpha=0.5)
    ax.set_xticks(xs)

    # y축 범위 여유
    mn, mx = min(temps), max(temps)
    ax.set_ylim(mn - (mx-mn)*0.15, mx + (mx-mn)*0.2)

# 범례 (전체)
normal_patch = mpatches.Patch(color='#2196F3', alpha=0.6, label='정상 범위')
alarm_patch  = mpatches.Patch(color='#F44336', alpha=0.8, label='이상 (임계 초과)')
fig.legend(handles=[normal_patch, alarm_patch],
           loc='upper right', fontsize=9, bbox_to_anchor=(1.0, 1.0))

plt.tight_layout()
plt.savefig('/tmp/mqtt_multiequip_chart.png', dpi=120, bbox_inches='tight')
plt.show()
print('\n✅ 시각화 완료 — 빨간 막대: 임계값 초과 (이상 알람 발생 조건)')
No description has been provided for this image
✅ 시각화 완료 — 빨간 막대: 임계값 초과 (이상 알람 발생 조건)
# ───────────────────────────────────────────────────
# 토픽별 메시지 수신 통계 시각화
# ───────────────────────────────────────────────────
fig2, (ax1, ax2) = plt.subplots(1, 2, figsize=(11, 4))
fig2.suptitle('MQTT 수신 통계', fontsize=13, fontweight='bold')

# 설비별 수신 건수
eq_ids   = [e['id'] for e in EQUIPMENTS]
eq_names = [e['name'] for e in EQUIPMENTS]
eq_counts = [len(collected.get(i, [])) for i in eq_ids]
if sum(eq_counts) == 0:
    eq_counts = [9, 9, 9]  # 샘플

bars = ax1.bar(eq_names, eq_counts,
               color=['#2196F3','#4CAF50','#FF9800'], alpha=0.8, width=0.5)
ax1.set_title('설비별 센서 수신 건수', fontsize=11)
ax1.set_ylabel('메시지 수', fontsize=9)
for bar, cnt in zip(bars, eq_counts):
    ax1.text(bar.get_x() + bar.get_width()/2, bar.get_height() + 0.1,
             str(cnt), ha='center', va='bottom', fontsize=11, fontweight='bold')
ax1.grid(axis='y', linestyle=':', alpha=0.5)

# QoS 레벨 비교 파이
qos_labels = ['QoS 0\n(Fire & Forget)', 'QoS 1\n(At Least Once)', 'QoS 2\n(Exactly Once)']
qos_use    = [1, 3, 1]   # 본 실습의 QoS 사용 비율
colors_qos = ['#78909C', '#42A5F5', '#AB47BC']
wedges, texts, autotexts = ax2.pie(
    qos_use, labels=qos_labels, colors=colors_qos,
    autopct='%1.0f%%', startangle=90,
    textprops={'fontsize': 9}, pctdistance=0.75
)
ax2.set_title('QoS 레벨 사용 분포 (실습 기준)', fontsize=11)

plt.tight_layout()
plt.show()
print('✅ 통계 시각화 완료')

실행 결과:

No description has been provided for this image
✅ 통계 시각화 완료

11. 장점·한계·제조 현장 적용

장점

항목내용
초경량 프로토콜고정 헤더 최소 2바이트 — 저전력 IoT 센서, 위성·협대역 통신 환경에서도 동작
QoS 3단계데이터 중요도에 따라 전달 신뢰성 수준을 선택적으로 적용 가능
Pub/Sub 다대다하나의 발행으로 수십 개의 구독자에게 동시 전달 — 폴링 불필요
Retained Message구독자 재시작 시 별도 요청 없이 최신 상태를 즉시 수신
Will Message설비 비정상 종료 자동 감지 — 운영자 즉각 알림 가능
클라우드 호환AWS IoT Core, Azure IoT Hub, Google Cloud IoT 모두 MQTT 기반

한계

항목내용
브로커 단일 장애점브로커 다운 시 전체 데이터 흐름 중단 → 고가용성 클러스터링 필수
보안 기본 설정 취약1883 포트는 평문 전송 → 운영 환경은 TLS(8883) + 인증 필수
토픽 설계 복잡도설비·라인 증가 시 토픽 구조 관리 어려움 → 초기 설계 표준화 필수
메시지 순서 미보장QoS 1 이상이어도 다중 구독자 간 순서는 보장 안 됨

제조 현장 도입 시 핵심 체크리스트

□ 브로커 이중화 (Active-Standby) 구성
□ TLS 인증서 적용 (포트 8883)
□ 클라이언트 ID 충돌 방지 (설비ID + MAC 조합)
□ 토픽 네이밍 규칙 문서화 및 팀 공유
□ QoS 레벨 정책 수립 (주기 센서 QoS 0, 알람 QoS 2)
□ Retained Message 활용 설비 상태 관리
□ Will Message 모든 설비 에이전트에 적용
□ 메시지 스키마(JSON) 버전 관리

현장 경험 기반

MQTT는 제조 데이터 수집 계층의 핵심 인프라이며 FDC, SCADA, MES 시스템이 아무리 정교해도 현장 센서 데이터를 신뢰성 있게, 실시간으로 수집하는 파이프라인이 없으면 모두 무용지물이 된다고 생각함.

특히 AI 이상 감지 시스템에서 MQTT는 단순한 데이터 전달 도구를 넘어, 이벤트 드리븐(Event-driven) AI 아키텍처의 트리거 역할을 한다. 센서 임계값 초과 → MQTT 이벤트 → LLM AI 분석 → 운영자 조치 가이드의 파이프라인이 현실적으로 구현 가능한 이유가 바로 MQTT의 낮은 지연시간과 Pub/Sub 설계 구조 덕분이라고 할 수 있다.

MQTTAPI실습 | 전병탁
참고: Eclipse Paho MQTT v1.6.x | Python 3.10+ | 공개 브로커: test.mosquitto.org

0개의 댓글