Spring Kafka — 무한 재전송 (GH-4465)
mergedsuspend @KafkaListener + DefaultErrorHandler 조합에서 비동기 retry가 무한 재전송되는 회귀(#4254)의 원인을 추적해 큐 중복을 제거했습니다. 버그를 재현하는 회귀 테스트를 함께 보냈습니다.
오픈소스 엔지니어 · OSS Contributor
Backend Engineer specializing in Distributed Systems, Spring, Kafka, and Event-Driven Architecture. OSS Contributor to Spring Kafka.
안녕하세요. 김빌(Bill Kim)입니다.
이 페이지는 제가 진행하고 있는 오픈소스(OSS) 활동과 프로젝트를 한곳에 정리하기 위해 만들었습니다.
주로 Apache Kafka와 Spring 생태계(Spring Kafka, Spring Modulith, Spring AI), Infinispan 등에 업스트림 기여를 하고 있으며, Kotlin과 Spring Boot 기반의 오픈소스 라이브러리를 직접 설계하고 운영하고 있습니다.
관심 분야는 분산 시스템(Distributed Systems), 분산 트랜잭션, 멱등성(Idempotency), 이벤트 기반 아키텍처(Event-Driven Architecture), 그리고 옵저버빌리티(Observability)입니다.
오픈소스 활동과 개발 과정에서 얻은 경험, 설계 고민, 기술적인 인사이트는 기술 블로그 DevBillLab에 기록하고 있습니다.
실제로 패치를 보낸 업스트림 PR입니다. 모두 GitHub @BK202503 계정에서 확인할 수 있습니다.
suspend @KafkaListener + DefaultErrorHandler 조합에서 비동기 retry가 무한 재전송되는 회귀(#4254)의 원인을 추적해 큐 중복을 제거했습니다. 버그를 재현하는 회귀 테스트를 함께 보냈습니다.
같은 파티션의 두 레코드가 같이 비동기 실패하면 earlier-offset 레코드가 한 번만 invoke되고 recoverer에 도달하지 못해 silent 유실되는 회귀. handleAsyncFailure가 같은 파티션 실패들을 도착 순서로 처리하면서 뒤 호출의 seek가 앞 호출의 seek를 덮어쓰던 문제를, offset 순 정렬 + partition-in-retry skip으로 수정. 후속으로 head-of-line amplification(burst quadratic)을 linear로 bound하는 PR과 4.0.x backport 회귀 fix PR도 진행 중.
@RetryableTopic 파이프라인에서 DLT에 안착한 레코드를 framework가 자동으로 retry chain entry로 다시 흘려보내는 first-class API 제안. eligibility policy(DltReplayPolicy)로 자동 재처리 대상 결정, cap 초과한 레코드는 parking queue에 종착시켜 human-only 검토 대상으로 격리. 모두 opt-in. 2022년 #2172가 'DIY로 하세요'로 닫혔던 use case를 framework 차원 contract로 lift.
Strimzi 조직에 보낸 네 건의 정리 — strimzi-kafka-operator 세 건과 kafka-access-operator 한 건. (1) GH-12847 — KafkaCluster만 IOException → RuntimeException으로 던지고, CruiseControl·ClusterCa는 warning 로그 후 빈 secret data / null CertAndKey로 진행하던 비대칭을 통일 (KafkaExporter 등 호출자 NPE 차단). (2) GH-12442 — Cluster Operator의 Cruise Control 클라이언트가 CC 서버 공개키를 핀하던 트러스트를 Cluster CA 기반으로 전환해 CA 롤 + 서버 cert 회전 사이 신뢰 불일치 제거. (3) GH-12870 — MirrorMaker 2 Grafana 대시보드 replication-consumer lag 패널이 Strimzi Metrics Reporter 환경에서 라벨 표기(clientid → client_id) 미스매치로 빈 그래프를 그리던 것을 SMR 컨벤션으로 정정. (4) kafka-access-operator GH-124 — KafkaListener가 ssl.truststore.crt에 Cluster CA Secret의 ca.crt만 담아 CA 키 교체 중 이전/신규 CA 중 한쪽만 신뢰하던 문제를 모든 .crt 엔트리 번들로 수정 (chain 항목은 trust anchor만).
Kafka Connect 전반에 퍼져 있던 ByteBuffer.array() 오용을 추적해 일관된 수정안을 제출했습니다. 직접 버퍼/슬라이스 버퍼 둘 다에서 안전하게 동작하도록 4개 컴포넌트에 같은 계열의 패치를 분리해 올렸습니다.
PR마다 PMD Copy/Paste Detector를 돌리는 GitHub Actions 워크플로를 추가했습니다. 기존 spotbugs.yml과 트리거/구조를 맞추고 CPD 결과를 아티팩트로 업로드합니다.
Spring Cloud AWS 3.x SQS 컨테이너에 IGNORE QueueNotFoundStrategy를 추가해, 큐가 없을 때 listener를 조용히 skip하고 application context 부팅을 차단하지 않도록 복원했습니다. 2.x spring-cloud-starter-aws-messaging의 기본 동작을 typed exception을 도입해 FAIL / CREATE / RetryableTopic 경로를 손대지 않고 분리해 적용했습니다.
contrib/btree_gist의 float4/float8 opclass가 raw C ==/> 연산자를 그대로 써서 IEEE 754 NaN을 잘못 처리. EXCLUDE 제약이 NaN 중복을 허용, RLS 정책이 NaN 행을 누수, GiST index scan이 같은 쿼리에 대해 seq scan과 다른 결과를 반환. utils/float.h의 NaN-aware float{4,8}_* / float{4,8}_cmp_internal()로 교체해 regular btree opclass와 같은 total order로 맞춤. core 245/245 + contrib 48 modules 전체 make check 통과. PostgreSQL은 GitHub PR을 안 받아서 pgsql-bugs 메일링리스트로 클레임 발송, upstream에 merged.
JpaEventPublicationAdapter#getStatus()가 영속 상태 대신 메모리 값을 돌려주던 버그를 수정했습니다. JdbcEventPublicationRepositoryV2가 이미 적용 중인 패턴과 동일하게 맞추고 markFailed 라운드트립 회귀 테스트를 추가했습니다.
스트리밍 응답에서 observation의 stop order가 잘못된 순서로 호출되는 이슈(#5971)에 대해 실패하는 회귀 테스트를 먼저 PR로 보냈습니다.
Lombok @Builder 스타일 fluent setter의 속성명이 JavaBean accessor와 어긋나, 두 번째 글자가 대문자인 필드의 매핑이 누락되던 #4000 버그를 DefaultAccessorNamingStrategy에서 수정했습니다. 회귀 테스트를 함께 보냈습니다.
Apache Kafka와 Spring Kafka의 구조, 메시지 흐름, 주요 소스 / 명령어를 한 페이지에 정리한 참고용 자료입니다.
브로커 / 토픽 / 파티션 / 컨슈머 그룹 아키텍처, produce → replicate → consume 메시지 흐름(acks=all + HW), kafka-topics / kafka-console-* / kafka-consumer-groups 주요 CLI 명령어 레퍼런스.
ContainerFactory → ConcurrentMessageListenerContainer → KafkaMessageListenerContainer / ListenerConsumer 컨테이너 계층, @KafkaListener 디스패치 + DefaultErrorHandler / SeekUtils / DeadLetterPublishingRecoverer 에러 처리 흐름, 주요 feature와 핵심 소스 파일 인덱스.
Cluster / Topic / User Operator 세 프로세스, Kafka·KafkaNodePool CR을 받아 StrimziPodSet · Service · Secret 로 reconcile하는 루프, 그리고 Kafka / KafkaNodePool / KafkaTopic / KafkaUser / KafkaConnect / KafkaMirrorMaker2 / KafkaBridge / KafkaRebalance CRD landscape 정리.
Apache Kafka 위에 Confluent가 얹어 파는 상용 컴포넌트(Schema Registry, ksqlDB, Cluster Linking, Tiered Storage, Confluent Cloud)의 라이선스 구분과 공개 API/시맨틱 정리. 소스가 공개되지 않은 부분은 다루지 않습니다.
Confluent Platform vs Confluent Cloud, 컴포넌트별 라이선스 매트릭스(Apache 2.0 / Confluent Community Licence / proprietary), Schema Registry의 subject·compatibility·wire format, ksqlDB의 push/pull 쿼리, Cluster Linking의 offset 보존 미러, Tiered Storage 핫/콜드 분리, Cloud의 cluster tier·eCKU·네트워크 옵션 정리.
Claude Code 워크플로우를 개선하는 오픈소스 도구.
Claude Code용 OSS 컨트리뷰션 워크플로우 툴킷. 업스트림 패치 발굴 → 제출 → 상태 추적을 위한 skill / agent / hook 세트.
Kotlin/Spring Boot 생태계에서 직접 설계하고 메인테이닝하는 오픈소스 라이브러리입니다.
Spring Boot용 Stripe 스타일 @Idempotent 어노테이션. coroutine-native idempotency-key 처리, JDBC/Redis 플러그러블 스토리지.
Spring Boot용 coroutine-native Saga 오케스트레이터. Kafka 이벤트 발행, JDBC 영속, 보상 트랜잭션, JVM 크래시 후 재개.
Spring Boot용 Transactional Outbox. Kotlin-first, coroutine-native, autoconfigured. bk-spring-saga와 짝으로 동작.
GitHub에 공개해 둔 프로젝트입니다.
Naver Local Search API 기반 카페 메타데이터 import 파이프라인.
오픈소스 협업, 코드 리뷰, 라이브러리 도입 문의 환영합니다.