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

# 파이프라인 아키텍처

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

이 문서에서는 플랜트펄스 플랫폼의 핵심 데이터 파이프라인 엔진 구조를 안내합니다. 파이프라인은 외부 프로토콜로부터 수신된 데이터를 샤딩 기반 큐에 분배하고, 6단계 DataFlow를 거쳐 최종 저장소 및 DDS(Data Distribution Service)로 전달하는 고성능 처리 엔진입니다.

## 전체 데이터 흐름 다이어그램 <a href="#data-flow-diagram" id="data-flow-diagram"></a>

```mermaid
flowchart LR
  IN[Ingestion<br/>Kafka · MQTT · STOMP · HTTP]
  SR[ShardRouter<br/>hash·tag_id]
  PQ[PipelineQueue<br/>MPMC + Offload]

  subgraph DF["DataFlow (6 stages)"]
    P[1 Prepare] --> V[2 Validate] --> C[3 Cache] --> S[4 Stream] --> ST[5 Store] --> D[6 DDS]
  end

  CASS[(Cassandra)]
  DDS[Kafka · MQTT<br/>재발행]

  IN --> SR --> PQ --> DF
  DF --> CASS
  DF --> DDS
```

## 샤딩 아키텍처 <a href="#sharding" id="sharding"></a>

파이프라인은 데이터 처리량을 극대화하기 위해 샤딩 기반 병렬 처리를 채택하고 있습니다.

### ShardRouter

`ShardRouter`는 수신된 데이터 포인트의 `tag_id`를 해시하여 샤드 인덱스를 결정합니다. 동일한 태그는 항상 동일한 샤드로 라우팅되므로, 태그 단위의 순서 보장이 가능합니다.

```
tag_id ──▶ hash(tag_id) ──▶ shard_index = hash % shard_count
```

### ShardedWorkerGroup

각 샤드는 독립적인 워커 스레드를 가지며, `ShardedWorkerGroup`이 전체 샤드 워커의 라이프사이클을 관리합니다.

| 항목     | 설명                                             |
| ------ | ---------------------------------------------- |
| 샤드 수   | CPU 코어 수 기반 자동 결정 (기본값: `availableProcessors`) |
| 워커 타입  | 각 샤드별 1개의 소비자 스레드                              |
| 라우팅 방식 | `tag_id` 해시 기반 고정 할당                           |
| 순서 보장  | 동일 태그는 동일 샤드에서 처리되어 순서가 보장됩니다                  |

## PipelineQueue 구조 <a href="#pipeline-queue" id="pipeline-queue"></a>

파이프라인 큐는 3계층 구조로 구성되어, 고부하 상황에서도 데이터 유실 없이 처리할 수 있도록 설계되었습니다.

```mermaid
flowchart TD
  T1[Tier 1<br/>MpmcArrayQueue<br/>Lock-free · In-memory]
  T2[Tier 2<br/>O3 Buffer<br/>Out-of-Order 정렬]
  T3[Tier 3<br/>FileOffloadQueue<br/>RocksDB · 디스크]

  T1 -->|큐 사용률 90 percent 초과| T2
  T2 -->|메모리 임계치 초과| T3
  T3 -.->|여유 회복 시 리로드| T1
```

| 계층     | 구현체                | 특징                                         |
| ------ | ------------------ | ------------------------------------------ |
| Tier 1 | `MpmcArrayQueue`   | 다중 생산자-다중 소비자 Lock-free 큐. 최고 처리 성능을 제공합니다 |
| Tier 2 | O3 Buffer          | 타임스탬프 기반 Out-of-Order 데이터 정렬 버퍼입니다         |
| Tier 3 | `FileOffloadQueue` | RocksDB 기반 디스크 오프로드 큐. 메모리 부족 시 자동 전환됩니다   |

## 오프로드/리로드 메커니즘 <a href="#offload-reload" id="offload-reload"></a>

파이프라인은 메모리 사용량에 따라 자동으로 디스크 오프로드 및 리로드를 수행합니다. 이를 통해 대량의 데이터 유입 시에도 OOM(Out of Memory) 없이 안정적으로 운영할 수 있습니다.

| 임계값            | 동작        | 설명                                       |
| -------------- | --------- | ---------------------------------------- |
| **90%** (오프로드) | 메모리 → 디스크 | 큐 사용률이 90%를 초과하면 RocksDB로 오프로드합니다        |
| **65%** (리로드)  | 디스크 → 메모리 | 큐 사용률이 65% 이하로 내려가면 디스크에서 다시 메모리로 리로드합니다 |

```
              오프로드 (90%)                  리로드 (65%)
메모리 큐 ────────────────▶ RocksDB ────────────────▶ 메모리 큐
  (가득 참)                  (디스크)                   (여유 확보)
```

오프로드된 데이터는 RocksDB에 순서가 보장된 상태로 저장되며, 리로드 시에도 원래의 처리 순서를 유지합니다.

## DataFlow 6단계 <a href="#dataflow-stages" id="dataflow-stages"></a>

파이프라인에서 디큐된 데이터는 6단계의 DataFlow 스테이지를 순차적으로 통과합니다.

### Stage 1: Prepare (준비) <a href="#stage-prepare" id="stage-prepare"></a>

도메인 정보 보강과 타임스탬프 보정을 수행하는 단계입니다.

* **도메인 보강**: `tag_id`를 기반으로 캐시에서 태그 메타정보(OPC, Asset, Site 등)를 조회하여 데이터에 부착합니다
* **타임스탬프 보정**: 에이전트/서버 간 시간 차이를 보정합니다. 미래 타임스탬프나 비정상적인 과거 타임스탬프를 감지하여 현재 서버 시간으로 교정합니다
* **기본값 설정**: 누락된 필드에 기본값을 할당합니다

### Stage 2: Validate (검증) <a href="#stage-validate" id="stage-validate"></a>

8가지 검증 코드를 통해 데이터 품질을 확인합니다.

| 검증 코드            | 설명                |
| ---------------- | ----------------- |
| `VALID`          | 정상 데이터입니다         |
| `TAG_NOT_FOUND`  | 등록되지 않은 태그입니다     |
| `TAG_DISABLED`   | 비활성화된 태그입니다       |
| `OPC_DISABLED`   | 비활성화된 OPC 연결입니다   |
| `VALUE_NULL`     | 값이 null입니다        |
| `TYPE_MISMATCH`  | 데이터 타입이 일치하지 않습니다 |
| `RANGE_EXCEEDED` | 설정된 범위를 초과하였습니다   |
| `DUPLICATE`      | 중복 데이터입니다         |

검증에 실패한 데이터는 카운터에 기록된 후 드롭되며, 이후 스테이지로 전달되지 않습니다.

### Stage 3: Cache (캐시) <a href="#stage-cache" id="stage-cache"></a>

검증을 통과한 데이터의 최신 값을 인메모리 캐시에 갱신합니다. 이 캐시는 실시간 현재값 조회, 대시보드, 알람 판정 등에 활용됩니다.

### Stage 4: Stream (스트리밍) <a href="#stage-stream" id="stage-stream"></a>

실시간 스트림 프로세서로 데이터를 전달하는 단계입니다. 두 가지 모드를 지원합니다.

| 모드          | 설명                                                   |
| ----------- | ---------------------------------------------------- |
| `DIRECT`    | 동기 방식으로 직접 전달합니다. 저지연이 필요한 환경에 적합합니다                 |
| `DISRUPTOR` | LMAX Disruptor 기반 비동기 링 버퍼를 통해 전달합니다. 고처리량 환경에 적합합니다 |

### Stage 5: Store (저장) <a href="#stage-store" id="stage-store"></a>

Cassandra 와이드 컬럼 스토어에 시계열 데이터를 영구 저장합니다. 비동기 배치 쓰기를 통해 높은 쓰기 처리량을 확보하고 있습니다.

### Stage 6: DDS (분산 데이터 서비스) <a href="#stage-dds" id="stage-dds"></a>

처리 완료된 데이터를 외부 시스템으로 분배합니다.

| 프로토콜  | 용도                                |
| ----- | --------------------------------- |
| Kafka | 분석 엔진, CEP, 외부 시스템 연동용 토픽으로 발행합니다 |
| MQTT  | 실시간 대시보드, 모바일 클라이언트 등으로 발행합니다     |

## 타임아웃 백업 시스템 <a href="#timeout-backup" id="timeout-backup"></a>

파이프라인 처리 중 타임아웃이 발생하면, 데이터 유실을 방지하기 위해 백업 시스템이 동작합니다.

```
DataFlow 처리 중
       │
       ├── 정상 완료 ──▶ 다음 스테이지
       │
       └── 타임아웃 ──▶ RocksDB 백업 저장
                              │
                              ▼
                        Redis 백업 인덱스 등록
                              │
                              ▼
                        복구 스케줄러가 주기적으로 재처리
```

| 항목     | 저장소     | 역할                       |
| ------ | ------- | ------------------------ |
| 데이터 본문 | RocksDB | 타임아웃된 데이터를 로컬 디스크에 저장합니다 |
| 백업 인덱스 | Redis   | 복구 대상 데이터의 키 목록을 관리합니다   |

## 설정 프로퍼티 <a href="#configuration" id="configuration"></a>

파이프라인 관련 주요 설정 항목입니다.

| 프로퍼티                                 | 기본값             | 설명                                 |
| ------------------------------------ | --------------- | ---------------------------------- |
| `pipeline.shard.count`               | `auto`          | 샤드 수. `auto`이면 CPU 코어 수 기반으로 결정됩니다 |
| `pipeline.queue.capacity`            | `65536`         | MpmcArrayQueue 용량입니다               |
| `pipeline.queue.offload.threshold`   | `0.9`           | 오프로드 임계값 (90%)입니다                  |
| `pipeline.queue.reload.threshold`    | `0.65`          | 리로드 임계값 (65%)입니다                   |
| `pipeline.offload.path`              | `/data/offload` | RocksDB 오프로드 저장 경로입니다              |
| `pipeline.dataflow.stream.mode`      | `DISRUPTOR`     | 스트림 모드 (`DIRECT` / `DISRUPTOR`)입니다 |
| `pipeline.dataflow.store.batch.size` | `500`           | Cassandra 배치 쓰기 크기입니다              |
| `pipeline.dataflow.store.async`      | `true`          | 비동기 저장 활성화 여부입니다                   |
| `pipeline.timeout.backup.enabled`    | `true`          | 타임아웃 백업 활성화 여부입니다                  |
| `pipeline.timeout.ms`                | `5000`          | DataFlow 스테이지 타임아웃 (밀리초)입니다        |

## 메트릭 목록 <a href="#metrics" id="metrics"></a>

파이프라인이 수집하는 주요 메트릭 목록입니다.

| 메트릭 이름                          | 타입      | 설명                          |
| ------------------------------- | ------- | --------------------------- |
| `pipeline.ingested.total`       | Counter | 수신된 전체 데이터 포인트 수입니다         |
| `pipeline.processed.total`      | Counter | 처리 완료된 데이터 포인트 수입니다         |
| `pipeline.dropped.total`        | Counter | 검증 실패로 드롭된 데이터 포인트 수입니다     |
| `pipeline.queue.size`           | Gauge   | 현재 큐에 대기 중인 데이터 수입니다        |
| `pipeline.queue.utilization`    | Gauge   | 큐 사용률 (0.0 \~ 1.0)입니다       |
| `pipeline.offloaded.total`      | Counter | 디스크로 오프로드된 데이터 수입니다         |
| `pipeline.reloaded.total`       | Counter | 디스크에서 리로드된 데이터 수입니다         |
| `pipeline.dataflow.latency`     | Timer   | DataFlow 6단계 전체 처리 지연 시간입니다 |
| `pipeline.stage.{name}.latency` | Timer   | 각 스테이지별 처리 지연 시간입니다         |
| `pipeline.timeout.backup.total` | Counter | 타임아웃 백업된 데이터 수입니다           |
| `pipeline.store.batch.latency`  | Timer   | Cassandra 배치 쓰기 지연 시간입니다    |
| `pipeline.dds.publish.latency`  | Timer   | DDS 발행 지연 시간입니다             |

## 패키지 구조 <a href="#package-structure" id="package-structure"></a>

파이프라인 관련 주요 패키지입니다.

```
com.plantpulse.pipeline
├── PipelineEngine                  # 파이프라인 엔진 메인 클래스
├── shard
│   ├── ShardRouter                 # 태그 해시 기반 샤드 라우터
│   └── ShardedWorkerGroup          # 샤드 워커 그룹 관리
├── queue
│   ├── PipelineQueue               # 3계층 큐 관리자
│   ├── MpmcArrayQueue              # Lock-free 인메모리 큐
│   ├── O3Buffer                    # Out-of-Order 정렬 버퍼
│   └── FileOffloadQueue            # RocksDB 디스크 오프로드 큐
├── dataflow
│   ├── DataFlowExecutor            # 6단계 DataFlow 실행기
│   ├── PrepareStage                # Stage 1: 도메인 보강, 타임스탬프 보정
│   ├── ValidateStage               # Stage 2: 8가지 검증
│   ├── CacheStage                  # Stage 3: 인메모리 캐시 갱신
│   ├── StreamStage                 # Stage 4: 스트림 프로세서 전달
│   ├── StoreStage                  # Stage 5: Cassandra 저장
│   └── DDSStage                    # Stage 6: Kafka/MQTT 발행
├── backup
│   ├── TimeoutBackupManager        # 타임아웃 백업 관리자
│   └── BackupRecoveryScheduler     # 백업 복구 스케줄러
└── metrics
    └── PipelineMetrics             # 메트릭 수집 및 리포팅
```
