선행 프로세스 그룹 두 개가 각각 DB 적재를 끝내야 후행 프로세스 그룹 하나가 실행되어야 한다. 하루 한 번 도는 배치이며, 선행과 후행이 데이터를 주고받지는 않는다. 후행은 "둘 다 끝났다" 는 신호만 받으면 된다.
NiFi 는 데이터 흐름 도구라 잡 스케줄러처럼 의존 관계를 선언하는 기능이 없다. 대신 Notify 와 Wait 프로세서로 신호를 주고받아 같은 효과를 낸다.
선행 그룹마다 마지막에 Notify 를 두고, 후행 그룹 앞에 Wait 를 둔다. 신호는 DistributedMapCache 에 저장된다.
선행 PG A : ... → PutDatabaseRecord → Notify (Signal Counter Name = pg_a)
선행 PG B : ... → PutDatabaseRecord → Notify (Signal Counter Name = pg_b)
후행 PG : GenerateFlowFile(CRON) → Wait (Target Signal Count = 2) → 실제 작업
Notify 설정은 다음과 같다.
| 속성 | 값 |
|---|---|
| Release Signal Identifier | ${now():format('yyyyMMdd', 'Asia/Seoul')} |
| Signal Counter Name | pg_a (선행 그룹마다 다르게) |
| Signal Counter Delta | 1 |
| Distributed Cache Service | DistributedMapCacheClientService |
| Signal Buffer Count | 기본값 |
Wait 설정은 다음과 같다.
| 속성 | 값 |
|---|---|
| Release Signal Identifier | ${now():format('yyyyMMdd', 'Asia/Seoul')} |
| Target Signal Count | 2 |
| Signal Counter Name | 비워 둔다. 비우면 모든 카운터의 합을 본다 |
| Wait Mode | Transfer to 'wait' relationship |
| Expiration Duration | 6 hours |
| Distributed Cache Service | 같은 클라이언트 서비스 |
신호 식별자를 날짜로 두면 그날의 신호끼리만 모인다. 배치 ID 를 공유하지 않는 상황에 맞는 방식이다.
Wait 의 wait 관계는 Wait 자신의 입력으로 되돌린다. 그래야 조건이 찰 때까지 다시 확인한다. expired 관계는 로그나 알림으로 보낸다. success 가 조건이 찬 경우다.
Wait 는 스케줄될 때마다 캐시를 확인하므로, Run Schedule 을 10 sec 정도로 두어 불필요한 조회를 줄인다. Penalty Duration 을 함께 쓰면 되돌린 FlowFile 이 곧바로 재시도되는 것을 막는다.
Notify 로 올라간 카운터는 Wait 가 조건을 충족해 소비할 때 삭제된다(Releasable FlowFile Count 기본 동작). 하루가 지나면 식별자가 바뀌므로 이전 날짜의 항목이 캐시에 남는다. DistributedMapCacheServer 는 자동 만료가 없으므로 항목이 계속 쌓인다. 주기적으로 정리하려면 캐시 크기를 넉넉히 잡고 서버를 주기적으로 재기동하거나, 만료가 있는 캐시 구현(예: Redis 기반 클라이언트 서비스)을 쓴다.
DistributedMapCacheServer 는 한 노드에서만 동작한다. 클라이언트 서비스의 Server Hostname 에 그 노드를 지정해야 하며, 해당 노드가 내려가면 신호 전달이 끊긴다. 가용성이 필요하면 외부 캐시를 쓰거나, 배치 창 안에서 사람이 개입할 여지를 남겨 둔다.
선행·후행이 모든 노드에서 각각 도는 구성이면 신호가 노드 수만큼 쌓인다. 배치 성격의 흐름은 Execution 을 Primary node only 로 두어 클러스터 전체에서 한 번만 실행되게 한다.
신호가 아니라 실제 데이터를 넘길 수 있다면 더 단순하다. 선행 그룹의 Output Port 를 후행 그룹의 Input Port 로 잇고, 후행 그룹의 FlowFile Concurrency 를 Single Batch Per Node, Outbound Policy 를 Batch Output 으로 두면 묶음 단위로 순차 실행된다.
건수가 고정돼 있으면 MergeContent 의 Minimum Number of Entries 를 그 수로 두고, 그것이 모였을 때만 다음으로 보내는 방법도 쓸 수 있다.
외부 스케줄러(Airflow 등)가 이미 있다면 NiFi REST API 로 후행 프로세스 그룹을 시작시키는 편이 흐름을 이해하기 쉽다. 의존 관계 관리는 스케줄러가 더 잘한다.
같은 대화에서 NiFi Registry 연동 중 Start Version Control 을 누르면 인증 오류가 나는 문제가 함께 다뤄졌으나 해결까지 이르지 못했다. 확인한 사실은 Registry 가 TLS 로 클라이언트 인증서를 요구하는 상태였다는 점뿐이다. 점검할 곳은 NiFi 쪽 Registry Client 설정의 URL 과 SSL Context Service, Registry 쪽에서 NiFi 노드 인증서의 DN 이 사용자로 등록돼 있는지, 그리고 해당 사용자에게 버킷 쓰기 정책이 있는지다.