> 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/data-flow.md).

# 데이터 처리 플로우

## 개요

이 문서에서는 플랜트펄스 플랫폼의 전체 데이터 처리 흐름을 안내합니다. 플랜트펄스 플랫폼은 산업 현장의 다양한 데이터 소스로부터 실시간 데이터를 수집하고, 변환/분석/저장/시각화하는 엔드투엔드 데이터 파이프라인을 제공합니다. 비동기 처리와 선 처리 → 후 저장 메커니즘을 통해 데이터 용량에 상관없이 Low-Latency를 확보하고 있습니다.

처음 플랫폼을 접하시는 분이라면, 데이터가 "수집 → 변환 → 분석 → 저장 → 공유 → 시각화"의 6단계를 거쳐 처리된다는 큰 흐름을 먼저 파악해 주시면 이해에 도움이 됩니다.

## 전체 데이터 플로우

```mermaid
flowchart LR
  subgraph SRC["1. 수집 (Collect)"]
    OPCUA[OPC-UA / DA]
    PLC[PLC]
    MODBUS[Modbus]
    DB[Database]
    HTTP[HTTP / REST]
    CSV[CSV / File]
    MQTT_IN[MQTT 디바이스]
  end

  subgraph TRN["2. 변환 (Transform)"]
    TAG_MAP[Tag 매핑]
    ASSET[Asset Tree]
    PROTO[Protocol Meta]
    BAND[Alert Band]
  end

  subgraph ANL["3. 분석 (Analysis)"]
    PIPE[36 스레드 파이프라인]
    CEP[CEP / Esper]
  end

  subgraph STO["4. 저장 (Store)"]
    BUF[Redis/Valkey 버퍼]
    CASS[(Cassandra 시계열)]
    PG[(PostgreSQL 메타)]
    MINIO[(MinIO 오브젝트)]
  end

  subgraph SHARE["5. 공유 (Trigger)"]
    STREAM[MQTT/STOMP 트리거]
    DBT[DB 트리거]
  end

  subgraph VIS["6. 시각화 (Visualize)"]
    CANVAS[캔버스 / 위젯]
    TREND[Trend 차트]
    ALARM_V[Alarm Console]
    GRAFANA[Grafana :3000]
  end

  OPCUA --> TAG_MAP
  PLC --> TAG_MAP
  MODBUS --> TAG_MAP
  DB --> TAG_MAP
  HTTP --> TAG_MAP
  CSV --> TAG_MAP
  MQTT_IN --> TAG_MAP

  TAG_MAP --> ASSET --> PROTO --> BAND --> PIPE
  PIPE --> CEP
  PIPE --> BUF --> CASS
  PIPE --> PG
  CEP --> STREAM
  CEP --> DBT
  CASS --> CANVAS
  CASS --> TREND
  CASS --> GRAFANA
  PG --> ALARM_V
```

## 1단계: 데이터 수집 (Collect)

데이터 수집은 산업 현장의 센서, PLC, MES 등 다양한 장비로부터 데이터를 가져오는 첫 번째 단계입니다.

### 지원 프로토콜

PlantPulse는 아래와 같이 다양한 산업 표준 프로토콜을 지원합니다. 현장 환경에 맞는 프로토콜을 선택하여 연결하시면 됩니다.

| 프로토콜           | 설명                            |
| -------------- | ----------------------------- |
| OPC-DA/UA      | OPC 표준 지원, OPC 서버와의 자동 연결/재연결 |
| PLC            | 국내외 다양한 제조사 PLC 직접 연결         |
| Modbus         | Modbus TCP/RTU 프로토콜 지원        |
| Database       | RDBMS 소스에서 데이터 수집             |
| File (CSV/LOG) | 파일 기반 데이터 소스 수집               |
| MQTT/REST/HTTP | IoT 표준 프로토콜 지원                |
| JMS            | Java Message Service 연동       |
| BLOB           | 이미지 및 동영상 바이너리 데이터            |

### 수집 특징

* **고속 이벤트 수집**: 고속/대용량 이벤트 소스를 적은 리소스로 수집합니다.
* **메시지 유실 방지**: 서버 다운 시 자동으로 데이터를 캐시하며, 복구 시 자동 재전송합니다.
* **자동 페일오버**: OPC 서버 다운 시 자동으로 Re-connect를 수행합니다.
* **커스텀 에이전트**: Java, .NET, C/C++ 라이브러리를 이용하여 커스텀 수집 에이전트를 개발하실 수 있습니다.

### 수집 데이터 흐름

```mermaid
flowchart LR
  DEV[Device · PLC · Sensor] --> AGENT[Edge Agent<br/>OPC-UA Plugin :11004]
  AGENT --> KAFKA[Kafka :9092]
  DEV2[IoT 디바이스] --> MQTT[MQTT :1883]
  MQTT --> KAFKA
  KAFKA --> SVR[plantpulse-server<br/>IIoT Engine]
  KAFKA --> CEP[plantpulse-cep :7400]
  SVR --> STOMP[STOMP :61000<br/>WebSocket 푸시]
  STOMP --> BROWSER[브라우저]
```

> **현재 표준**: 데이터 수집은 `plantpulse-plugin/opc-ua` (포트 11004) 또는 외부 Edge Agent → Kafka(:9092) / MQTT(:1883) 경로를 사용합니다. 과거 별도 모듈로 존재하던 `plantpulse-agent` 는 현재 `plantpulse-plugin` 으로 통합되었습니다.

## 2단계: 데이터 변환 (Transform)

수집된 원시 데이터를 플랫폼 내부 모델로 변환하는 단계입니다. 이 과정을 통해 원시 데이터가 플랫폼에서 활용 가능한 형태로 정리됩니다.

### 변환 프로세스

| 단계            | 설명                                |
| ------------- | --------------------------------- |
| TAG 매핑        | 물리적 태그를 플랫폼 내부 태그 ID로 매핑합니다       |
| ASSET TREE 매핑 | ISA-95 기반 논리적 자산 트리에 태그를 연결합니다    |
| SITE 매핑       | 지역/라인/설비 계층 구조에 데이터를 할당합니다        |
| PROTOCOL 식별   | 수집 프로토콜 유형 메타데이터를 추가합니다           |
| ALERT BAND 적용 | 태그별 알람 임계값(HH/H/L/LL) 설정을 적용합니다   |
| ATTRIBUTE 부여  | 태그 속성 정보(단위, 설명, 데이터 타입 등)를 연결합니다 |

### 시계열 태그 데이터 모델

변환이 완료된 데이터는 아래와 같은 태그 데이터 모델 형태로 저장됩니다.

```
┌──────────────────────────────────────────┐
│           TAG DATA MODEL                  │
├──────────┬───────────┬─────────┬─────────┤
│ RESOURCE │ DATE/TIME │  NAME   │  VALUE  │
│          │           │         │         │
│ ATTRIBUTE│           │         │ RESULT  │
└──────────┴───────────┴─────────┴─────────┘
```

## 3단계: 스트리밍 분석 (Analysis)

### 실시간 분석 (Streaming Analysis)

IIoT 엔진 내부의 36스레드 파이프라인에서 실시간으로 데이터를 처리합니다. 아래의 분석 유형을 조합하여 다양한 실시간 분석 시나리오를 구현하실 수 있습니다.

| 분석 유형                | 설명                                          |
| -------------------- | ------------------------------------------- |
| **Filtering**        | 조건 기반 데이터 필터링, 불필요한 중복 데이터를 제거합니다           |
| **Aggregation**      | 윈도우 기반 집계 (평균, 최대, 최소, 합계, 표준편차)를 수행합니다     |
| **Pattern Matching** | CEP 엔진(Esper) 기반으로 이벤트 패턴을 매칭합니다            |
| **Data Trigger**     | 조건 충족 시 외부 시스템으로 데이터를 전달합니다 (MQTT/STOMP/DB) |

### CEP (Complex Event Processing) 엔진

초당 500,000 이벤트를 처리할 수 있는 인메모리 기반 CEP 엔진입니다. 복잡한 이벤트 패턴을 실시간으로 감지하여 알람 발생, 트리거 실행 등의 동작을 수행합니다.

```
입력 이벤트                     CEP 엔진 (Esper)                    출력 이벤트
─────────────────────────────────────────────────────────────────────────────

TEMP EVENT ──┐                  EQL (Event Query Language)
SPEED EVENT ─┼──▶  ┌─────────────────────────────────┐  ──▶ 알람 발생
BUY EVENT  ──┘     │ SELECT AVG(TEMP)                 │  ──▶ 트리거 실행
                   │ FROM TempEvent.win:time(10 min)  │  ──▶ DB 저장
                   │ WHERE TEMP_ID = 'MAIN'           │  ──▶ 메시지 전송
                   └─────────────────────────────────┘
```

### EQL 분석 기능

EQL(Event Query Language)을 활용하면 아래와 같은 다양한 분석 기능을 사용하실 수 있습니다.

| 기능                       | 설명                        |
| ------------------------ | ------------------------- |
| Filtering                | 조건에 맞는 이벤트만 선별합니다         |
| Correlation (JOIN)       | 여러 이벤트 스트림 간 상관관계를 분석합니다  |
| Database Lookup          | 외부 DB 데이터를 참조합니다          |
| Hierarchical Events      | 계층적 이벤트 구조를 처리합니다         |
| Event Pattern Matching   | 시간/순서 기반으로 이벤트 패턴을 탐지합니다  |
| In-Memory Caching        | 메모리 기반 고속 캐싱을 활용합니다       |
| Aggregation over Windows | 시간/개수 윈도우 기반으로 집계합니다      |
| Dynamic Query            | 런타임에 쿼리를 추가하거나 변경할 수 있습니다 |

### 히스토리컬 분석 (Historical Analysis)

과거 데이터를 기반으로 심층 분석이 필요한 경우, 아래의 방법을 활용하실 수 있습니다.

| 분석 유형                | 설명                          |
| -------------------- | --------------------------- |
| Wide Column Store    | Cassandra 기반 대용량 시계열 데이터 조회 |
| Analysis Script      | Python/R 기반 분석 스크립트         |
| Statistic Package    | 통계 패키지 연동                   |
| Predictive Algorithm | 머신러닝 기반 예측 알고리즘             |

## 4단계: 데이터 저장 (Data Store)

### 시계열 데이터 저장

분석이 완료된 데이터는 Redis/Valkey 쓰기 버퍼를 거쳐 Cassandra 클러스터에 저장됩니다. Micro Batch 방식으로 주기적으로 플러시하여 저장 성능을 최적화하고 있습니다.

```mermaid
flowchart TB
  PIPE[IIoT Pipeline<br/>36 스레드] --> BUFFER[Write Buffer<br/>Valkey :6379]
  BUFFER -->|Micro Batch Flush| CASS[Cassandra :9042]
  CASS --> HOT[Hot Data<br/>최근 · 빠른 조회]
  CASS --> COLD[Cold Data<br/>과거 · TTL]
  COLD -.->|warehouse 아카이빙| ICE[Iceberg / MinIO]
```

### 저장소 구성

PlantPulse는 데이터의 특성에 맞게 여러 저장소를 조합하여 사용합니다.

| 저장소           | 용도                           | 기술                 |
| ------------- | ---------------------------- | ------------------ |
| **메타데이터 저장소** | 사용자, 사이트, 에셋, 태그 설정, 알람 규칙 등 | PostgreSQL         |
| **시계열 저장소**   | 태그 포인트, 알람 이력, 에셋 데이터 등      | Cassandra/ScyllaDB |
| **캐시**        | 쓰기 버퍼링, 실시간 데이터, 세션 등        | Redis              |
| **메시지 버스**    | 컴포넌트 간 데이터 전달                | Kafka, MQTT, STOMP |

### 시계열 엔진 특징

* **순차적 데이터 저장**: Write 최적화된 시계열 데이터 모델을 사용합니다.
* **구간 검색 최적화**: 날짜/시간 기반 구간 검색을 고속으로 처리합니다.
* **분산 처리**: 클러스터 기반 수평 확장(Scale-out)을 지원합니다.
* **Schema-Free**: 유연한 스키마 구조를 채택하고 있습니다.
* **Spark 연동**: Spark/Cassandra Connector를 통한 분산 SQL 분석이 가능합니다.

## 5단계: 데이터 공유 (Trigger)

수집된 데이터를 외부 시스템으로 실시간 전달하는 단계입니다. 트리거를 설정하면 특정 조건에 맞는 데이터를 자동으로 외부로 내보낼 수 있습니다.

### 트리거 유형

| 유형                    | 설명                                |
| --------------------- | --------------------------------- |
| **Streaming Trigger** | MQTT, STOMP 등 메시지 큐 형태로 실시간 출력합니다 |
| **DB Trigger**        | 정의한 컬럼 정보를 외부 DB 테이블로 저장합니다       |

### 데이터 공유 흐름

```mermaid
flowchart LR
  ENG[PlantPulse Engine] -->|Streaming Trigger| MQTT_OUT[MQTT :1883<br/>STOMP :61000]
  ENG -->|DB Trigger| PG[PostgreSQL :5432]
  ENG -->|Kafka 재발행| KAFKA_OUT[Kafka :9092]
  MQTT_OUT --> EXT[외부 시스템<br/>실시간 구독]
  PG --> ERP[ERP / MES 배치 연계]
  KAFKA_OUT --> DOWN[Downstream<br/>Flow Engine · CEP · BI]
```

### 데이터 샘플링

대용량 시계열 데이터를 실시간으로 샘플링 처리하여 통계 작업을 효율적으로 수행합니다.

| 샘플링 옵션      | 설명                       |
| ----------- | ------------------------ |
| Time Series | 1분 \~ 1시간 단위 집계          |
| Value       | 처음/마지막 값, 최저/평균/최고, 표준편차 |

## 6단계: 시각화 (Visualization)

최종적으로 처리된 데이터는 다양한 형태의 시각화를 통해 사용자에게 제공됩니다.

| 시각화 유형            | 설명                                         |
| ----------------- | ------------------------------------------ |
| **실시간 모니터링**      | 에셋 트리 기반 태그 현재 상태값 시계열 차트                  |
| **시계열 트렌드 분석**    | 다중 태그 트렌드 차트, Multi-Y축, 줌인/아웃              |
| **실시간 알람**        | 중요도별 통계, 비정상 알람 TAG 랭킹, Full-Text 검색       |
| **캔버스**           | SVG 기반 위젯 시각화, 실시간 데이터 바인딩, OEE/RAM/EMS 분석 |
| **Dynamic Query** | 현재 발생 데이터에 대한 실시간 다양한 분석                   |

## 데이터 처리 성능

아래 표는 PlantPulse 플랫폼의 데이터 처리 성능 지표입니다. 운영 환경에 따라 차이가 있을 수 있으니 참고용으로 활용해 주세요.

| 지표         | 성능                    |
| ---------- | --------------------- |
| 파이프라인 스레드  | 36개 병렬 처리             |
| 메시지 처리량    | 초당 40,000건            |
| 큐 용량       | 1,200,000건            |
| CEP 이벤트 처리 | 초당 500,000건 (독립 환경)   |
| 일일 데이터 처리  | 100 PLC 기준 약 8억 7천만 건 |
| 연간 데이터 관리  | 약 100TB 이상            |
