
@Slf4j
@RestController
@RequiredArgsConstructor
@RequestMapping("/api")
public class KafkaConsumerController {
private final KafkaConsumerService kafkaConsumerService;
private final Long TIMEOUT = 300_000L;
@KafkaListener(topics = "normal-topic", groupId = "kafkaGroup1", errorHandler = "kafkaErrorHandler")
@SendTo("error-topic")
public void listener(String message) {
log.info("listener {} ", message);
kafkaConsumerService.listen(message);
}
@KafkaListener(topics = "error-topic", groupId = "kafkaGroup2" , errorHandler = "kafkaErrorHandler")
@SendTo("error-topic")
public void errorListener(String message) {
log.info("errorListener {} ", message);
kafkaConsumerService.errorListen(message);
}
}
리스너의 로직에서 에러가 발생했을 때 그 에러를 kafkaErrorHandler 라는 빈이 처리한다는 로직이다.
@SendTo와 함께 쓸 수 있다.
@Bean
public KafkaListenerErrorHandler kafkaErrorHandler() {
return (m, e) -> {
try {
// m.getPayload()를 JsonNode로 파싱
JsonNode jsonNode = objectMapper.readTree(m.getPayload().toString());
// "message" 필드의 값 추출
String message = jsonNode.path("message").asText();
// 추출된 message 값 로그로 출력
log.error("[KafkaErrorHandler] Extracted message=[" + message + "], errorMessage=[" + e.getMessage() + "]");
// 이후 처리할 메시지를 반환
return message; // message 값만 반환하여 sendTo 토픽으로 전송
} catch (Exception ex) {
log.error("[KafkaErrorHandler] Failed to parse JSON payload, error: " + ex.getMessage());
return m.getPayload(); // 파싱 실패 시 원본 payload 반환
}
};
}
KafkaListenerErrorHandler에서 반환되는 결과 값이 @send to 에서 설정하는 토픽으로 자동 보내지게 된다.
KafkaListenerErrorHandler는 리스너 부분에서 에러가 났을 때 동작 하며 backOff 와 같은 기능을 지원하고 있지 않기 때문에 defaultErrorHandler를 쓰는 것이 좋다.
defaultErrorHandler는 다른 블로그 레퍼런스들이 많이 없어서 공식문서를 보고 이해한 바를 다음 글에서 살펴보도록 하겠다.