> For the complete documentation index, see [llms.txt](https://kopens.gitbook.io/plantpulse-platform/llms.txt). Markdown versions of documentation pages are available by appending `.md` to page URLs; this page is available as [Markdown](https://kopens.gitbook.io/plantpulse-platform/developer/messaging.md).

# 메시징 아키텍처

## 개요 <a href="#overview" id="overview"></a>

이 문서에서는 플랜트펄스 플랫폼의 멀티 프로토콜 메시징 엔진 구조를 안내합니다. 메시징 엔진은 Kafka, MQTT, STOMP, HTTP 등 다양한 프로토콜을 통해 수신된 데이터를 통합하여 파이프라인으로 전달하는 역할을 담당합니다.

## 아키텍처 다이어그램 <a href="#architecture-diagram" id="architecture-diagram"></a>

```mermaid
flowchart LR
  subgraph IN["Ingestion Layer"]
    KC[Kafka Consumer<br/>4종 Runner]
    ML[MQTT Listener<br/>HiveMQ v5]
    SL[STOMP Listener<br/>JMS]
    HL[HTTP Listener<br/>AsyncServlet]
  end

  M2P[MessageToDataFlowPipe<br/>역직렬화 + 변환]
  ENG[Pipeline Engine<br/>36 스레드]

  subgraph OUT["Publishing Layer"]
    DDS[DDS Stage]
    KP[Kafka Publisher<br/>LZ4 압축]
    MP[MQTT Publisher<br/>async]
    SP[STOMP Publisher<br/>CAS + Lock]
  end

  KC --> M2P
  ML --> M2P
  SL --> M2P
  HL --> M2P
  M2P --> ENG
  ENG --> DDS
  DDS --> KP
  DDS --> MP
  DDS --> SP
```

## 멀티 프로토콜 브로커 <a href="#multi-protocol-broker" id="multi-protocol-broker"></a>

`MultiProtocolMessageServerBroker`는 플랜트펄스의 메시징 허브로서, 모든 프로토콜의 리스너와 퍼블리셔를 통합 관리합니다.

### 포트 목록

| 프로토콜  | 포트           | 설명                         |
| ----- | ------------ | -------------------------- |
| Kafka | `9092`       | Apache Kafka 브로커 포트입니다     |
| MQTT  | `1883`       | MQTT v5 브로커 포트입니다          |
| STOMP | `61000`      | STOMP over WebSocket 포트입니다 |
| HTTP  | `80` (서버 포트) | 비동기 서블릿 엔드포인트입니다           |

### 브로커 초기화 순서

```mermaid
flowchart TD
  ST[MultiProtocolMessageServerBroker.start]
  ST --> K[1 Kafka Consumer<br/>4종 Runner]
  ST --> M[2 MQTT Client<br/>4 clients]
  ST --> S[3 STOMP Listener<br/>JMS]
  ST --> H[4 HTTP Listener<br/>AsyncServlet]
```

## Kafka Consumer 아키텍처 <a href="#kafka-consumer" id="kafka-consumer"></a>

Kafka Consumer는 데이터 유형에 따라 4종의 Runner로 구분되어 운영됩니다.

### 4종 Runner

| Runner         | 토픽                  | 설명                      |
| -------------- | ------------------- | ----------------------- |
| `PointRunner`  | `plantpulse.point`  | 실시간 데이터 포인트를 소비합니다      |
| `RowRunner`    | `plantpulse.row`    | 로우(행) 단위 데이터를 소비합니다     |
| `ImportRunner` | `plantpulse.import` | 대량 Import 데이터를 소비합니다    |
| `BLOBRunner`   | `plantpulse.blob`   | 바이너리(이미지/영상) 데이터를 소비합니다 |

### 주요 설정

| 프로퍼티                               | 기본값      | 설명                            |
| ---------------------------------- | -------- | ----------------------------- |
| `kafka.consumer.threads`           | `4`      | Runner당 소비자 스레드 수입니다          |
| `kafka.consumer.poll.timeout.ms`   | `100`    | poll() 타임아웃 (밀리초)입니다          |
| `kafka.consumer.max.poll.records`  | `500`    | 한 번의 poll()로 가져오는 최대 레코드 수입니다 |
| `kafka.consumer.fetch.min.bytes`   | `1`      | 최소 fetch 바이트입니다               |
| `kafka.consumer.auto.offset.reset` | `latest` | 오프셋 리셋 정책입니다                  |

### 백프레셔 (Backpressure)

파이프라인 큐 사용률이 높아지면 Kafka Consumer의 poll 간격을 자동으로 조절하여 시스템 과부하를 방지합니다.

```
큐 사용률 < 70%  ──▶ 정상 poll 간격
큐 사용률 70~90% ──▶ poll 간격 증가 (점진적 감속)
큐 사용률 > 90%  ──▶ poll 일시 중지 (오프로드 완료 대기)
```

## MQTT Listener <a href="#mqtt-listener" id="mqtt-listener"></a>

HiveMQ v5 클라이언트 라이브러리 기반의 MQTT 리스너입니다.

### 구성

| 항목      | 값                   | 설명                                     |
| ------- | ------------------- | -------------------------------------- |
| 프로토콜 버전 | MQTT v5             | QoS 1 기반 메시지 수신을 지원합니다                 |
| 클라이언트 수 | 4개                  | 병렬 처리를 위한 다중 클라이언트입니다                  |
| 구독 방식   | Shared Subscription | `$share/plantpulse/topic` 형식의 공유 구독입니다 |
| 자동 재연결  | 지원                  | 연결 끊김 시 자동으로 재연결합니다                    |

### Shared Subscription

공유 구독(Shared Subscription)을 통해 4개의 MQTT 클라이언트가 동일 토픽의 메시지를 분산 소비합니다. 이를 통해 단일 클라이언트 병목 없이 높은 처리량을 확보할 수 있습니다.

```
MQTT Broker
    │
    ├──▶ Client-1 (shared) ──▶ MessageToDataFlowPipe
    ├──▶ Client-2 (shared) ──▶ MessageToDataFlowPipe
    ├──▶ Client-3 (shared) ──▶ MessageToDataFlowPipe
    └──▶ Client-4 (shared) ──▶ MessageToDataFlowPipe
```

## STOMP Listener <a href="#stomp-listener" id="stomp-listener"></a>

JMS(Java Message Service) 기반의 STOMP 리스너입니다. 레거시 시스템 및 Java 기반 클라이언트와의 연동에 활용됩니다.

| 항목     | 설명                     |
| ------ | ---------------------- |
| 프로토콜   | STOMP over WebSocket   |
| 포트     | `61000`                |
| 메시지 형식 | JSON (JMS TextMessage) |
| 구독 모델  | Queue (Point-to-Point) |

## HTTP Listener <a href="#http-listener" id="http-listener"></a>

비동기 서블릿(Async Servlet) 기반의 HTTP 리스너입니다. REST API를 통해 데이터를 직접 전송하는 외부 시스템과 연동할 때 사용됩니다.

| 항목     | 설명                               |
| ------ | -------------------------------- |
| 방식     | 비동기 서블릿 (`AsyncContext`)         |
| 엔드포인트  | `/api/v1/points`, `/api/v1/rows` |
| 데이터 형식 | JSON (배열 지원)                     |
| 타임아웃   | 30초                              |

비동기 서블릿을 사용하여 요청 처리 스레드를 즉시 반환하므로, 동시에 많은 HTTP 요청이 유입되어도 서블릿 스레드 풀이 고갈되지 않습니다.

## 메시지 역직렬화 <a href="#deserialization" id="deserialization"></a>

모든 프로토콜로부터 수신된 메시지는 `MessageToDataFlowPipe`에서 통합 역직렬화됩니다.

### Jackson Streaming API

대용량 JSON 배열 메시지의 효율적인 파싱을 위해 Jackson Streaming API(`JsonParser`)를 사용합니다. DOM 방식 대비 메모리 사용량이 현저히 적으며, 수만 건의 포인트가 포함된 배열도 스트리밍 방식으로 하나씩 파싱하여 파이프라인에 투입합니다.

```
JSON byte[] ──▶ JsonParser (Streaming) ──▶ PointData 객체 ──▶ Pipeline.enqueue()
```

## Publisher <a href="#publisher" id="publisher"></a>

파이프라인 처리가 완료된 데이터를 외부로 발행하는 퍼블리셔입니다.

### Kafka Publisher

| 항목    | 설명                            |
| ----- | ----------------------------- |
| 압축    | LZ4 압축을 적용하여 네트워크 대역폭을 절약합니다  |
| 배치    | `linger.ms` 기반 마이크로 배칭을 수행합니다 |
| 직렬화   | JSON (byte\[])                |
| 에러 처리 | 콜백 기반 비동기 에러 핸들링을 수행합니다       |

### MQTT Publisher

| 항목    | 설명                                           |
| ----- | -------------------------------------------- |
| 방식    | 비동기(async) 발행입니다                             |
| QoS   | QoS 0 (At most once)을 기본으로 사용합니다             |
| 토픽 형식 | `plantpulse/dds/{site_id}/{opc_id}/{tag_id}` |

### STOMP Publisher

| 항목     | 설명                                               |
| ------ | ------------------------------------------------ |
| 동시성 제어 | CAS(Compare-And-Swap) + `ReentrantLock` 이중 잠금입니다 |
| 메시지 형식 | JSON (JMS TextMessage)                           |
| 대상     | Queue / Topic 모두 지원합니다                           |

CAS와 `ReentrantLock`의 이중 잠금을 통해 다중 스레드 환경에서도 안전하게 메시지를 발행합니다.

## 상태 추적 <a href="#status-tracking" id="status-tracking"></a>

### MessageListenerStatus

각 리스너의 상태를 카운터로 추적합니다.

| 카운터             | 설명                 |
| --------------- | ------------------ |
| `received`      | 수신된 메시지 수입니다       |
| `processed`     | 처리 완료된 메시지 수입니다    |
| `failed`        | 처리 실패한 메시지 수입니다    |
| `dropped`       | 드롭된 메시지 수입니다       |
| `backpressured` | 백프레셔로 지연된 메시지 수입니다 |

### MessageRateTimer

초당 메시지 처리량을 측정하는 타이머입니다. 1초 간격으로 수신/처리/실패 건수를 샘플링하여 실시간 처리량 지표를 제공합니다.

## 지연 시간 측정 <a href="#latency-measurement" id="latency-measurement"></a>

메시지 수신 시점부터 파이프라인 투입까지의 종단 간(end-to-end) 지연 시간을 측정합니다.

```
메시지 수신 (ts1) ──▶ 역직렬화 ──▶ 파이프라인 enqueue (ts2)

지연 시간 = ts2 - ts1
```

지연 시간 메트릭은 프로토콜별로 분리되어 수집되므로, 각 프로토콜의 처리 성능을 독립적으로 모니터링할 수 있습니다.

## 보안 <a href="#security" id="security"></a>

### JAAS (Java Authentication and Authorization Service)

Kafka 브로커와의 통신에 JAAS 기반 인증을 지원합니다.

| 항목    | 설명                       |
| ----- | ------------------------ |
| 인증 방식 | SASL/PLAIN 또는 SASL/SCRAM |
| 설정 파일 | `kafka_jaas.conf`        |
| 암호화   | SSL/TLS 지원 (선택적)         |

MQTT 브로커의 경우 사용자명/비밀번호 기반 인증을 지원하며, STOMP는 JMS 보안 컨텍스트를 활용합니다.

## 에러 처리 및 복원력 <a href="#error-handling" id="error-handling"></a>

메시징 엔진은 다양한 장애 시나리오에 대비한 복원력 메커니즘을 갖추고 있습니다.

| 장애 유형        | 처리 방식                         |
| ------------ | ----------------------------- |
| Kafka 브로커 다운 | Consumer 자동 재연결 및 리밸런싱을 수행합니다 |
| MQTT 브로커 다운  | 자동 재연결 (지수 백오프)을 수행합니다        |
| 역직렬화 실패      | 에러 로깅 후 해당 메시지를 스킵합니다         |
| 파이프라인 큐 가득 참 | 백프레셔 적용 (poll 간격 조절)을 수행합니다   |
| 메모리 부족       | 파이프라인 오프로드 트리거를 발생시킵니다        |
| 네트워크 단절      | 자동 재연결 + 오프셋 기반 재소비를 수행합니다    |

## 토픽 설정 <a href="#topic-configuration" id="topic-configuration"></a>

### Kafka DDS 토픽

| 토픽                      | 용도                |
| ----------------------- | ----------------- |
| `plantpulse.dds.point`  | 실시간 데이터 포인트 발행입니다 |
| `plantpulse.dds.alarm`  | 알람 이벤트 발행입니다      |
| `plantpulse.dds.event`  | 시스템 이벤트 발행입니다     |
| `plantpulse.dds.status` | 장비/연결 상태 변경 발행입니다 |
| `plantpulse.dds.report` | 리포트 데이터 발행입니다     |

### MQTT DDS 토픽

| 토픽 패턴                                              | 용도              |
| -------------------------------------------------- | --------------- |
| `plantpulse/dds/point/{site_id}/{opc_id}/{tag_id}` | 태그별 실시간 값 발행입니다 |
| `plantpulse/dds/alarm/{site_id}`                   | 사이트별 알람 발행입니다   |
| `plantpulse/dds/status/{opc_id}`                   | OPC 연결 상태 발행입니다 |
