프라이머리 노드와 클러스터 조정자는 선출로 정해진다. 관리자가 특정 노드를 프라이머리로 못 박는 설정은 없다. 선출은 ZooKeeper 의 리더 선출(또는 NiFi 2.x 의 Kubernetes Lease)이 담당하며, 현재 프라이머리가 사라지면 남은 노드 중 하나가 자동으로 승계한다.
현재 어느 노드가 프라이머리인지는 UI 우측 상단 메뉴의 Cluster 화면에서 본다. 노드 목록의 Primary 열과 Coordinator 열이 그것이다.
프라이머리를 옮겨야 한다면 현재 프라이머리 노드를 Disconnect 하거나 재기동한다. 그러면 남은 노드에서 새로 선출된다. 다시 붙인 노드가 프라이머리로 돌아오지는 않는다.
선출이 이뤄지지 않고 프라이머리가 없는 상태가 이어지면 ZooKeeper 쪽을 본다. 앙상블이 과반을 잃었거나, 노드마다 Connect String · Root Node 가 다르거나, 이전 설치의 znode 가 남아 있는 경우다.
echo stat | nc zk01.example.com 2181
/opt/zookeeper/bin/zkCli.sh -server zk01.example.com:2181 ls /nifi
Primary node only 로 설정된 프로세서는 프라이머리가 바뀌는 순간 옛 노드에서 멈추고 새 노드에서 시작된다. 이 전환 중에는 잠시 아무 노드도 실행하지 않는 구간이 생기므로, 주기가 짧은 수집 프로세서에서는 한두 주기를 건너뛸 수 있다.
프로세스 그룹 설정에는 FlowFile Concurrency 와 Outbound Policy 가 있다. 그룹 안으로 한 번에 얼마나 들여보낼지, 그리고 밖으로 언제 내보낼지를 정한다.
| FlowFile Concurrency | 의미 |
|---|---|
| Unbounded | 제한 없이 들어온다. 기본값 |
| Single FlowFile Per Node | 노드마다 한 번에 하나의 FlowFile 만 그룹 안으로 들인다. 그 하나가 그룹을 빠져나가야 다음이 들어온다 |
| Single Batch Per Node | 입력 포트에 대기 중이던 묶음을 한 번에 들이고, 그 묶음이 전부 나갈 때까지 다음 묶음을 들이지 않는다 |
| Outbound Policy | 의미 |
| --- | --- |
| Stream When Available | 준비되는 대로 내보낸다 |
| Batch Output | 그룹 안의 처리가 모두 끝난 뒤 한꺼번에 내보낸다 |
노드마다 독립적으로 적용된다는 점이 중요하다. 3 노드 클러스터에서 Single FlowFile Per Node 라면 클러스터 전체로는 동시에 3 건이 흐른다. 클러스터 전체에서 하나만 흐르게 하려면 그 앞 연결을 Single node 분배로 몰거나 Primary node only 로 수집해야 한다.
NiFi 가 보장하는 것은 하나의 연결 안에서의 순서뿐이며, 그것도 큐의 우선순위 설정에 따라 달라진다. 여러 노드, 여러 스레드로 갈라진 뒤에는 처리 순서가 보장되지 않는다. 따라서 순서가 중요한 흐름은 다음 중 하나를 쓴다.
첫째, 흐름을 한 노드로 모은다. 연결의 분배 전략을 Single node 로 두거나 소스 프로세서를 Primary node only 로 둔다. 그리고 해당 프로세서의 Concurrent Tasks 를 1 로 유지한다.
둘째, 키 단위 순서만 필요하면 Partition by attribute 로 같은 키를 같은 노드에 모은다. Kafka 파티션과 같은 개념이다.
셋째, 순서 번호가 데이터에 들어 있으면 EnforceOrder 프로세서를 쓴다. 지정한 속성의 값을 순번으로 보고, 앞 순번이 도착할 때까지 뒤 순번을 대기시킨다. 직접 정렬 스크립트를 짜는 것보다 안전하다.
| EnforceOrder 속성 | 설명 |
|---|---|
| Group Identifier | 순서를 따질 묶음 키. 예: ${batch.id} |
| Order Attribute | 순번이 담긴 속성 이름 |
| Initial Order | 시작 순번 |
| Maximum Order | 마지막 순번. 비워 두면 제한 없음 |
| Wait Timeout | 앞 순번을 기다리는 한도. 넘으면 overtook · failure 로 보낸다 |
순번을 붙이는 쪽은 별도의 처리가 필요하다. 소스가 순번을 주지 않으면 Primary node only 로 도는 한 프로세서에서 DistributedMapCache 에 카운터를 두고 붙인다. 이때도 그 프로세서의 Concurrent Tasks 는 1 이어야 한다.
Groovy 스크립트로 큐를 통째로 읽어 정렬하는 방법도 쓰이지만, 한 번에 세션으로 가져온 범위 안에서만 정렬되므로 전체 순서를 보장하지 못한다. 배치 경계에 걸친 순서가 필요하면 쓰지 않는다.
PutKudu 는 FlowFile 하나를 레코드 여러 건으로 읽어 Kudu 세션에 넣고, Batch Size 만큼 모이면 flush 한다. 즉 FlowFile → 레코드 → 배치의 관계이며, 배치 경계는 FlowFile 경계와 일치하지 않을 수 있다. 순서가 중요한 적재에서는 Flush Mode 와 Batch Size 를 함께 본다.
| 속성 | 설명 |
|---|---|
| Record Reader | FlowFile 내용을 레코드로 읽는 컨트롤러 서비스 |
| Flush Mode | AUTO_FLUSH_SYNC · AUTO_FLUSH_BACKGROUND · MANUAL_FLUSH. 순서와 오류 확인이 중요하면 동기 방식 |
| Batch Size | 한 번에 flush 할 레코드 수 |
AUTO_FLUSH_BACKGROUND 는 처리량이 높지만 오류가 비동기로 돌아와 어느 레코드가 실패했는지 추적하기 어렵다.