Kafka 재시도가 정상 파티션까지 멈출 때
목차
벡터 저장이 실패했는데 오류 결과도 발행할 수 없었다. 입력에 결과 메시지의 key를 만들 필수 식별자가 없었기 때문이다. 결과 발행 실패를 모두 일시적인 통신 장애로 처리하자 같은 입력이 재시도 자리를 계속 차지했다. 다른 파티션의 정상 입력도 멈출 수 있었다.
9월에 소비 루프를 검토하면서 로컬 재현으로 확인한 문제다. 운영에서 같은 장애가 발생했다는 뜻은 아니다. 입력을 잃지 않으려 강화한 재시도가 진행 가능성을 해친다는 점이 중요한 발견이었다.
다시 시도하면 무엇이 달라지는가
브로커 연결이 끊겼다면 시간이 지나 연결이 복구될 수 있다. 필수 key가 없는 입력은 시간이 지나도 그대로다. 둘을 같은 실패로 묶으면 재시도 횟수나 지연을 조정해도 문제를 해결할 수 없다.
| 실패 종류 | 같은 입력으로 재시도하면 | 필요한 대응 |
|---|---|---|
| 일시적인 브로커 연결 실패 | 연결 복구 후 성공할 수 있다 | 전달을 재시도한다 |
| 결과 key를 만들 식별자 누락 | 같은 실패가 반복된다 | 작업 등록 전에 격리한다 |
이번 경로는 저장 완료나 저장 실패를 결과 이벤트로 알리고 그 전달을 확인한 뒤 입력을 완료한다. 그런데 저장 identity 검사가 실패한 입력에 오류 이벤트를 만들려고 해도 같은 식별자가 필요했다. 문제는 출력 계약을 만족하지 못하는 입력을 worker에 먼저 넣으면서 시작됐다.
수정은 작은 편이었다. 기존 결과 key 계산 함수를 작업 등록 전에 호출했다. 결과를 발행해야 하는 이벤트에 필요한 key가 없다면 worker budget을 쓰지 않고 해당 파티션을 되감아 멈춘다. 원본 offset은 커밋하지 않는다.
입력 수신
→ 처리할 이벤트인지 확인
→ 필요한 결과 key를 만들 수 있는지 확인
→ 불가능: 해당 파티션 미커밋·격리
→ 가능: worker 등록·처리·결과 전달 확인저장할 필요가 없어 무시하는 이벤트에는 결과 key를 강제하지 않았다. key가 있는 정상 형식의 입력이 저장 중 실패하면 기존 오류 결과를 발행한다. 오류 종류마다 진행 조건이 다르다.
완료 순서와 커밋 순서는 다르다
한 파티션에서 작업을 병렬로 처리하면 뒤의 작업이 먼저 끝날 수 있다. 다음 예시에서 A가 실패 중인데 C까지 커밋하면 재시작한 소비자는 A를 다시 받지 못할 수 있다.
수신 순서: A → B → C
완료 상태: 대기 완료 완료
커밋 가능: 아직 없음소비 위치는 poll하면서 앞으로 갈 수 있지만 복구 시작점인 committed position은 안전한 완료 범위까지만 전진해야 한다. Kafka의 offset 숫자는 중간에 비어 있을 수 있으므로 숫자가 연속인지보다 실제 수신한 기록의 앞부분이 모두 끝났는지를 봐야 한다. Apache Kafka 소비자 문서
아래는 실제로 받은 기록 중 앞에서부터 모두 완료된 구간(prefix)을 구하는 실행 예제다. 실제 소비자는 파티션 소유권, commit 응답과 재시도까지 관리한다.
def next_offset(delivered, done):
result = None
for offset in delivered:
if offset not in done:
break
result = offset + 1
return result
assert next_offset([10, 12, 15], {12, 15}) is None
assert next_offset([10, 12, 15], {10, 15}) == 11
assert next_offset([10, 12, 15], {10, 12, 15}) == 1611이라는 값은 offset 11의 레코드를 처리했다는 뜻이 아니다. 여기서는 offset 10까지 완료했다고 나타내는 복구 위치다. 사용한 Python 클라이언트의 commit(message=...)도 메시지 offset에 1을 더해 저장한다. Confluent Python API
앞의 A 때문에 B와 C가 완료 상태로 남는 현상은 head-of-line blocking의 한 형태다. 파티션 순서를 지키려면 같은 파티션의 이 제약을 없애기 어렵다. 이번 변경은 그 영향을 정상적인 다른 파티션까지 넓히지 않는 데 초점을 뒀다.
실행 슬롯과 보관 공간을 따로 센다
실행 중인 작업 수만 세면 완료했지만 커밋하지 못한 기록이 계속 쌓일 수 있다. 반대로 그 기록까지 하나의 토픽 공통 한도에 넣으면 느린 파티션 하나가 다른 파티션의 진입까지 막는다.
구현에서는 실행 중인 작업과 커밋을 기다리는 기록을 따로 셌다.
| 한도 | 적용 단위 | 포함하는 대상 |
|---|---|---|
| L | 토픽 | 진행 중인 작업 |
| 2L | 파티션 | 실행 중이거나 완료 후 커밋을 기다리는 기록 |
영구 입력 오류는 실행 슬롯에 등록하기 전에 격리한다.
이 한도는 새로운 비용도 드러낸다. 파티션별 상한을 쓰면 파티션 수 P에 따라 기록 수의 최대치가 대략 P × 2L로 늘 수 있다. 각 기록의 payload 크기도 별도다. 파티션별 격리가 프로세스 전체 메모리 상한이나 공유 executor의 공정성을 자동 보장하는 것은 아니다.
격리는 복구 완료가 아니다
깨진 입력을 보존하고 파티션을 멈추면 운영자가 원인을 조사할 수 있다. 동시에 그 파티션의 정상 후속 메시지도 기다린다. 유실 방지와 가용성 사이의 선택을 숨기면 안 된다.
DLQ로 보내고 진행하는 정책도 가능하지만, 이 구현에서는 자동 폐기하지 않았다. DLQ를 택한다면 원문 전달의 영속성, 실패 시 원본 커밋 금지, 재처리 identity와 중복 효과를 먼저 정해야 한다. 다른 저장소에 결과를 쓰는 소비자는 Kafka offset과 그 저장을 자동으로 한 트랜잭션에 묶을 수 없다. Kafka의 외부 저장소와 offset 설명
파티션의 seek나 pause가 실패했을 때도 계속 진행하지 않게 했다. 소비 루프를 멈추고 readiness를 실패시켜 미확정 입력을 건너뛰지 않도록 했다. 미커밋 입력은 재할당이나 재시작으로 소비가 정상 재개되면 다시 검사한다.
다만 이후 버전 검토에서는 재할당 중 일시 오류 뒤 polling이 재개되지 않고 readiness가 정상으로 남는 문제가 확인됐다. 이 복구 경로는 미해결이며, 여기서 확인한 poison 입력 격리와 구분한다.
무엇을 검증했나
회귀 테스트는 잘못된 입력과 같은 토픽의 정상 다른 파티션, 다른 토픽의 정상 입력을 함께 넣는다. 확인 대상은 세 가지다. 잘못된 입력의 offset이 보존되는지, 정상 파티션이 진행하는지, 유효한 key를 가진 저장 실패는 오류 이벤트로 완료되는지다.
소비 루프와 테스트 코드, 당시 로컬 검증 기록을 대조했다. 기존 기록에는 집중 재현과 전체 Python 테스트 통과가 남아 있다. 이번 글에서는 위 prefix 예제를 실행했으며 실제 Kafka·벡터 DB·Gateway를 연결한 장애 시험이나 배포 상태는 확인하지 않았다.
재시도 루프를 점검할 때는 실패 로그의 개수보다 같은 입력이 어떤 자원을 붙잡고 있는지 보면 좋다. 입력을 바꾸지 않고 백 번 실행해도 달라질 것이 없다면 통신 재시도와 분리할 후보다. 잘못된 입력 한 개를 넣고 다른 파티션의 정상 작업이 끝나는지 확인하는 테스트로 시작할 수 있다.