Kafka 기초 개념을 정리하며 경험했었던 Consumer Group과 Offset의 개념, 그리고 Kafka에서의 Commit이란 무엇인가에 대해 정리해보겠다.
☑️ Offset과 Consumer Group이 뭐였었지?
Offset이란?
- Offset은 각 파티션에서 메시지에 붙는 번호
- 파티션마다 0,1,2,3 과 같이 순서대로 증가
- Consumer는 어떤 파ㅣㅌ셔넹서 어디까지 읽었는지에 대한 기준을 Offset을 통해 기억
Consumer Group과 Offset
- Offset은 Consumer Group 단위로 관리
- 어디까지 읽었는지에 대해서는 Consumer Group + Partition 조합으로 구분된다.
- Consumer Group이 다르다면 같은 토픽을 읽더라도 어디까지 읽었는가에 대해서는 서로 독립적이라고 할 수 있다.
✅ Commit은 무엇인가?
Offset은 번호, Commit은 그 번호를 저장하는 행위라고 말할 수 있겠다.
(GitHub로 작업하면서 수없이 진행했던 커밋의 의미와 크게 다르지 않다!)
Commit의 의미
- Commit은 이 Offset까지 읽었습니다. 라고 Kafka에 기록해 두는 것
- Consumer가 다시 시작되면 마지막으로 커밋한 Offset 이후부터 메시지를 계속 읽는다.
왜 Commit이 중요한가?
- Commit이 없다면 Consumer가 재시작될 때마다 처음부터 다시 읽어야 한다.
- Commit을 너무 빨리 해버리면 아직 처리하지 못한 메시지도 읽었다고 표시될 수 있고, 이는 데이터 유실 위험이 생기게 된다.
- Commit을 너무 늦게 해버리면 이미 처리한 메시지를 다시 읽어버릴 수 있고, 이는 중복 처리 위험이 생기게 된다.
Kafka는 기본적으로 중복되더라도 일단 한 번 이상은 꼭 전달하자라는 철학을 갖고 있다. (at-least-once 전달 방식)
✅ Spring Kafka에서 Commit은 어떻게 동작하는가?
Spring Kafka 는 @KafkaListener 뒤에서 Listener Container를 사용하며
이 컨테이너가 메시지를 가져와서 @KafkaListener 메서드에 넘겨주고, 적절한 시점에 Offset을 Commit 해주는 것
자동으로만 진행되었던 Commit을 수동으로 구현 및 실제로 진행해보며 어떤 흐름으로 동작하는지를 확인할 수 있었는데 그 내용은 아래와 같았다.
수동 Commit용 Listener 설정 추가
import org.springframework.kafka.listener.ContainerProperties;
import org.springframework.kafka.config.ConcurrentKafkaListenerContainerFactory;
// ...
// 수동 Commit(Ack)용 ConsumerGroup 생성
@Bean
public ConsumerFactory<String, String> manualConsumerFactory() {
Map<String, Object> props = new HashMap<>();
// Kafka 서버 위치
props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers);
// 이 Consumer가 속할 그룹 아이디
props.put(ConsumerConfig.GROUP_ID_CONFIG, "manual-ack-group");
// key, value를 문자열로 역직렬화
props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
// 처음 시작할 때, 토픽 맨 앞(earliest)부터 읽도록 설정
props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest");
return new DefaultKafkaConsumerFactory<>(props);
}
// 수동 Commit(Ack)용 ListenerContainerFactory
@Bean
public ConcurrentKafkaListenerContainerFactory<String, String>
manualAckKafkaListenerContainerFactory() {
ConcurrentKafkaListenerContainerFactory<String, String> factory =
new ConcurrentKafkaListenerContainerFactory<>();
factory.setConsumerFactory(manualConsumerFactory());
// AckMode를 MANUAL로 설정
factory.getContainerProperties().setAckMode(ContainerProperties.AckMode.MANUAL);
return factory;
}
- AckMode.MANUAL 설정을 진행하여
- @KafkaListener 메서드에서 Acknowledgement 객체를 받아서
- 직접 ack.acknowledge()를 호출할 수 있다. 이 호출 시점에 Offset이 커밋된다.
수동 Commit Listener 작성
import lombok.extern.slf4j.Slf4j;
import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.springframework.kafka.annotation.KafkaListener;
import org.springframework.kafka.support.Acknowledgment;
import org.springframework.stereotype.Component;
@Slf4j
@Component
public class ManualAckConsumerListener {
@KafkaListener(
topics = "simple-messages",
groupId = "manual-ack-group",
containerFactory = "manualAckKafkaListenerContainerFactory"
)
public void consume(ConsumerRecord<String, String> record,
Acknowledgment ack) {
String message = record.value();
log.info("[manual-ack] 받은 메시지: {}", message);
if (message.equals("error")) {
throw new RuntimeException("kafka 에러 발생");
}
// 처리에 성공했다고 판단되면 Offset을 Commit합니다.
ack.acknowledge();
}
}
- ConsumerRecord<String, String>을 사용하면 메시지 값 뿐만 아니라 Offset, 파티션 정보를 함께 확인할 수 있다.
- Acknowledgement ack 를 통해 메시지의 처리 완료 상태와 Offset을 Commit해도 좋다고 직접 알려줄 수도 있다.
진행해보며 ack.acknowledge()를 호출하지 않아도 커밋되는 상황이 있었다.
Error가 들어오면 여러번 재시도 된 뒤 재시도 횟수를 넘기면 메시지가 들어오지 않는 상황이 있었다.
이런 현상이 발생하는 이유로는
- @KafkaListener에서 예외가 발생하면 메시지 처리 제어권은 에러 핸들러에게 넘어간다.
- Spring Kafka의 기본 에러 핸들러는 설정된 횟수만큼 동일한 메시지를 재시도한다.
- 재시도 후에도 계속 실패하면, 이 메시지는 더이상 처리할 수 없다고 판단하고 그 레코드를 포기하는 과정에서 해당 Offset을 커밋할 수 있다.
→ 정상 처리 흐름에서는 ack.acknowledge() 를 호출한 시점에 커밋이 일어나고, 에러 흐름에서는 에러 핸들러가 재시도를 포기하는 순간 Offset이 커밋될 수 있다.