pile·
백엔드·마켓컬리마켓컬리 Hello World·

분산 시스템 환경에서 Kafka Consumer 오프셋 이동하기

다른 조직이 운영하는 Kafka 클러스터의 배송 토픽을 구독하는 주문팀이, 권한 없이도 컨슈머를 중단하지 않고 오프셋을 이동해 메시지를 재처리하는 방법을 찾은 과정이다. CLI/Admin API가 모두 컨슈머 그룹 비활성을 요구하는 제약을 Spring Kafka의 seek 기능으로 넘고, 분산 환경까지 확장한다.

핵심 포인트
  • kafka-consumer-groups.sh reset-offsets와 Admin API의 alterConsumerGroupOffsets는 모두 컨슈머 그룹이 empty/비활성이어야 해서 무중단 요건을 못 채운다.
  • Spring Kafka의 ConsumerSeekAware/ConsumerSeekCallback(내부적으로 Consumer.seek API)을 쓰면 애플리케이션을 멈추지 않고 오프셋을 이동할 수 있다.
  • AbstractConsumerSeekAware를 상속하면 파티션별 seek 콜백 등록·제거가 쉬워지고, 컨슈머는 poll 전에 seek 대기열을 처리한다.
  • Spring Kafka는 개별 호스트 단위로만 동작하므로, HTTP API + Redis Pub/Sub로 seek 요청을 분산 서버 전체에 전파해 컨슈머 그룹 레벨로 확장했다.
상세 정리
  • 제약 배경: 물류팀 Kafka 클러스터에 관리자 권한이 없어 오프셋 이동 때마다 타 팀 도움이 필요했고, 기존 방식은 명령 실행 전 컨슈머를 반드시 중지해야 했다.
  • 중단의 비용: 실시간 배송(컬리나우, 주문 1시간 내 배달) 특성상 가용성 저하, Lag 처리 지연, 한 그룹 내 다른 컨슈머까지 동반 중단되는 문제가 있었다.
  • 대안1(권한 획득): 조직 간 권한 공유는 관리가 어렵고, 얻어도 CLI reset-offsets가 비활성 상태를 요구해 무중단에 실패한다.
  • 대안2(Admin API): alterConsumerGroupOffsets도 "group must be empty"를 명시해 무중단에 실패한다.
  • 대안3(Spring Kafka): Java 21 + Spring Boot 3.3.4 + Spring for Apache Kafka 3.2.4 스택에서 seek 기능으로 무중단 이동이 가능해 채택.
  • 코드: AbstractConsumerSeekAware를 상속한 리스너에서 파티션별 콜백으로 seekToBeginning 등을 호출해, 앱·컨슈머 중단 없이 오프셋이 처음으로 이동하는 로그를 확인.
  • 버전 주의: getTopicsAndCallbacks는 3.3.0부터 지원되고, 이전 getSeekCallbacks는 같은 토픽을 다른 그룹이 구독하면 콜백이 누락될 수 있다(GH-3328).
  • 분산 확장: seek 요청용 HTTP API(topics/partitions/seekAt)를 정의해 분산 서버 중 하나가 받아 Redis 채널에 게시하고, 각 서버의 Redis 리스너가 수신해 자기 오프셋을 이동한다. 이미 쓰던 Redis라 추가 인프라가 불필요했다.
  • 부가 확장·기여: 같은 플로우로 컨슈머 시작/중지 기능도 확장했고, seek 조사 중 발견한 Spring Kafka 이슈를 메인테이너와 논의해 오픈소스로 기여했다.
왜 읽나분산 환경에서 Kafka 컨슈머를 무중단으로 재처리해야 하는 백엔드 개발자에게 CLI/Admin API의 한계와 Spring Kafka seek + Redis Pub/Sub 전파 설계를 보여주는 사례.
마켓컬리
마켓컬리 Hello World 블로그
원문은 여기서 이어서 읽을 수 있어요
원문 읽기
읽음 (0)

이 글과 비슷한

  1. 백엔드·github-engGitHub Engineering·

    조기 종료를 없애야 벡터화된다 — 메모리 속도 소스 코드 케이스 폴딩

    GitHub의 코드 검색 엔진 Blackbird는 480TB 이상의 소스 코드를 인덱싱하기 전 모든 바이트에 case folding을 적용한다. 이 글은 Rust로 구현한 case folding을 메모리 대역폭 한계(45+ GiB/s)까지 끌어올린 두 가지 반직관적 최적화를 상세히 다룬다. 핵심은 루프 조기 종료(break) 제거로 LLVM 벡터화를 유도하고, UTF-8을 디코딩하지 않고 바이트 공간 산술만으로 fold를 수행하는 것이다.

    #rust#unicode#simd+2
  2. 백엔드·여기어때 (GC컴퍼니)여기어때 (GC컴퍼니)·

    트랜잭션 스크립트에서 숙소 메타 + 가격 계산 모듈로 — 전시 아키텍처 개선기 (2/3)

    여기어때 전시개발팀이 숙소 상세(PDP) API를 해부한 결과, 코드상으로는 DB 호출 3번처럼 보이던 요청이 실제로는 MongoDB $lookup 체인으로 컬렉션을 19회 접근하는 구조였다. 이 트랜잭션 스크립트 방식의 핵심 문제는 "aggregation이 I/O를 가린다"는 점으로, 독립적인 쿼리 10개가 단일 파이프라인에 직렬화되어 병렬화 기회를 잃고, 가격 때문에 거의 안 바뀌는 이미지까지 매 요청마다 읽어야 하는 읽기 증폭이 발생했다. V3에서는 "조회 시점 조립"을 "쓰기 시점 사전 조립"으로 전환하고, 화면별로 복제되던 가격 계산 로직을 goodsprice 단일 모듈로 수렴했다. 4개 API(PLP/PDP/RDP/ILP)의 반복 마이그레이션은 Claude Code skill로 절차를 고정하고 쉐도잉 + 동일성 검증으로 안전망을 마련하는 방식으로 진행됐다.

    #architecture#migration#caching+2