이 문서는 MQTT 메시징 프로토콜을 활용한 사설 클라우드 환경 구축을 위한 서버 및 클라이언트 핵심 코드를 다룹니다. 끊김 없는 통신을 위해 자동 재연결 기능을 안정적으로 지원하며, 공인 IP 서버, 사내망 서버 또는 클라우드 서버 등 다양한 환경에 배포할 수 있습니다. MQTT 통신 및 데이터 저장 기능을 구현하는 방법을 소개합니다.
최근 MQTT를 이용한 사설 클라우드 구축에 관심을 갖고 있습니다. MQTT(Message Queuing Telemetry Transport)는 경량 메시지 큐잉 및 원격 측정을 위한 발행/구독 메시지 전송 프로토콜로, 특히 사물 인터넷(IoT) 장치 간의 효율적인 통신에 적합합니다. 본 글에서는 안정적인 자동 재연결 기능을 갖춘 MQTT 서버 및 클라이언트 구축을 위한 핵심 코드를 다루겠습니다.
서버 구축 (Mosquitto)
MQTT 메시지를 중개하는 역할을 수행할 MQTT 브로커가 필요합니다. 여기서는 다양한 플랫폼을 지원하는 오픈 소스 MQTT 브로커인 Eclipse Mosquitto를 사용합니다.
sudo apt update
sudo apt install mosquitto mosquitto-clients
설치가 완료되면 Mosquitto 서비스를 시작합니다.
sudo systemctl start mosquitto
Mosquitto는 기본적으로 1883번 포트를 사용합니다. 포트 변경이나 기타 설정은 `/etc/mosquitto/mosquitto.conf` 설정 파일을 수정하여 조정할 수 있습니다.
클라이언트 구현 (Python)
다음은 Python으로 작성된 MQTT 클라이언트 예제 코드입니다. 이 코드는 MQTT 브로커와의 연결이 끊어졌을 때 자동으로 재연결하는 기능을 포함하고 있습니다.
import paho.mqtt.client as mqtt
import time
import sys
BROKER_ADDRESS = "your_broker_address" # 실제 브로커 주소로 변경하세요.
PORT = 1883
RECONNECT_DELAY = 5 # 재연결 시도 간격 (초)
def on_connect(client, userdata, flags, rc):
"""MQTT 브로커 연결 성공/실패 콜백 함수"""
if rc == 0:
print("브로커에 성공적으로 연결되었습니다.")
# 연결 성공 시 구독할 토픽
client.subscribe("my/data/topic")
else:
print(f"연결 실패, 응답 코드: {rc}")
# 연결 실패 시 재연결 시도 (on_disconnect에서 처리)
def on_disconnect(client, userdata, rc):
"""MQTT 브로커 연결 끊김 콜백 함수"""
print(f"연결이 끊어졌습니다. 코드: {rc}. 재연결 시도 중...")
# 자동 재연결 로직
while not client.is_connected():
try:
client.reconnect()
print("재연결 시도...")
time.sleep(RECONNECT_DELAY)
except ConnectionRefusedError:
print("연결이 거부되었습니다. 잠시 후 다시 시도합니다.")
time.sleep(RECONNECT_DELAY * 2)
except Exception as e:
print(f"재연결 중 오류 발생: {e}. 잠시 후 다시 시도합니다.")
time.sleep(RECONNECT_DELAY)
def on_message(client, userdata, msg):
"""메시지 수신 콜백 함수"""
print(f"주제: {msg.topic}, 메시지: {msg.payload.decode()}")
# 여기에 메시지 처리 로직 추가 (예: 데이터베이스 저장)
# MQTT 클라이언트 객체 생성
client = mqtt.Client(client_id="", clean_session=True, userdata=None, protocol=mqtt.MQTTv311)
# 콜백 함수 설정
client.on_connect = on_connect
client.on_disconnect = on_disconnect
client.on_message = on_message
try:
# 브로커에 연결 시도
print(f"{BROKER_ADDRESS}:{PORT} 에 연결 중...")
client.connect(BROKER_ADDRESS, PORT, 60)
# 네트워크 루프 시작 (별도 스레드에서 실행)
client.loop_start()
# 메시지 발행 루프
topic_to_publish = "device/status"
message_counter = 0
while True:
message = f"Heartbeat {message_counter}"
result = client.publish(topic_to_publish, message)
if result.rc == mqtt.MQTT_ERR_SUCCESS:
print(f"'{topic_to_publish}' 토픽에 '{message}' 발행 성공")
else:
print(f"'{topic_to_publish}' 토픽에 메시지 발행 실패, 코드: {result.rc}")
message_counter += 1
time.sleep(15) # 15초마다 메시지 발행
except KeyboardInterrupt:
print("프로그램 종료 요청 감지.")
except ConnectionRefusedError:
print("연결이 거부되었습니다. 브로커 주소 및 포트를 확인하세요.")
sys.exit(1)
except Exception as e:
print(f"연결 또는 실행 중 오류 발생: {e}")
sys.exit(1)
finally:
# 프로그램 종료 시 루프 중지 및 연결 해제
if client.is_connected() or client.is_listening():
print("정리 작업 중...")
client.loop_stop()
client.disconnect()
print("연결이 해제되었습니다.")
위 Python 코드는 `on_connect`와 `on_disconnect` 콜백 함수를 정의하여 각각 연결 성공 및 실패 시의 동작을 처리합니다. `on_disconnect` 함수 내에서 연결이 끊어졌을 때 재연결을 시도하는 로직이 구현되어 있습니다. `client.connect()` 메서드는 MQTT 브로커에 연결하며, `client.loop_start()`는 백그라운드 스레드를 시작하여 네트워크 통신을 관리합니다. `client.publish()` 메서드를 사용하여 지정된 주제로 메시지를 발행합니다.
데이터 저장
MQTT 브로커 자체는 메시지 영속성을 제공하지 않습니다. 따라서 메시지를 데이터베이스나 다른 저장소 시스템에 저장하려면, 특정 주제를 구독하고 수신된 메시지를 처리하는 별도의 로직을 구현해야 합니다.
import paho.mqtt.client as mqtt
import sqlite3
import sys
BROKER_ADDRESS = "your_broker_address" # 실제 브로커 주소로 변경하세요.
PORT = 1883
DB_FILE = 'mqtt_data_log.db'
TABLE_NAME = 'received_messages'
def initialize_database():
"""SQLite 데이터베이스 및 테이블 초기화"""
try:
conn = sqlite3.connect(DB_FILE)
cursor = conn.cursor()
cursor.execute(f'''
CREATE TABLE IF NOT EXISTS {TABLE_NAME} (
id INTEGER PRIMARY KEY AUTOINCREMENT,
topic TEXT NOT NULL,
payload TEXT,
received_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP
)
''')
conn.commit()
print(f"데이터베이스 '{DB_FILE}' 및 테이블 '{TABLE_NAME}'이 준비되었습니다.")
except sqlite3.Error as e:
print(f"데이터베이스 초기화 오류: {e}")
sys.exit(1)
finally:
if conn:
conn.close()
def on_message_received(client, userdata, msg):
"""메시지 수신 시 데이터베이스에 저장하는 콜백 함수"""
try:
conn = sqlite3.connect(DB_FILE)
cursor = conn.cursor()
# SQL Injection 방지를 위해 매개변수화된 쿼리 사용
cursor.execute(f"INSERT INTO {TABLE_NAME} (topic, payload) VALUES (?, ?)",
(msg.topic, msg.payload.decode('utf-8', 'ignore'))) # 디코딩 오류 방지
conn.commit()
print(f"메시지 저장 완료 - 주제: {msg.topic}, 페이로드: {msg.payload.decode('utf-8', 'ignore')[:50]}...") # 일부만 출력
except sqlite3.Error as e:
print(f"데이터베이스 저장 오류: {e}")
except Exception as e:
print(f"메시지 처리 중 예외 발생: {e}")
finally:
if conn:
conn.close()
def on_connect_db(client, userdata, flags, rc):
"""DB 저장용 클라이언트 연결 콜백"""
if rc == 0:
print("데이터베이스 저장용 클라이언트가 브로커에 연결되었습니다.")
# 데이터베이스 저장용으로 구독할 토픽
topic = userdata['subscribe_topic']
client.subscribe(topic)
print(f"'{topic}' 토픽 구독 시작.")
else:
print(f"데이터베이스 저장용 클라이언트 연결 실패, 코드: {rc}")
# 데이터베이스 초기화
initialize_database()
# 데이터베이스 저장용 MQTT 클라이언트 객체 생성
# userdata를 사용하여 구독할 토픽을 전달
db_client = mqtt.Client(userdata={'subscribe_topic': "sensor/data"})
db_client.on_connect = on_connect_db
db_client.on_message = on_message_received
try:
# 브로커에 연결
print(f"{BROKER_ADDRESS}:{PORT} 에 데이터베이스 저장용 클라이언트로 연결 중...")
db_client.connect(BROKER_ADDRESS, PORT, 60)
# 메시지를 영구적으로 처리하기 위해 loop_forever() 사용
# 이 함수는 블로킹되며, 별도 스레드에서 실행해야 할 경우 loop_start() 사용
db_client.loop_forever()
except ConnectionRefusedError:
print("연결이 거부되었습니다. 브로커가 실행 중인지 확인하세요.")
sys.exit(1)
except KeyboardInterrupt:
print("프로그램 종료.")
except Exception as e:
print(f"오류 발생: {e}")
finally:
if db_client.is_connected():
db_client.disconnect()
print("데이터베이스 저장용 클라이언트 연결 해제.")
이 코드는 `on_message_received` 콜백 함수 내에서 수신된 메시지를 SQLite 데이터베이스에 저장하는 과정을 보여줍니다. 먼저 데이터베이스에 연결하고, SQL 문의를 실행하여 데이터를 삽입한 후, 트랜잭션을 커밋하고 연결을 닫습니다. `ignore` 옵션을 사용하여 잘못된 UTF-8 바이트 시퀀스로 인한 디코딩 오류를 방지합니다.
배포 고려사항
MQTT 시스템은 다양한 환경에 배포될 수 있습니다. 원격 접근이 필요한 경우 공인 IP를 가진 서버를 사용하고, 로컬 네트워크 내 통신이 목적이라면 사내망 서버를 활용할 수 있습니다. 비용 효율성과 성능을 고려할 때 클라우드 기반의 경량 서버 인스턴스도 좋은 선택이 될 수 있습니다.