카프카의 재시도 처리 방식

박태현·2025년 8월 15일

기타

목록 보기
1/12

Kafka에서는 장애로 인한 문제 혹은 의도적으로 이미 처리한 메세지를 재처리해야 할 때 Kafka의 Offset을 이용하여 메세지를 재처리할 수 있습니다.

분산 환경에서 컨슈머 애플리케이션을 운영할 때, 컨슈머 그룹을 중단하지 않고 Offset을 이동하는 방법

메세지 재처리


컨슈머 리밸런싱 혹은 장애 또는 문제 발생으로 인해 메세지를 처리되지 않은 상황에서는 메세지를 재처리해야 합니다.

특히 컨슈머 로직이 실패했지만, Offset이 커밋되는 상황이 발생하는 경우는 반드시 재처리 해야 합니다.

⇒ 해당 메세지가 처리되지 않았지만 Offset이 커밋되어 이후 재처리할 기회도 잃어버리기 때문

컨슈머 리밸런싱 : 같은 컨슈머 그룹 안에서 각 컨슈머가 파티션의 소유권을 재조정하는 과정


리밸런싱을 하는 경우

  • 컨슈머 인스턴스가 추가되는 경우
  • 컨슈머 인스턴스가 종료되는 경우
  • 구독하는 토픽을 변경하는 경우
  • 파티션의 수를 변경하는 경우

Kafka Offset 커밋 방식 2가지

  1. 자동 커밋 ( enable.auto.commit=true → default)

    기본 설정은 자동 커밋이고, auto.commit.interval.ms 주기마다 최근 읽은 레코드의 오프셋을 커밋

    **자동 커밋의 동작 방식**
    
    컨슈머는 주기적으로 poll() 메서드를 호출하여 메세지를 가져오며, 내부적으로 Kafka 클라이언트는 interval.ms
    주기마다 마지막 poll()에서 반환된 레코드들의 오프셋을 커밋함
    
    반환받은 오프셋 값은 __consumer_offsets 이라는 토픽에 저장하며 이는 Kafka 클러스터 자체에 저장되고
    , 모든 브로커가 공유합니다.
    ⇒ 저장 형식 : Key: (consumer group id, topic, partition), Value: 커밋된 오프셋 값 + 메타데이터
    
    컨슈머가 실행되면, 자신이 속한 consumer group id로 __consumer_offsets 토픽에서 오프셋을 조회하고
    , 조회한 Offset부터 메세지를 읽음

    메세지 처리가 끝났는지 여부는 확인하지 않고 처리한 것을 가정하고 커밋하며, 로직이 중간에 실패하더라도 Offset은 증가할 수 있습니다.

  2. 수동 커밋

    메세지 처리가 정상적으로 끝난 뒤 직접 commitSync() 또는 commitAsync() 명령어를 호출하는 방식

    커밋 위치를 신경써서 작성해야 함

메세지를 재처리 하기 위해서는 Offset 재설정을 통해 할 수 있으며, 카프카 CLI 오프셋 재설정 명령 혹은 애플리케이션 로직을 사용하는 방식이 있습니다.

메세지 재처리 방법


카프카 CLI로 Offset을 조정하기 위해서는 컨슈머를 반드시 중지해야 합니다.

kafka-consumer-groups.sh --reset-offsets 명령은 Offset을 변경하려는 컨슈머 그룹비활성 상태여야 실행되기 때문에 해당 그룹을 사용하는 컨슈머가 실행중이면 명령이 실패하거나 경고가 발생하게 됩니다.

따라서, 해당 컨슈머 그룹을 반드시 중지해야 합니다 ( 해당 컨슈머 그룹을 실행 중인 애플리케이션을 중지해야 함 )

CLI로 오프셋을 이동하려고 하면, Kafka는 현재 오프셋이 어떤 파티션에 적용될지 알아야 하는데, 컨슈머가 실행 중이면 그 시점에도 오프셋이 커밋될 수 있어서 충돌이 발생하게 되기 때문

따라서, 카프카 CLI 도구는 그룹 비활성 상태에서만 오프셋 재설정을 허용합니다.

컨슈머를 중지하고 재가동할 때의 문제

  • 데이터 처리 지연

    중지, 재가동 하는 과정에서 카프카 Lag가 증가함에 따라 처리 지연이 발생할 수 있습니다.

    카프카 Lag :  Topic의 가장 최신 Offset 값과 컨슈머  Offset 값의 차이



  • 다른 consumer에도 영향을 끼침

    특정 애플리케이션이나 컨슈머 그룹에 여러 컨슈머가 포함되어 있으면, 특정 파티션의 Offset만 조정하려는 경우에도 애플리케이션을 중지, 재기동하는 과정에서 동일한 컨슈머 그룹 내 다른 컨슈머들의 메시지 처리까지 중단되는 문제가 발생할 수 있습니다.


  • 가용성 저하 컨슈머를 중단하게 되면 메세지 처리가 불가하기에 서비스의 가용성이 저하되게 됩니다.

    가용성 : 시스템이나 서비스가 필요할 때 정상적으로 동작하고 접근 가능한 정도



Offset 이동 방법


To reset offsets of a consumer group, "–reset-offsets" option can be used. This option supports one consumer group at the time. It requires defining following scopes: –all-topics or –topic. One scope must be selected, unless you use '–from-file' scenario. Also, first make sure that the consumer instances are inactive. See KIP-122 for more details.

Apache Kafka 문서를 보면 컨슈머 그룹을 비활성 상태로 만들기 위해서는 반드시 애플리케이션을 중단해야 한다고 되어 있으므로 애플리케이션을 중단하지 않고 Offset을 변경할 수는 없습니다.

Apache kafka Admin API 사용

org.apache.kafka.clients.admin 패키지 Admin 인터페이스의 Offset 이동 메서드 사용

alterConsumerGroupOffsets(...) 메서드

/*
	Alters offsets for the specified group. In order to succeed, the group must be empty.
	This operation is not transactional so it may succeed for some partitions while fail for others.
	Params:
	groupId – The group for which to alter offsets.
	offsets – A map of offsets by partition with associated metadata. Partitions not specified in the map are ignored.
	options – The options to use when altering the offsets.
	Returns:
	The AlterOffsetsResult.
*/
AlterConsumerGroupOffsetsResult alterConsumerGroupOffsets(
	String groupId,
	Map<TopicPartition, OffsetAndMetadata> offsets,
	AlterConsumerGroupOffsetsOptions options
);

하지만, 이 방식에도 In order to succeed, the group must be empty 문구를 보면 컨슈머 그룹을 비활성 상태로 만들어야 한다고 명시되어 있음

⇒ 이 방식은 사용할 수 없음 .. ..

Spring Kafka의 Offset 이동 기능 사용

Spring Kafka에서 제공하는 Offset 이동 기능을 담은 ConsumerSeekAware 인터페이스와 ConsumerSeekCallback을 사용

  • ConsumerSeekAware & ConsumerSeekCallback ConsumerSeekAware : 오프셋을 어디서부터 읽을지 결정할 타이밍을 제공하는 역할 ConsumerSeekCallback : 오프셋을 실제로 이동( Seek ) 시키는 API
    public interface ConsumerSeekAware {
    
    	default void registerSeekCallback(ConsumerSeekCallback callback) {
    	}
    
    	default void onPartitionsAssigned(Map<TopicPartition, Long> assignments, ConsumerSeekCallback callback) {
    	}
    
    	default void onPartitionsRevoked(Collection<TopicPartition> partitions) {
    	}
    
    	default void onIdleContainer(Map<TopicPartition, Long> assignments, ConsumerSeekCallback callback) {
    	}
    
    	default void onFirstPoll() {
    	}
    
    	default void unregisterSeekCallback() {
    	}
    	
    	interface ConsumerSeekCallback {
    
    		void seek(String topic, int partition, long offset);
    
    		void seek(String topic, int partition, Function<Long, Long> offsetComputeFunction);
    
    		void seekToBeginning(String topic, int partition);
    
    		default void seekToBeginning(Collection<TopicPartition> partitions) {
    			throw new UnsupportedOperationException();
    		}
    
    		void seekToEnd(String topic, int partition);
    
    		default void seekToEnd(Collection<TopicPartition> partitions) {
    			throw new UnsupportedOperationException();
    		}
    
    		void seekRelative(String topic, int partition, long offset, boolean toCurrent);
    
    		void seekToTimestamp(String topic, int partition, long timestamp);
    
    		void seekToTimestamp(Collection<TopicPartition> topicPartitions, long timestamp);
    		
    		@Nullable
    		default String getGroupId() {
    			return null;
    		}
    	}
    
    }

두 인터페이스를 함께 구현 해놓은 AbstractConsumerSeekAware 헬퍼 추상 클래스를 상속하면, 컨슈머에서 토픽, 파티션별 오프셋 이동 콜백을 쉽게 등록하고 제거하는 기능을 구현할 수 있습니다.

두 인터페이스를 직접 구현하면 파티션별 오프셋 관리, 콜백 등록/제거 로직을 매번 작성해야 해서 번거로움
AbstractConsumerSeekAware가 이런 등록, 제거, 저장 로직을 기본 제공
개발자는 오프셋을 언제·어디로 이동할지만 구현하면 됨



또한, Spring 카프카 컨슈머는 polling 전에 할당된 토픽 파티션에 대해 Offset 이동 콜백을 자동으로 추적하고 처리하므로, 개발자는 오프셋 제어에 대한 부담을 덜 수 있습니다

// 스프링 카프카 컨슈머가 실행되면 pollAndInvoke() 메서드가 주기적으로 호출됨
protected void pollAndInvoke() {

		// seeks: 오프셋 이동 요청 대기열
		// 이 큐에 요청이 있으면 processSeeks()를 호출하여 실제 KafkaConsumer의 오프셋을 원하는 위치로 변경
		// 이렇게 하면 다음 poll 시점부터 해당 오프셋 이후의 메시지를 읽음
    if (!this.seeks.isEmpty()) { 
        processSeeks(); // 대기열의 요청을 처리하여 오프셋을 미리 옮깁니다.
    }
    
		// 변경된 오프셋부터 메시지를 얻음
    ConsumerRecords<K, V> records = doPoll(); 

		// 읽어온 메세지가 있다면 리스너를 호출
		// 리스너 : 사용자가 구현한 메시지 처리 로직
    invokeIfHaveRecords(records); 
}

pollAndInvoke()는 계속 반복되며, 오프셋 이동 요청이 있으면 먼저 처리 → 메시지 읽기 → 리스너 호출
과정을 계속 자동으로 수행

AbstractConsumerSeekAware를 상속하도록 하여, 토픽 파티션 별로 Offset을 맨 처음으로 옮기는 기능

@Component
public class SimpleSeekAwareListener extends AbstractConsumerSeekAware {

    @KafkaListener(topics = "my-topic", groupId = "my-group-id")
    public void listen(String message) {
        System.out.println("Received: " + message);
    }

    // 오프셋을 맨 처음으로 옮기는 메서드
    public void seekToEarliest() {

        // 토픽 파티션 별
        this.getTopicsAndCallbacks()
            .forEach((topicPartition, callbacks) -> {
                callbacks.forEach(cb -> cb.seekToBeginning(topicPartition.topic(),
									topicPartition.partition()));
            });
    }
 }

ConsumerSeekCallback에서 제공하는 Offset 이동 메서드 종류

  1. seek(String topic, int partition, long offset)

    지정한 토픽, 파티션의 Offset을 원하는 숫자 위치로 이동

    Ex ) cb.seek("my-topic", 0, 15) → 0번 파티션을 오프셋 15부터 읽기 시작

  2. seekToBeginning(String topic, int partition)

    지정한 파티션의 가장 처음 Offset으로 이동

    Ex ) cb.seekToBeginning("my-topic", 0) → 특정 파티션 번호의 파티션을 처음부터 읽어라

  3. seekToEnd(String topic, int partition)

    지정한 파티션의 가장 마지막 오프셋( 현재까지 데이터의 끝 )으로 이동

    Ex ) cb.seekToEnd("my-topic", 0) → 특정 파티션 번호의 파티션을 마지막부터 읽어라

위 방식을 사용한다면 애플리케이션을 중단하거나 컨슈머를 중지하지 않고도 Offset을 변경할 수 있습니다.

분산 환경에서의 Offset 변경


Kafka 컨슈머 그룹은 일반적으로 확장성과 내결함성을 위해 여러 호스트에 분산된 컨슈머들로 구성

내결함성 : 하나의 서버가 장애로 다운되더라도 다른 서버의 컨슈머가 이어받아 처리

확장성 : 여러 서버에서 동시에 파티션을 나눠 병렬적으로 처리하므로 처리 속도가 증가



⭐️⭐️

Spring Kafka의 기능은 개별 호스트 수준에서만 적용된다는 한계가 존재하는데, 물론 각 호스트별로 기능을 수행하도록 할 수는 있지만 번거롭기 때문에 개선이 필요

분산 환경에서 Offset 이동 요청을 각 서버에 전파할 수 있는 방법

  1. 각 서버에 오프셋 이동 요청을 받을 수 있는 HTTP API를 생성

    이 API를 호출하면 해당 서버에서 동작 중인 컨슈머가 seek()을 실행하도록 하는 방식

    ⇒ 해당 API가 호출되었을 때 실행될 seek() 메서드가 존재해야함

    • topics: 토픽 목록
    • partitions: 파티션 목록 ( 설정하지 않는다면 모든 파티션에서 오프셋 이동 )
    • seekAt: 오프셋 이동 시작 일시 ( 특정 시간 이후에 발행된 첫 번째 메시지부터 처리를 다시 시작하고 싶을 때 사용 )
    POST /kafka/seek
    Content-Type: application/json
    
    {
      "topics": ["order-topic", "payment-topic"],
      "partitions": [0, 1],
      "seekAt": "2025-08-14T10:00:00Z"
    }

  2. Redis Pub/Sub 사용

    모든 서버가 Redis의 특정 채널을 구독하도록 하고 사용자가 오프셋 이동을 요청하면 분산 서버 중 하나가 이 요청을 Redis 채널에 전송하여 이를 Redis가 수신하면 각 서버로 컨슈머 오프셋 이동 메서드를 호출하도록 하는 방식

이를 통해 개별 호스트에서 처리되던 Offset 이동을 분산 환경에서 컨슈머 그룹 레벨로 확장 가능합니다.

profile
꾸준하게

0개의 댓글