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

컬리 검색이 카프카를 들여다본 이야기 2

컬리 검색/추천팀이 서로 다른 토픽의 메시지를 같은 키로 조합해 하나의 색인 메시지로 만들기 위해 카프카 스트림즈를 도입한 경험을 다룬다. 레디스를 중간 저장소로 쓰던 초기 방식을 스트림즈의 KTable/KStream join으로 대체하고, 조용히 죽는 스트림즈 스레드를 스프링 카프카로 관리·모니터링한 과정을 정리한다.

핵심 포인트
  • 두 토픽 메시지를 레디스에 적재 후 병합하던 방식을 카프카 스트림즈 KTable-KStream join으로 단순화
  • 순수 자바로 StreamsBuilder를 직접 만들면 스트림즈 스레드가 죽어도 스프링 앱은 정상으로 표시되는 사각지대 발생
  • 스프링 카프카의 @EnableKafkaStreams로 StreamsBuilder를 @Autowired 주입해 생명주기를 스프링에 위임
  • KafkaStreams.state()를 읽는 HealthIndicator로 스트림즈 상태를 실시간 헬스체크
  • Error 상태는 재시작 외 복구가 불가라 Down으로 판정해 조기 감지
상세 정리
  • 요구사항: 1편의 레디스 기반 조합(각 토픽 메시지를 정제해 레디스에 적재하고 한 상품의 데이터가 다 모이면 읽어 병합·색인)을 더 단순한 방식으로 바꾸고 싶었다.
  • 스트림즈 선택: 카프카에 저장된 데이터를 처리하는 클라이언트 라이브러리로 진입장벽이 낮은 카프카 스트림즈를 발견했다.
  • 기본 파이프라인: sourceTopicA를 KTable로, sourceTopicB를 KStream으로 읽어 같은 키로 join한 뒤 null을 필터링해 MSG-MERGED 토픽으로 내보내는 토폴로지를 순수 자바로 구현했다.
  • 1차 스프링 연동: ApplicationStartedEvent 리스너에서 buildPipeline을 호출해 KafkaStreams를 직접 생성하고 start 하는 방식으로 붙였다.
  • 장애: 테스트 환경의 규격 어긋난 메시지가 처리 오류를 냈고, All stream threads have died 로그를 남기며 스트림즈 스레드가 죽었다.
  • 사각지대: 스레드가 죽어도 스프링 어플리케이션 상태는 정상으로 표시돼, 별도 헬스체크 없이는 장애를 인지할 수 없었다.
  • 가시성 확보: 스프링 카프카가 카프카 스트림즈를 지원한다는 점을 활용해 @EnableKafkaStreams 설정과 KafkaStreamsConfiguration 빈을 구성했다.
  • 구현 전환: KafkaStreams를 직접 만들지 않고 StreamsBuilder를 @Autowired로 주입해 스프링이 스트림즈 클라이언트 생성과 생명주기를 관리하게 바꿨고, num.stream.threads는 3으로 설정했다.
  • 헬스체크: Created·Running·Re-Balancing은 up, Pending Shutdown·Not Running·Error는 down으로 보고, StreamsBuilderFactoryBean에서 KafkaStreams를 꺼내 state()를 확인하는 HealthIndicator를 구현했다.
  • 결론: 스트림즈로 메시지 병합을 쉽게 처리하고, 관리와 상태 모니터링은 스프링에 일임할 수 있게 됐다.
왜 읽나스프링 환경에서 카프카 스트림즈를 운영하며 스레드 사망 가시성 문제로 고생한 백엔드 개발자에게 @EnableKafkaStreams와 HealthIndicator 패턴 레퍼런스.
마켓컬리
마켓컬리 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