NiFi 1.x 의 ExecuteScript 프로세서에서 Script Engine 으로 python 을 고르면 CPython 이 아니라 JVM 위에서 도는 Jython 이 돈다. 결과적으로 두 가지 제약이 따라온다. 첫째, C 확장에 의존하는 패키지(pandas · numpy · requests 의 일부 등)를 쓸 수 없다. 둘째, Jython 이 지원하는 언어 수준이 Python 2.7 에 머물러 있다.
그래서 실무에서는 같은 자리에 Groovy 를 쓰는 경우가 많다. JVM 객체를 그대로 다루고 성능도 낫다.
Jython 으로 FlowFile 내용을 고치는 기본 형태는 다음과 같다.
from org.apache.nifi.processor.io import StreamCallback
from org.python.core.util import StringUtil
class Upper(StreamCallback):
def process(self, inputStream, outputStream):
text = StringUtil.toString(inputStream.readFully())
outputStream.write(StringUtil.toBytes(text.upper()))
flowFile = session.get()
if flowFile is not None:
flowFile = session.write(flowFile, Upper())
session.transfer(flowFile, REL_SUCCESS)
session.get() 이 None 을 돌려주는 경우를 반드시 처리한다. 이 검사를 빼면 큐가 비었을 때 예외가 반복된다.
NiFi 2.x 는 Jython 기반 스크립팅을 걷어내고, CPython 으로 도는 Python Processor API 를 새로 넣었다. 프로세서를 파이썬 클래스로 작성해 python/extensions 아래에 두고, NiFi 가 별도 파이썬 프로세스로 띄워 통신하는 구조다. 1.x 의 ExecuteScript 파이썬 코드는 그대로 옮겨지지 않으므로 이관 대상으로 잡는다.
GenerateFlowFile 프로세서로 내용을 고정해 찍어낼 수 있다. 여기서 순번을 넣을 때 쓰는 것은 NiFi 표현 언어의 nextInt() 함수다. JVM 안에서 1씩 올라가는 값을 돌려준다.
Record ${nextInt()}: sequential test line
| 속성 | 값 |
|---|---|
| Custom Text | 위 문자열 |
| File Size | 0 B (Custom Text 를 쓰면 무시된다) |
| Batch Size | 한 번에 만들 FlowFile 수 |
| Run Schedule | 만들어 내는 주기 |
nextInt() 는 NiFi 인스턴스가 재기동되면 0 부터 다시 시작하고 클러스터 노드마다 따로 센다. 중복 없는 전역 순번이 필요하면 UpdateAttribute 의 상태 저장(Stateful) 기능이나 ExecuteScript 에서 프로세서 상태를 쓰는 쪽으로 간다.
여러 줄이 든 파일이 필요하다면 GenerateFlowFile 로 한 줄짜리를 만들고 MergeContent 로 묶는 편이 다루기 쉽다. 이후 PutFile · PutHDFS · PublishKafka 로 흘려보낸다.