Kafka 토픽의 JSON 메시지를 ConsumeKafkaRecord 로 받아 svc 값에 따라 여러 Kudu 테이블로 분기 적재하는 PG 에서, ExtractText 정규식으로 잡히지 않아 유실되던 메시지를 ExecuteScript(Groovy) 파서로 바꾼 패턴과, 파서 교체 후 과거 메시지를 다시 받기 위한 컨슈머 그룹 offset 리셋 절차를 정리한다.
메시지 하나를 읽어 svc 로 route 를 정하고 필요한 값을 속성으로 올린다. 대상 외 메시지는 여기서 버리고, 특정 svc 안에 여러 로그 종류가 섞여 있으면 라우터에서 message 내용 조건까지 걸러 파서는 파싱만 하게 한다.
import groovy.json.JsonSlurper
def flowFile = session.get()
if (!flowFile) return
try {
def content = ""
session.read(flowFile, { ins -> content = ins.text } as org.apache.nifi.processor.io.InputStreamCallback)
def json = new JsonSlurper().parseText(content)
def msg = json instanceof List ? json[0] : json // SplitJson 뒤에는 단일 객체
def svc = (msg.svc ?: "").toString()
def message = (msg.message ?: "").toString()
def regn = (msg.regn ?: "").toString()
def routeMap = ["log-in":"login", "log-out":"logout", "was_send":"was_send", "was_desktop":"desktop_mail"]
def route = routeMap[svc] ?: "unmatched"
if (route == "unmatched") { session.remove(flowFile); return }
def regnCd = regn?.trim() ? regn.trim().toUpperCase() : "UNKNOWN"
flowFile = session.putAllAttributes(flowFile, ["svc":svc, "message":message, "regn_cd":regnCd, "route":route])
session.transfer(flowFile, REL_SUCCESS)
} catch (Exception e) {
log.error("라우팅 파싱 오류: ${e.message}", e)
session.transfer(flowFile, REL_FAILURE)
}
^(\d{4}-\d{2}-\d{2}\s\d{2}:\d{2}:\d{2}\.\d{3}))와 regn_cd 처럼 항상 존재해야 하는 PK 가 비면 parse_status=PK_MISSING 을 붙여 REL_FAILURE 로 보낸다."NA" 로 치환해 적재한다. Kudu 는 PK NULL 을 허용하지 않아 그대로 두면 행 전체가 유실된다.parse_status 값(NO_MESSAGE / PK_MISSING / NO_LOGGER / FIELD_SHORT / PARSE_ERROR)을 남겨 실패 원인을 큐에서 바로 볼 수 있게 한다.ExecuteScript 의 Script Engine 이 Groovy 로 되어 있는지 확인한다. 엔진을 바꾸지 않으면 곧바로 파싱 오류가 난다.SELECT substr(ts,1,10), count(*) ... GROUP BY 1 로 일자별 건수를 보고, 원천이 준 "원본 건수 / PK 기준 건수" 와 대조한다. PutKudu 가 UPSERT 이므로 PK 가 같은 행은 집약된다.파서를 바꾼 뒤 retention 안에 남아 있는 메시지를 처음부터 다시 받으려면 --to-earliest 를 쓴다. 특정 시각(--to-datetime)은 그 시각이 retention 경계와 겹치면 어차피 log-start-offset 으로 붙는다.
BROKERS="broker1:9093,broker2:9093,broker3:9093"
CFG=/tmp/kafka-ssl.properties # security.protocol=SSL, truststore 경로, 비밀번호는 ${KAFKA_TRUSTSTORE_PW}
kafka-consumer-groups --bootstrap-server "$BROKERS" --command-config "$CFG" --list
kafka-consumer-groups --bootstrap-server "$BROKERS" --command-config "$CFG" --group <group> --describe
kafka-consumer-groups --bootstrap-server "$BROKERS" --command-config "$CFG" --group <group> --topic <topic> \
--reset-offsets --to-earliest --dry-run
kafka-consumer-groups --bootstrap-server "$BROKERS" --command-config "$CFG" --group <group> --topic <topic> \
--reset-offsets --to-earliest --execute
순서는 해당 PG 의 ConsumeKafka 정지 → --describe 에서 CONSUMER-ID / HOST 가 - 인지 확인 → --dry-run → --execute → ConsumeKafka 재시작이다. 컨슈머가 살아 있으면 리셋이 거부된다. kafka-consumer-groups 의 설정 옵션은 --command-config 이고 kafka-console-consumer 는 --consumer.config 다.
Consumer group '0' does not exist 가 나면 배열을 정의한 셸과 실행한 셸이 달라 변수가 빈 것이다. 4개뿐이면 그룹명을 리터럴로 직접 쓴다.regn_cd 가 GB 처럼 예상 밖 코드로 들어오는 것은 원본 값이므로 파서에서 임의 변환하지 말고 원천 담당자와 확인한다.