1. 컨슈머와 컨슈머 그룹
컨슈머란
- 토픽을 구독(subscribe)하고, 파티션에서 메시지를 읽어 검증,저장,변환 등의 처리를 하는 클라이언트(소비자)
- 프로듀서가 쓰는 속도가 컨슈머가 처리하는 속도보다 빠르면 컨슈머가 뒤처진다(랙 증가)
- 이때 파티션과 컨슈머를 늘려 수평 확장 할 수 있다
컨슈머 그룹
- 같은 group.id를 가진 컨슈머들의 집합
- 그룹 내 컨슈머들이 토픽 파티션을 나눠서 읽는다
- 한 파티션은 그룹 내 하나의 컨슈머에만 할당됩니다
- 한 컨슈머는 여러 파티션을 읽을 수 있습니다
- 컨슈머 수 > 파티션 수이면 초과분은 아무것도 읽지 못한다
- 파티션 수가 그룹 병렬성의 상한
- 컨슈머 수 < 파티션 수 의 경우는 하나의 컨슈머가 한 개 이상의 파티션을 담당하게 된다
- 컨슈머는 (비교적 느릴 수 있는) 외부 작업(DB 쓰기, 무거운 계산 등)을 하는 경우가 많으므로, 처리 시간이 긴 컨슈머일수록 그룹 확장이 중요할 수 있다
- 서로 다른 애플리케이션은 서로 다른 group.id를 사용 한다
- 각 그룹이 독립적으로 토픽 전체 메시지를 읽으며, 한 그룹의 소비가 다른 그룹에 영향을 주지 않는다
- 그룹 A(2개 컨슈머), 그룹 B(4개 컨슈머)가 같은 토픽을 각각 읽음
- 각 그룹이 독립적으로 토픽 전체 메시지를 읽으며, 한 그룹의 소비가 다른 그룹에 영향을 주지 않는다
2. 컨슈머 그룹과 파티션 리밸런스
리밸런스(Rebalance): 컨슈머에서 다른 컨슈머로 파티션이 할당 되는 것
리밸런스가 일어나는 경우
- 컨슈머가 그룹에 추가될 때
- 기존 컨슈머가 읽던 파티션 일부를 새 컨슈머가 가져감
- 컨슈머가 종료 또는 크래시할 때
- 해당 컨슈머의 파티션이 남은 컨슈머에게 재할당됨
- 토픽에 파티션이 추가되는 등 토픽 메타데이터가 바뀔 때
리밸런스는 컨슈머 그룹의 고가용성과 확장성을 가능하게 하지만, 리벨런싱 중에는 지연이 발생할 수 있어 주의가 필요하다
조급한(Eager) 리밸런스
- 모든 컨슈머가 읽기를 멈추고 자신의 파티션 할당을 해제 한다
- 그룹에 다시 참여(rejoin)하고, 새로 파티션을 할당 한다
- 이 동안 전체 그룹이 읽지 못하는 stop-the-world가 발생 한다
- 컨슈머가 많을수록 이 중단 시간이 길어질 수 있다
- 재할당 결과가 이전과 같더라도 모든 파티션을 일단 내려놓고 다시 받는다
협력적(Cooperative, Incremental) 리밸런스
- 컨슈머 그룹 리더가 다른 컨슈머들에게 재할당이 필요한 파티션만 알려 준다
- 해당 파티션을 가진 컨슈머만 그 파티션의 읽기를 멈추고 할당을 해제 한다
- 할당 해제된 파티션은 새 컨슈머에게 할당 된다
- 점진적으로 진행되며, 나머지 파티션은 계속 읽기가 가능해 전체 중단이 없습니다
- 컨슈머가 많은 큰 그룹에서 유리하다
활성화 방법 쓸 것
리밸런스 동작 방식
- 컨슈머는 그룹 코디네이터(group coordinator) 역할을 하는 브로커에 하트비트(heartbeat) 를 보내 멤버십과 할당된 파티션 할당을 유지 한다
- 하트비트는 컨슈머의 백그라운드 스레드에서 전송되며, 일정 시간 하트비트가 없으면 코디네이터가 그 컨슈머를 죽은 것으로 간주하고 리밸런스 시작한다
- 컨슈머가 정상 종료(close())하면 코디네이터에 그룹 탈퇴를 알려, 세션 타임아웃을 기다리지 않고 즉시 리밸런스가 일어난다
- 파티션 할당 과정
- 컨슈머가 처음 그룹에 참여하려 할 때 코디네이터에 JoinGroup 요청
- 그룹에 처음 참여한 컨슈머가 그룹 리더가 된다
- 리더는 코디네이터로부터 그룹 멤버 목록을 받아 PartitionAssignor 구현체를 사용해 파티션 할당을 계산한다
- 리더가 할당 결과를 코디네이터에 전달하고, 코디네이터가 각 컨슈머에게 자기 할당 정보만 전달한다
- 각 컨슈머는 자기 할당만 알고, 전체 할당은 리더만 알고 있다- 이 과정은 리밸런스가 일어날 때마다 반복된다
정적 그룹 멤버십 (Static Group Membership)
- 기본적으로 컨슈머의 그룹 멤버십은 일시적이다 컨슈머가 그룹을 떠나면 파티션이 회수되고, 재참여하면 새 멤버 ID와 새 파티션이 할당된다
- group.instance.id를 지정하면 정적 멤버가 된다
- 컨슈머가 종료되어도 session.timeout.ms 이내에 같은 group.instance.id로 다시 참여하면 리밸런스 없이 이전 파티션을 그대로 돌려받는다
- 그룹 코디네이터가 그룹 상태를 캐시하므로, 다시 참여하는 컨슈머는 이전 할당 정보를 받는다
- session.timeout.ms를 넘기면 리밸런스가 일어나고 파티션은 다른 컨슈머에게 넘어간다
- 유용한 경우: 각 컨슈머가 로컬 상태나 캐시를 유지하는 애플리케이션(재구축 비용이 큰 경우)
- 주의점
- group.instance.id는 그룹 내에서 고유해야 한다 같은 ID로 두 컨슈머가 뜨면 나중에 시작한 쪽이 오류를 받는다
- 컨슈머가 죽었는데 재시작이 늦으면, 그동안 그 파티션은 아무도 읽지 않아 지연이 커진다
- 따라서 session.timeout.ms는 "재시작을 기다려줄 시간"과 "장애 감지 속도" 사이에서 균형을 맞춰 정한다
3. 컨슈머 생성과 토픽 구독
필수 속성
속성 설명
| bootstrap.servers | 카프카 클러스터 연결용 브로커 목록(host:port) 최소 2개 이상 지정 권장(하나가 죽어도 연결 가능) |
| key.deserializer | 키를 바이트 배열에서 객체로 변환하는 클래스 |
| value.deserializer | 값을 바이트 배열에서 객체로 변환하는 클래스 |
| group.id | 컨슈머 그룹 ID 엄밀히는 필수가 아니지만(그룹 없이 쓰는 방법도 있음) 거의 항상 지정 |
생성 예제
java
Properties props = new Properties();
props.put("bootstrap.servers", "broker1:9092,broker2:9092");
props.put("group.id", "CountryCounter");
props.put("key.deserializer",
"org.apache.kafka.common.serialization.StringDeserializer");
props.put("value.deserializer",
"org.apache.kafka.common.serialization.StringDeserializer");
KafkaConsumer<String, String> consumer = new KafkaConsumer<>(props);
토픽 구독
java
// 토픽 목록 구독
consumer.subscribe(Collections.singletonList("customerCountries"));
// 정규식 구독
consumer.subscribe(Pattern.compile("test.*"));
- 정규식 구독은 나중에 생성되는 토픽도 이름이 일치하면 자동으로 구독하고 리밸런스가 일어난다
- 여러 시스템의 데이터를 한 컨슈머가 읽을 때 유용하다
- 정규식 구독 시 컨슈머(또는 브로커 쪽)는 주기적으로 전체 토픽, 파티션 목록을 요청한다 토픽과 파티션이 매우 많으면 컨슈머 그룹이 많은 클러스터에서 브로커 네트워크 부하가 커질 수 있고, 클라이언트 쪽에서 목록을 필터링하는 비용도 있다
- 정규식 구독은 컨슈머에 토픽 목록 조회 권한(ACL의 describe 등)이 필요하다 특정 토픽에만 권한이 있는 환경에서는 정규식 구독이 실패할 수 있다
4. 폴링 루프
기본 구조
java
Duration timeout = Duration.ofMillis(100);
try {
while (true) {
ConsumerRecords<String, String> records = consumer.poll(timeout);
for (ConsumerRecord<String, String> record : records) {
log.debug("topic = {}, partition = {}, offset = {}, customer = {}, country = {}",
record.topic(), record.partition(), record.offset(),
record.key(), record.value());
// 처리 로직
}
}
} finally {
consumer.close();
}
poll()이 하는 일
- 컨슈머 API의 핵심이다 처음 호출 시 다음을 수행한다
- 그룹 코디네이터를 찾고, 컨슈머 그룹에 참여하고, 파티션을 할당받는다
- 이후 호출에서는 리밸런스, 하트비트 확인, 데이터 가져오기를 처리한다
- 반환 값은 ConsumerRecords로, 레코드마다 토픽, 파티션, 오프셋, 키, 값 정보를 담고 있다
- poll() 인자(timeout)는 컨슈머 버퍼에 데이터가 없을 때 블록할 최대 시간이다 0이면 즉시 반환하고, 데이터가 있으면 timeout과 관계없이 즉시 반환한다
- 하트비트는 별도 스레드가 보내지만, poll()을 오래 호출하지 않으면 max.poll.interval.ms 초과로 컨슈머가 그룹에서 제외된다
- 루프 안에서 오래 걸리는 처리(외부 호출, 대량 계산)를 하면 poll() 간격이 길어져 리밸런스가 발생할 수 있다 해결책
- max.poll.records 감소
- max.poll.interval.ms 증가
- 처리를 별도 스레드로 분리
스레드 안전성
- 하나의 스레드에 하나의 컨슈머가 원칙이다
- 같은 그룹의 컨슈머 여러 개를 한 스레드에서 돌리거나, 한 컨슈머를 여러 스레드가 공유하는 것은 지원하지 않는다
(KafkaConsumer는 스레드 안전하지 않으며, 동시에 접근하면 ConcurrentModificationException이 발생한다) - 같은 애플리케이션에서 같은 그룹의 컨슈머를 여러 개 돌리려면 각각 별도 스레드에서 실행한다
(ExecutorService에 컨슈머 하나씩 제출) - wakeup()은 다른 스레드에서 호출해도 안전한 유일한 메서드
- 이벤트 처리 병렬화가 필요하면 컨슈머 하나가 poll()로 데이터를 가져오고, 처리는 스레드 풀에 넘기는 패턴이 있다 이 경우 커밋 시점과 처리 완료 시점을 맞추는 데 주의해야 한다
- 컨슈머 그룹을 스레드/프로세스에 걸쳐 확장하는 방법이 필요하면 Kafka Streams 등 프레임워크 사용도 고려한다
5. 컨슈머 설정
데이터 가져오기 관련
설정 기본값 설명
| fetch.min.bytes | 1 | 브로커가 응답하기 전 모으는 최소 데이터량 늘리면 브로커, 컨슈머의 요청 처리 부하가 줄고 처리량이 올라가지만, 지연이 늘어날 수 있음 |
| fetch.max.wait.ms | 500 | fetch.min.bytes 조건을 기다리는 최대 시간 둘 중 먼저 만족되는 조건에서 응답 즉 지연의 상한이 됨 |
| fetch.max.bytes | 50MB | 컨슈머가 브로커에 요청할 때 브로커가 반환하는 최대 데이터 단, 첫 배치가 이보다 크면 그대로 반환(진행이 막히지 않도록) |
| max.partition.fetch.bytes | 1MB | 파티션당 서버가 반환하는 최대 바이트 컨슈머 메모리 산정 시 파티션 수 × 이 값을 고려 일반적으로 fetch.max.bytes를 쓰는 편이 관리하기 쉬움 |
| max.poll.records | 500 | poll() 한 번에 반환되는 최대 레코드 수 바이트가 아니라 레코드 수 제한 폴링 루프 한 번의 처리량을 조절 |
| receive.buffer.bytes / send.buffer.bytes | 64KB / 128KB | 데이터 읽기, 쓰기 시 TCP 소켓 버퍼 크기 -1이면 OS 기본값 사용 브로커와 컨슈머가 서로 다른 데이터센터에 있을 때 늘리면 유용 |
생존(그룹 멤버십) 관련
설정 기본값 설명
| session.timeout.ms | 45초 | 하트비트 없이 컨슈머가 살아있다고 볼 최대 시간 초과 시 코디네이터가 죽은 것으로 판단 (예전 기본값은 10초) |
| heartbeat.interval.ms | 3초 | 하트비트 전송 주기 보통 session.timeout.ms의 1/3로 설정 둘은 함께 조정 |
| max.poll.interval.ms | 5분 | poll() 호출 간 허용 최대 시간 하트비트는 백그라운드 스레드가 보내므로 컨슈머 메인 스레드가 멈춘 것을 이 값으로 감지 초과하면 컨슈머가 그룹에서 이탈 |
| default.api.timeout.ms | 1분 | timeout 파라미터를 명시하지 않은 컨슈머 API 호출의 타임아웃(단, poll() 제외) 예: commitSync() 등 |
| request.timeout.ms | 30초 | 컨슈머가 브로커 응답을 기다리는 최대 시간 초과하면 연결을 닫고 재연결 시도 가급적 기본값 유지 권장 |
- session.timeout.ms를 낮추면: 죽은 컨슈머를 빨리 감지하지만, 처리 지연이나 GC pause 등으로 불필요한 리밸런스가 늘 수 있다
- session.timeout.ms를 높이면: 불필요한 리밸런스는 줄지만, 실제로 죽은 컨슈머 감지가 느려진다
그 외 설정
설정 기본값 설명
| auto.offset.reset | latest | 커밋된 오프셋이 없거나 유효하지 않을 때(오프셋이 삭제된 경우 포함) 동작 latest: 가장 최신 레코드부터, earliest: 처음부터, none: 예외 발생 |
| enable.auto.commit | true | 오프셋 자동 커밋 여부 중복/유실을 직접 제어하려면 false |
| auto.commit.interval.ms | 5초 | 자동 커밋 주기 |
| partition.assignment.strategy | Range + CooperativeSticky | 파티션 할당 전략 (아래 참고) |
| client.id | 임의 문자열 | 브로커가 요청 출처를 식별하는 값 로깅, 메트릭, 쿼터에 사용 |
| client.rack | (없음) | 컨슈머의 위치(AZ 등) 같은 랙의 팔로워 레플리카에서 읽도록 해 네트워크 비용, 지연 감소 브로커의 replica.selector.class 설정과 함께 사용 |
| group.instance.id | (없음) | 정적 그룹 멤버십용 고유 ID |
| offsets.retention.minutes | 7일 (브로커 설정) | 컨슈머 그룹이 활동하지 않을 때(빈 그룹) 커밋 오프셋을 보관하는 기간 이보다 오래 멈춘 뒤 재시작하면 오프셋이 사라져 auto.offset.reset 동작을 따름 |
파티션 할당 전략 (partition.assignment.strategy)
전략 동작
| Range | 컨슈머가 구독하는 각 토픽에 대해, 토픽 파티션을 연속된 범위로 나눠 할당 컨슈머 수로 나누어떨어지지 않으면 앞쪽 컨슈머가 파티션을 더 받음 여러 토픽을 구독하면 앞쪽 컨슈머에 편중되는 문제 |
| RoundRobin | 모든 구독 토픽의 모든 파티션을 컨슈머에 순환하며 균등하게 배분 컨슈머 간 파티션 수 차이가 최대 1 |
| Sticky | RoundRobin처럼 균형을 맞추되, 리밸런스 시 기존 할당을 최대한 유지하여 이동을 최소화 조급한 리밸런스 프로토콜 사용 |
| CooperativeSticky | Sticky와 같지만 협력적 리밸런스 프로토콜 지원 컨슈머가 처리 중이던 파티션을 리밸런스 중에도 계속 읽을 수 있음 |
- 커스텀 할당자는 PartitionAssignor 인터페이스를 구현하여 지정할 수 있다
6. 커밋과 오프셋
개념
- 카프카는 다른 JMS 계열 큐와 달리 컨슈머의 ack를 받지 않는다 대신 컨슈머가 파티션의 현재 읽은 위치(오프셋)를 기록한다
- 오프셋을 업데이트하는 동작을 커밋(commit)이라 하며, 컨슈머는 오프셋을 내부 토픽 __consumer_offsets에 메시지로 보낸다(파티션별로 커밋된 오프셋)
- 컨슈머가 크래시하거나 그룹에 새 컨슈머가 들어와 리밸런스가 일어나면, 각 컨슈머는 마지막으로 커밋된 오프셋부터 읽기를 이어간다
- 커밋 오프셋과 실제 처리 위치가 다르면 문제가 생긴다
- 커밋된 오프셋이 클라이언트가 실제 처리한 마지막 오프셋보다 작으면 → 그 사이 메시지가 중복 처리된다
- 커밋된 오프셋이 클라이언트가 실제 처리한 마지막 오프셋보다 크면 → 그 사이 메시지가 유실된다(처리되지 않은 채 건너뜀)
- 따라서 오프셋 관리 방식이 애플리케이션에 큰 영향을 준다
자동 커밋
- enable.auto.commit=true이면 auto.commit.interval.ms(기본 5초)마다 컨슈머가 poll()이 반환한 가장 마지막 오프셋을 커밋한다
- 커밋은 poll() 호출 안에서, 마지막 커밋 이후 interval이 지났는지 확인해 수행한다
- 장점: 간단하다
- 단점: 커밋 후 3초 뒤 크래시가 나면 리밸런스 후 3초 분량 메시지가 중복 처리된다 커밋 간격을 줄여도 중복 윈도만 줄일 뿐 없앨 수는 없다
- 주의: 다음 poll() 전에 반환된 모든 레코드 처리를 끝내야 한다(다음 poll 시점에 커밋되므로, 처리 중 예외가 나면 안전성이 깨질 수 있다) close() 시에도 자동 커밋된다
현재 오프셋 수동 커밋: commitSync()
java
while (true) {
ConsumerRecords<String, String> records = consumer.poll(timeout);
for (ConsumerRecord<String, String> record : records) {
// 처리
}
try {
consumer.commitSync();
} catch (CommitFailedException e) {
log.error("commit failed", e);
}
}
- enable.auto.commit=false로 설정한다
- 가장 최근 poll()이 반환한 마지막 오프셋을 커밋한다 따라서 poll()이 반환한 모든 레코드를 처리한 후 호출해야 한다(그렇지 않으면 유실 가능)
- 복구 불가능한 오류가 아니면 성공할 때까지 재시도한다 복구 불가능하면 예외를 던진다(CommitFailedException 등)
- 브로커 응답이 올 때까지 블록되므로 처리량이 떨어진다 커밋 빈도를 줄이면(배치 여러 개마다) 완화되지만 중복 가능성이 늘어난다
비동기 커밋: commitAsync()
java
consumer.commitAsync();
// 콜백 사용
consumer.commitAsync(new OffsetCommitCallback() {
public void onComplete(Map<TopicPartition, OffsetAndMetadata> offsets, Exception e) {
if (e != null)
log.error("Commit failed for offsets {}", offsets, e);
}
});
- 브로커 응답을 기다리지 않고 요청만 보내므로 처리량이 높다
- 재시도하지 않는다 재시도 시 순서 문제가 생기기 때문이다
- 예: 오프셋 2000 커밋이 일시 실패 → 그동안 3000 커밋이 성공 → 2000이 재시도로 성공하면 커밋 오프셋이 뒤로 돌아가 중복 처리를 유발한다
- 재시도를 직접 구현하려면 단조 증가하는 시퀀스 번호를 두고, 콜백에서 재시도 시점에 시퀀스가 최신일 때만 재시도한다(이미 더 나중 커밋이 있다면 재시도하지 않는다) 다만 이 패턴은 복잡하므로 아래 조합 방식을 더 자주 쓴다
- 콜백은 다음 poll()이나 commitSync() 호출 중에 실행된다
동기 + 비동기 조합 (권장 패턴)
java
try {
while (!closing) {
ConsumerRecords<String, String> records = consumer.poll(timeout);
for (ConsumerRecord<String, String> record : records) {
// 처리
}
consumer.commitAsync(); // 평상시: 빠르게
}
consumer.commitSync(); // 종료 직전: 확실하게
} catch (Exception e) {
log.error("Unexpected error", e);
} finally {
consumer.close();
}
- 평상시의 일시적 커밋 실패는 다음 커밋이 대체하므로 무시해도 된다
- 하지만 마지막 커밋(종료, 리밸런스 직전)은 다음 커밋이 없으므로 commitSync()로 성공을 보장한다
특정 오프셋 커밋
- commitSync()/commitAsync()는 기본적으로 마지막 poll이 반환한 최신 오프셋을 커밋하므로, 큰 배치 처리 중간에 커밋하고 싶을 때는 오프셋 맵을 직접 전달한다
java
private Map<TopicPartition, OffsetAndMetadata> currentOffsets = new HashMap<>();
int count = 0;
// ... 루프 안에서 레코드 처리 후
currentOffsets.put(
new TopicPartition(record.topic(), record.partition()),
new OffsetAndMetadata(record.offset() + 1, "no metadata")); // 다음에 읽을 오프셋
if (count % 1000 == 0)
consumer.commitAsync(currentOffsets, null);
count++;
- 커밋하는 값은 마지막으로 처리한 오프셋 + 1(= 다음에 읽을 오프셋)이다
- 컨슈머가 여러 파티션을 읽는 경우 파티션별로 추적해야 한다
7. 리밸런스 리스너
컨슈머 API는 subscribe() 호출 시 ConsumerRebalanceListener를 전달해 리밸런스 전후에 코드를 실행할 수 있게 한다
메서드 호출 시점 용도
| onPartitionsAssigned(partitions) | 파티션이 컨슈머에 재할당된 후, 컨슈머가 읽기를 시작하기 전 | 파티션별 상태 로드, 오프셋 조회/seek() |
| onPartitionsRevoked(partitions) | 컨슈머가 읽기를 멈춘 후, 리밸런스 시작 전 | 오프셋 커밋, 리소스 정리 조급한 리밸런스에선 모든 파티션, 협력적 리밸런스에선 회수되는 파티션만 대상 |
| onPartitionsLost(partitions) | 협력적 리밸런스에서, 파티션이 회수 절차 없이 이미 다른 컨슈머에게 할당된 경우(예외적 상황) | 상태, 리소스 정리만 이미 다른 컨슈머가 소유하므로 오프셋 커밋은 하지 않음 기본 구현은 onPartitionsRevoked() 호출 |
예제: 리밸런스 직전에 처리한 오프셋 커밋
java
private Map<TopicPartition, OffsetAndMetadata> currentOffsets = new HashMap<>();
private class HandleRebalance implements ConsumerRebalanceListener {
public void onPartitionsAssigned(Collection<TopicPartition> partitions) { }
public void onPartitionsRevoked(Collection<TopicPartition> partitions) {
System.out.println("Lost partitions in rebalance. Committing current offsets: " + currentOffsets);
consumer.commitSync(currentOffsets); // 여기서는 동기 커밋
}
}
consumer.subscribe(topics, new HandleRebalance());
- 리밸런스 중에는 commitSync()로 확실하게 커밋해야 한다(비동기는 리밸런스가 끝나고 나서야 처리될 수 있다)
- 리스너 콜백은 컨슈머 스레드에서 poll() 도중에 호출되므로, 처리가 길어지면 다시 max.poll.interval.ms 문제가 생길 수 있다
8. 특정 오프셋의 레코드 읽기
poll()은 마지막 커밋 오프셋부터 읽지만, 다른 위치에서 읽고 싶을 때 다음 메서드를 사용한다
- seekToBeginning(partitions): 파티션의 맨 처음부터
- seekToEnd(partitions): 파티션의 맨 끝(최신)부터
- seek(partition, offset): 특정 오프셋부터
- offsetsForTimes(Map<TopicPartition, Long>): 특정 타임스탬프에 해당하는 오프셋을 조회 예를 들어 "한 시간 전부터 다시 처리"할 때 사용한다
java
Long oneHourEarlier = Instant.now().atZone(ZoneId.systemDefault()).minusHours(1).toEpochSecond();
Map<TopicPartition, Long> partitionTimestampMap = consumer.assignment()
.stream()
.collect(Collectors.toMap(tp -> tp, tp -> oneHourEarlier));
Map<TopicPartition, OffsetAndTimestamp> offsetMap = consumer.offsetsForTimes(partitionTimestampMap);
for (Map.Entry<TopicPartition, OffsetAndTimestamp> entry : offsetMap.entrySet()) {
consumer.seek(entry.getKey(), entry.getValue().offset());
}
정확히 한 번(exactly-once) 처리
- 카프카 커밋 대신 결과와 오프셋을 같은 트랜잭션으로 외부 저장소(DB)에 저장하면 중복 없이 한 번만 처리하는 효과를 낼 수 있다
java
public class SaveOffsetsOnRebalance implements ConsumerRebalanceListener {
public void onPartitionsRevoked(Collection<TopicPartition> partitions) {
commitDBTransaction(); // 결과+오프셋을 함께 커밋
}
public void onPartitionsAssigned(Collection<TopicPartition> partitions) {
for (TopicPartition partition : partitions)
consumer.seek(partition, getOffsetFromDB(partition)); // DB의 저장 위치로 이동
}
}
consumer.subscribe(topics, new SaveOffsetsOnRebalance());
consumer.poll(Duration.ofMillis(0)); // 파티션 할당 유도
for (TopicPartition partition : consumer.assignment())
consumer.seek(partition, getOffsetFromDB(partition));
while (true) {
ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(100));
for (ConsumerRecord<String, String> record : records) {
processRecord(record);
storeRecordInDB(record);
storeOffsetInDB(record.topic(), record.partition(), record.offset());
}
commitDBTransaction();
}
- 재시작과 리밸런스 후에는 항상 DB에 저장된 오프셋으로 seek()해야 한다
9. 안전하게 종료하기
- 폴링 루프는 무한 루프이므로, 종료하려면 다른 스레드(예: shutdown hook)에서 consumer.wakeup()을 호출한다
- wakeup()은 다른 스레드에서 호출해도 안전한 유일한 컨슈머 메서드다
- 컨슈머가 poll()에서 대기 중이면 즉시 WakeupException을 던지며 빠져나온다 처리 중이었다면 다음 poll() 호출 시 예외가 발생한다
- 메인 스레드에서 WakeupException을 잡아 루프를 빠져나오며, 예외 자체는 별도 처리가 필요 없다
- 종료 전에 반드시 consumer.close()를 호출한다
- 필요하면 오프셋을 커밋하고,
- 그룹 코디네이터에 컨슈머 탈퇴를 알려 세션 타임아웃을 기다리지 않고 즉시 리밸런스가 일어나게 한다
java
Runtime.getRuntime().addShutdownHook(new Thread() {
public void run() {
System.out.println("Starting exit...");
consumer.wakeup();
try {
mainThread.join(); // 메인 스레드가 정리를 마칠 때까지 대기
} catch (InterruptedException e) {
e.printStackTrace();
}
}
});
try {
// ... 컨슈머 구독 및 폴링 루프
while (true) {
ConsumerRecords<String, String> records = consumer.poll(timeout);
// 처리 및 commitAsync()
}
} catch (WakeupException e) {
// 무시: 종료 신호
} finally {
consumer.commitSync();
consumer.close();
System.out.println("Closed consumer and we are done");
}
10. 역직렬화(Deserializer)
- 프로듀서의 시리얼라이저와 반드시 짝이 맞아야 한다
- 종류: StringDeserializer, IntegerDeserializer, ByteArrayDeserializer 등
Custom Deserializer
- Deserializer<T> 인터페이스(configure(), deserialize(), close())를 구현한다
- 단점
- 프로듀서와 컨슈머의 직렬화 형식이 강하게 결합된다 필드가 추가, 변경되면 두 쪽 코드를 동시에 바꾸고 호환성을 신경 써야 한다(이전 데이터와 새 데이터가 한 토픽에 섞여 있을 때 특히 어렵다)
- 팀 간에 같은 코드를 공유해야 하고 유지보수가 어렵다
- 따라서 일반적으로 JSON, Thrift, Protobuf, Avro 같은 범용 포맷과 직렬화 프레임워크를 쓴다
Avro Deserializer + 스키마 레지스트리
- Avro는 스키마 변경에 유연하다 스키마를 별도 저장소(스키마 레지스트리)에 두고, 메시지에는 스키마 ID만 담는다 컨슈머는 그 ID로 스키마를 조회해 역직렬화한다
- 스키마가 변경 되어도(하위/상위 호환 규칙 내에서) 컨슈머가 이전 데이터와 새 데이터를 모두 읽을 수 있다
java
Properties props = new Properties();
props.put("bootstrap.servers", "broker1:9092,broker2:9092");
props.put("group.id", "CountryCounter");
props.put("key.deserializer", "io.confluent.kafka.serializers.KafkaAvroDeserializer");
props.put("value.deserializer", "io.confluent.kafka.serializers.KafkaAvroDeserializer");
props.put("specific.avro.reader", "true");
props.put("schema.registry.url", schemaUrl);
String topic = "customerContacts";
KafkaConsumer<String, Customer> consumer = new KafkaConsumer<>(props);
consumer.subscribe(Collections.singletonList(topic));
while (true) {
ConsumerRecords<String, Customer> records = consumer.poll(Duration.ofMillis(1000));
for (ConsumerRecord<String, Customer> record : records) {
System.out.println("current customer name is: " + record.value().getName());
}
consumer.commitSync();
}
- schema.registry.url: 스키마 레지스트리 주소
- specific.avro.reader=true: 생성된 특정 클래스(Customer)로 역직렬화한다 지정하지 않으면 GenericRecord로 반환된다
- 스키마 레지스트리는 Avro 외에 Protobuf, JSON Schema도 지원한다
11. Standalone 컨슈머
컨슈머 그룹과 리밸런스 없이 특정 파티션만 직접 읽고 싶은 경우다
- 예: 특정 토픽의 특정 파티션을 읽는 프로세스가 항상 하나만 있으면 될 때
java
List<PartitionInfo> partitionInfos = consumer.partitionsFor("topic");
if (partitionInfos != null) {
List<TopicPartition> partitions = new ArrayList<>();
for (PartitionInfo partition : partitionInfos)
partitions.add(new TopicPartition(partition.topic(), partition.partition()));
consumer.assign(partitions);
while (true) {
ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(1000));
for (ConsumerRecord<String, String> record : records) {
// 처리
}
consumer.commitSync();
}
}
- subscribe() 대신 assign()으로 파티션을 직접 지정한다 그룹 관리, 리밸런스가 없다
- subscribe()와 assign()은 함께 사용할 수 없다
- 오프셋 커밋은 여전히 가능하며, 커밋할 때는 group.id를 지정해 두는 것이 일반적이다
- 새 파티션이 추가되어도 자동으로 감지되지 않는다 주기적으로 partitionsFor()를 호출하거나, 파티션이 추가되면 컨슈머를 다시 시작해야 한다
'공부 > Kafka' 카테고리의 다른 글
| [카프카] Kafka Producer (1) | 2026.09.13 |
|---|---|
| [Kafka] Kafka CLI (0) | 2025.08.29 |
| [Kafka] Topic, Partition, Producer, Consumer (0) | 2025.08.29 |
| [Kafka] Zookeeper와 Broker (0) | 2025.06.29 |
| [Kafka] Kafka의 특징 (0) | 2021.06.13 |





