ConsumeKafkaRecord → EvaluateJsonPath → UpdateAttribute → RouteOnAttribute 로 나뉘어 있던 단계를 ExecuteScript(Groovy) 하나로 합친다. Filebeat 가 보내는 메시지는 [{"@timestamp":..., "svc":"ct", "category":"messenger", "regn":"KR", "host":{"name":...}, "message":"..."}] 형태다. 앞 단 UpdateAttribute 에서 Parameter #{regn_cd} 를 속성 regn_cd 로 넣어 두고, 스크립트가 메시지의 regn 과 비교해 다른 리전 메시지를 거른다.
import groovy.json.JsonSlurper
def flowFile = session.get()
if (!flowFile) return
try {
def regnCd = flowFile.getAttribute("regn_cd") ?: ""
def content = ""
session.read(flowFile, { inputStream ->
content = inputStream.text
} as org.apache.nifi.processor.io.InputStreamCallback)
def json = new JsonSlurper().parseText(content)
def msg = json instanceof List ? json[0] : json
def svc = msg.svc ?: ""
def category = msg.category ?: ""
def hostName = msg.host?.name ?: ""
def message = msg.message ?: ""
def timestamp = msg.'@timestamp' ?: ""
def regn = msg.regn ?: ""
def route = "unmatched"
if (regn == regnCd && category == "messenger") {
if (svc in ["ct", "la", "ms", "rm"]) route = svc
}
flowFile = session.putAllAttributes(flowFile, [
"svc": svc, "category": category, "host_name": hostName,
"message": message, "timestamp": timestamp, "regn": regn, "route": route
])
session.transfer(flowFile, REL_SUCCESS)
} catch (Exception e) {
log.error("파싱 오류: ${e.message}", e)
session.transfer(flowFile, REL_FAILURE)
}
ExecuteScript 의 관계는 success · failure 뿐이라 session.transfer(flowFile, 'unmatched') 처럼 임의 관계로 보낼 수 없다(Relationship 'unmatched' is not known). 라우팅 결과는 속성 route 에 담고 뒤에 RouteOnAttribute 를 하나 두어 ${route:equals('ct')} 식으로 분기한다. unmatched 는 RouteOnAttribute 의 기본 관계라 그쪽에서 처리한다.
| 케이스 | regn | category | svc | route |
|---|---|---|---|---|
| 정상 | KR | messenger | ct | ct |
| svc 미매칭 | KR | messenger | smailetc | unmatched |
| regn 불일치 | JP | messenger | ct | unmatched |
| category 불일치 | KR | smail | ct | unmatched |
| 배열 형태 입력 | KR | messenger | ct | ct |
| regn 없음 | 빈값 | messenger | ct | unmatched |
#{etl_db} 같은 Parameter 는 ExecuteScript 의 동적 속성으로 등록한 뒤 context.getProperty('database').evaluateAttributeExpressions(flowFile).getValue() 로 읽을 수 있다. 동적 속성 등록 없이 context.getProperty('#{etl_db}') 로는 읽히지 않으므로, 간단히 하려면 앞 단 UpdateAttribute 에서 속성으로 만들어 넘긴다.
message 속성을 컬럼으로 나누는 정규식.