새 Consumer Group 은 커밋된 offset 이 없어 auto.offset.reset 값에 따라 시작 위치가 정해진다. 기본값 latest 면 그룹 등록 이후 들어오는 메시지만 받고 기존 메시지는 건너뛴다. 과거 메시지부터 받으려면 earliest 로 바꾼다.
확인 순서는 다음과 같다.
kafka-console-consumer.sh --bootstrap-server <broker> --topic <topic> --from-beginning 으로 본다.READ 권한을 준다.kafka-consumer-groups.sh --bootstrap-server <broker> --describe --group <group> 으로 그룹이 등록됐는지 본다.토픽마다 Group ID 를 나누는 것은 offset 관리가 독립돼 오히려 권장되는 패턴이다.
auto.offset.reset 은 커밋된 offset 이 없을 때만 적용된다. earliest 로 정상 수신을 확인하고 offset 이 커밋된 뒤 latest 로 되돌려도 커밋된 offset 부터 이어 읽으므로 유실이 없다. 다만 그룹이 오래 비활성이면(기본 7 일) offset 이 삭제되어 다시 auto.offset.reset 이 적용되므로, 장기 운영에서는 earliest 로 두는 편이 안전하다.
한국시간 매일 03:00 은 UTC 기준 전날 18:00 이다. NiFi 의 6 자리 cron 은 0 0 18 * * ? 이다.
ConsumeKafkaRecord → EvaluateJsonPath → UpdateAttribute → RouteOnAttribute → ReplaceText → ExtractText → AttributesToJSON → PutKudu 파이프라인에서 각 프로세서의 failure 와 unmatched 를 auto-terminate 해 두면 실패 건이 그냥 버려진다. 로그를 남기려면 다음 순서로 바꾼다.
failure · unmatched 체크를 해제한다.error.source 등 출처 속성을 붙이고 LogAttribute 또는 PutKudu(에러 테이블)로 보낸다.관계를 해제한 뒤 연결하기 전까지는 프로세서가 시작되지 않으므로, 운영 중이면 한 프로세서씩 순차로 작업한다.
ConsumeKafkaRecord 는 FlowFile 속성에 kafka.topic 을 넣어 준다. RouteOnAttribute 에 토픽별 조건을 두면 된다.
| 속성 이름 | 값 |
|---|---|
topic_a |
${kafka.topic:equals('topicA')} |
topic_b |
${kafka.topic:equals('topicB')} |