홈
CCDAK
벼락치기 요약
외워야 할 숫자 25개
시간이 없으면 이 다섯 개만 확실히 하세요.
linger.ms=5 ·
in-flight ≤ 5 ·
max.poll.interval.ms=300000 / max.poll.records=500 ·
message.max.bytes=1048588 ·
session.timeout.ms=45000 .
설정 기본값 총정리
설정값 관계도 — 프로듀서 · 브로커 · 토픽 · 컨슈머 설정이 서로 어떻게 얽히는지 한 장으로.
acks ↔ min.insync.replicas ↔ replication.factor,
max.request.size ↔ message.max.bytes ↔ fetch.max.bytes ↔
replica.fetch.max.bytes,
linger.ms ↔ batch.size ↔ delivery.timeout.ms 묶음이 보입니다.
시험 직전 한 장 요약으로 가장 유용한 그림입니다.
프로듀서
컨슈머
브로커 · 토픽
Streams · Connect
Application Development — 28%
프로듀서
send()는 비동기 즉시 반환 . 실제 전송은 Sender 스레드
onCompletion()은 Sender 스레드 에서 실행. 블로킹 작업 금지
metadata와 exception 중 하나는 항상 null
같은 파티션의 콜백은 순서 보장 . 파티션이 다르면 보장 없음
멱등성 = PID + 파티션별 시퀀스 번호 . 범위는 한 세션 × 한 파티션
멱등성 켜면 in-flight ≤ 5 . 6 이상이면 기동 시 ConfigException
키 있음 → murmur2(key) % N 결정적 / 키 없음 → 배치 단위 sticky
파티션 수를 늘리면 같은 키의 배치가 바뀌어 순서가 깨짐
실질 재시도 상한은 retries가 아니라 delivery.timeout.ms
linger.ms + request.timeout.ms ≤ delivery.timeout.ms
버퍼 고갈 시 send()가 블로킹 되고 max.block.ms 초과 시 TimeoutException
컨슈머
하트비트 축 = heartbeat.interval.ms · session.timeout.ms
poll 축 = max.poll.interval.ms · max.poll.records
"프로세스는 살아 있는데 그룹에서 빠짐" → poll 축 . 처방은 max.poll.records 축소
subscribe() = 그룹 참여·리밸런스·group.id 필수
assign() = 수동 지정·리밸런스 없음·group.id 없어도 동작·파티션 증가 미감지
둘을 섞으면 IllegalStateException
리밸런스 조건: 합류/이탈, 크래시, poll 초과, 파티션 수 증가 , 패턴 매칭 새 토픽, 코디네이터 장애
auto.offset.reset은 커밋이 없거나 범위 밖일 때만 발동
커밋 오프셋은 offsets.retention.minutes(7일) 후 만료 → latest면 구간 건너뜀
권장 커밋: 루프에서 commitAsync(), 종료 finally에서 commitSync()
처리 전 커밋 = at-most-once(유실), 처리 후 커밋 = at-least-once(중복)
group.protocol=consumer면 session·heartbeat·assignment.strategy가 무시
전달 보장 · 트랜잭션
기본 동작은 at-least-once
EOS = 프로듀서 transactional.id + 컨슈머 read_committed + sendOffsetsToTransaction
isolation.level 기본값 read_uncommitted — 컨슈머 설정
read_committed는 LSO까지만 읽음. 열린 트랜잭션이 뒤 메시지까지 막음
컨트롤 레코드는 애플리케이션에 안 보이지만 오프셋을 소비 → 오프셋 빈 칸 발생
EOS는 Kafka 경계 안에서만 . 외부 DB·API는 멱등 키를 직접 구현
트랜잭션 사용 시 컨슈머 enable.auto.commit=false
Fundamentals — 23%
복제 · 오프셋
순서 보장 단위는 파티션 . 전역 오프셋은 없음
파티션 수는 늘릴 수만 있고 늘리면 리밸런스 + 키 배치 변경
acks=all은 "현재 ISR 전원"의 응답. 모든 레플리카가 아님
쓰기 가능 장애 대수 = RF − minISR (RF=3, minISR=2 → 1대)
ISR 이탈 기준은 replica.lag.time.max.ms(30000)
unclean.leader.election.enable=false(기본) → 유실 방지, 가용성 희생
오프셋 4종: log-start / committed / high watermark / LEO
lag = high watermark − committed offset (LEO 아님)
컨슈머는 high watermark를 넘어 읽지 못함
메시지 크기 5+1개
상한은 압축 후 레코드 배치 기준
브로커 값(1048588)만 1MiB가 아님 . 12바이트 차이
fetch 3종(replica.fetch.max.bytes 등)은 하드 상한이 아님 — 빼먹어도 복제는 되고 처리량만 무너짐 . 하드 게이트는 max.request.size와 message.max.bytes 둘뿐
리텐션 · 컴팩션
delete = 기간·크기 초과분을 세그먼트 단위 로 삭제
compact = 키별 최신 값만. 키가 필수
활성 세그먼트는 컴팩션 대상이 아님 → 옛 값이 보이는 것이 정상
컴팩션은 즉시가 아니라 백그라운드 (dirty ratio 0.5 초과 시)
tombstone = value=null . 수명은 delete.retention.ms(1일)
retention.ms를 줄여도 segment.ms가 크면 안 지워짐
retention.bytes는 파티션당
디스크 = 유입량 × 일수 × RF × (1 − 압축률) × 여유
스키마 호환성
BACKWARD: 필드 추가에는 default 필수 . 삭제는 자유
FORWARD: 필드 추가 자유. default 없는 필드 삭제는 거부
타입 변경·이름 변경은 어느 모드에서도 위험 (alias 필요)
wire format = magic byte(1) + schema id(4) + payload
subject 전략: TopicName(기본) / RecordName / TopicRecordName
프로덕션에서 auto.register.schemas는 끄기
보안
PLAINTEXT(둘 다 없음) / SSL(암호화+선택 mTLS) /
SASL_PLAINTEXT(인증만, 평문 ) / SASL_SSL (권장)
SASL 메커니즘: PLAIN, SCRAM-SHA-256/512, GSSAPI, OAUTHBEARER
인증 ≠ 인가. ACL이 없으면 TopicAuthorizationException
컨슈머는 Topic Read + Group Read 두 ACL 필요
Kafka Connect — 15%
worker(프로세스) → connector(정의) → task(병렬성 단위)
standalone = properties로 커넥터 전달, 로컬 상태 / distributed = REST만 , Kafka 토픽에 상태
내부 토픽 3종 모두 cleanup.policy=compact
config.storage.topic은 파티션 1개 (순서 중요)
워커 group.id = Connect 클러스터 이름. 컨슈머 그룹과 겹치면 안 됨
source 오프셋 → offset.storage.topic
sink 오프셋 → __consumer_offsets , 그룹은 connect-{커넥터명}
컨버터는 key·value 각각 . JsonConverter는 schemas.enable
Avro·Protobuf 컨버터는 Confluent 제공 (Schema Registry 필요)
SMT 위치: source는 converter 앞 , sink는 converter 뒤
SMT는 메시지 1건 · stateless · 체인 순서가 결과를 바꿈 . 집계·조인 불가
DLQ 3조건: errors.tolerance=all + 토픽 이름 지정 + sink 전용
errors.deadletterqueue.context.headers.enable=true여야 원인이 남음
errors.retry.timeout 기본 0 = 재시도 없음 , -1은 무한
태스크가 FAILED가 되면 자동 재시작하지 않음
409 Conflict = 리밸런스 진행 중 또는 이름 중복
tasks.max는 상한. sink 실질 상한은 입력 파티션 수
확장 순서: tasks.max → 파티션 증설 → 워커 추가. 커넥터 복제는 중복 유발
REST 핵심 8개
Application Observability — 13%
client.id 를 지정하지 않으면 메트릭에서 구분 불가
프로듀서: record-error-rate(알림 1순위) · record-retry-rate ·
record-queue-time-avg · request-latency-avg · batch-size-avg
컨슈머 fetch: records-lag-max · records-lag-avg ·
records-consumed-rate · fetch-latency-avg
컨슈머 poll: time-between-poll-max ·
last-poll-seconds-ago · poll-idle-ratio-avg
컨슈머 코디네이터: assigned-partitions ·
rebalance-rate-per-hour · failed-rebalance-rate-per-hour ·
commit-latency-avg · last-rebalance-seconds-ago
리밸런스 예측은 time-between-poll-max 로. lag은 결과 지표라 늦음
assigned-partitions=0 → 컨슈머 수 > 파티션 수 또는 리밸런스 미완
lag 측정 3가지
브로커 메트릭 (개발자 관점)
UnderReplicatedPartitions = 0이어야 함 (경고 )
UnderMinIsrPartitionCount = 0이어야 함
(이미 acks=all 쓰기 실패 → NotEnoughReplicasException)
AtMinIsrPartitionCount — 한 대만 더 빠지면 실패
OfflinePartitionsCount = 0 (리더 없음 = 읽기·쓰기 불가)
ActiveControllerCount — 클러스터 전체 합이 1
RequestHandlerAvgIdlePercent · NetworkProcessorAvgIdlePercent — 낮으면 포화
IsrShrinksPerSec · IsrExpandsPerSec — 잦으면 불안정 (LeaderElectionRateAndTimeMs는 KRaft에서 제거됨)
Connect · Streams 관측
connector-failed-task-count > 0 → 즉시 조사
total-records-skipped · deadletterqueue-produce-requests —
errors.tolerance=all의 조용한 유실 을 잡는 지표
sink-record-lag-max는 sink에만 있음
Connect 로그는 태스크를 실행한 워커 노드의 로컬 파일 .
status의 worker_id로 찾음
Streams MBean: stream-thread-metrics · stream-task-metrics ·
stream-processor-node-metrics · stream-state-metrics
Kafka Streams — 12%
KStream = append-only 사실 / KTable = upsert 상태 (null = 삭제) /
GlobalKTable = 전체 복제
판단 기준: 같은 키가 반복될 때 두 값이 모두 의미가 있는가
리파티션 유발 : map · flatMap · selectKey ·
groupBy · repartition
유발하지 않음 : mapValues · flatMapValues ·
filter · peek · groupByKey
stateful: count · reduce · aggregate · 모든 윈도우 ·
모든 조인 · suppress
상태 저장소 = RocksDB + changelog 토픽(compact) .
리파티션 토픽은 delete
state.dir 기본이 임시 디렉터리 . num.standby.replicas는 0
태스크 수 = 입력 파티션 수 . 컨슈머 그룹 ID = application.id
processing.guarantee 기본 at_least_once. EOS는 브로커 3대 전제
조인 4종
co-partitioning 요건은 2개 : 파티션 수 동일 · 파티셔너 동일 . “키 동일”은 요건이 아니라 equi-join의 전제 입니다
파티션 수 불일치 → TopologyException
파티셔너 불일치는 감지되지 않아 조인이 조용히 누락
맞출 때는 적은 쪽을 많은 쪽에 . KStream-KTable이면 KStream을 리파티션
윈도우 4종
기본 시간 축은 event time . 추출기 기본은 FailOnInvalidTimestamp
…WithNoGrace면 지연 도착 레코드가 버려짐
중간 결과를 막고 최종만 내보내려면 suppress()
Application Testing — 8%
createInputTopic → serializer , createOutputTopic → deserializer
event-time punctuation은 자동 , wall-clock은 advanceWallClockTime()
상태 저장소는 getKeyValueStore()로 미리 채우거나 검증 가능
TopologyTestDriver로 리밸런스·다중 인스턴스·브로커 장애는 검증 불가
MockConsumer는 updateBeginningOffsets() 없이는 poll이 동작하지 않음
단일 브로커 테스트: transaction.state.log.replication.factor=1,
transaction.state.log.min.isr=1
--throughput -1 은 최대 처리량 측정용. 지연은 목표 값을 지정해 측정
지연은 99th · 99.9th 백분위 로 판단
Streams 테스트 재실행 전 kafka-streams-application-reset.sh 또는 cleanUp()
perf 도구의 --producer-props → --command-property ,
--messages → --num-records ,
--producer.config/--consumer.config → --command-config .
기존 옵션은 deprecated이며 Kafka 5.0에서 제거 예정입니다.
API · 워크플로 호출 순서
list order 유형 대비. 소리 내어 순서를 읊어 보세요.
트랜잭션 (consume-transform-produce)
1. initTransactions() ← 애플리케이션 시작 시 1회만
2. beginTransaction()
3. send(...) ← 결과 레코드
4. sendOffsetsToTransaction(...) ← 읽은 오프셋을 같은 트랜잭션에
5. commitTransaction()
실패 시 abortTransaction()
ProducerFenced / OutOfOrderSequence / Authorization → close() 후 종료
파티션 재할당
1. --generate (계획 JSON 생성)
2. --execute (--throttle 로 대역폭 제한)
3. --verify ← 반드시 실행. 스로틀이 여기서 해제됨
진행 중 취소는 --cancel
※ 파티션 수는 바꿀 수 없음 (레플리카 배치만 이동)
KRaft 부트스트랩
1. kafka-storage.sh random-uuid (클러스터 UUID 생성)
2. kafka-storage.sh format -t {UUID} -c server.properties
3. kafka-server-start.sh server.properties
4. kafka-topics.sh --create --bootstrap-server localhost:9092
※ --zookeeper 는 3.0 제거(KIP-604). kafka-configs.sh 만 3.9까지 남았고 4.x 엔 없음
큰 메시지 허용
1. 토픽 max.message.bytes
2. 프로듀서 max.request.size
3. 컨슈머 max.partition.fetch.bytes
4. 브로커 replica.fetch.max.bytes ← 하드 상한 아님. 빼먹으면 복제 처리량만 급감
스키마 진화 배포 (BACKWARD)
1. 새 스키마를 등록해 호환성 검사 통과 확인
2. 컨슈머 배포 ← BACKWARD 는 컨슈머 먼저
3. 프로듀서 배포
※ FORWARD 는 2와 3의 순서가 반대
Connect 실패 조사
1. GET /connectors/{name}/status → 실패 태스크와 trace, worker_id
2. 그 worker_id 노드의 로그 파일 확인 ← 리더 워커가 아님
3. PUT /admin/loggers/{logger} → 필요하면 런타임 로그 레벨 상향
4. 원인 수정 (컨버터·권한·스키마 등)
5. POST /connectors/{name}/restart?includeTasks=true&onlyFailed=true
6. 로그 레벨을 원래대로 되돌림
Streams 상태 초기화
1. 애플리케이션 정지 (모든 인스턴스)
2. kafka-streams-application-reset.sh --application-id {id}
3. 각 인스턴스에서 KafkaStreams.cleanUp() 또는 state.dir 삭제
4. 재기동
업그레이드 경로
2.x → 3.9 (마지막 ZooKeeper 지원 버전) → KRaft 마이그레이션 → 4.x
※ 4.3 업그레이드는 KRaft 필수, 소프트웨어·메타데이터 최소 3.3.x
버전 사실 10개
Java: 17 · 21 · 25 완전 지원 . 11은 clients/streams 등 일부. 8은 4.0에서 제거
배포판은 kafka_2.13-4.3.1.tgz 단일 . 2.13은 Scala 버전
Kafka 2.13이라는 버전은 존재하지 않습니다 (2.8 다음이 3.0)
시험장 들어가기 전 마지막 확인
문항 유형 3종 — multiple-choice · matching · list order. 객관식만이 아닙니다
linger.ms=5 , group.protocol=classic — 시중 자료가 가장 많이 틀리는 두 개
in-flight ≤ 5 , 초과하면 기동 실패
하트비트 축과 poll 축은 별개 . 처리 지연이면 max.poll.records
isolation.level은 컨슈머 설정 , 기본값은 read_uncommitted
EOS는 Kafka 경계 안에서만
lag = high watermark − committed offset
message.max.bytes=1048588 , max.request.size=1048576
BACKWARD=컨슈머 먼저 , FORWARD=프로듀서 먼저
sink 오프셋은 __consumer_offsets , source는 offset.storage.topic
DLQ는 errors.tolerance=all + 토픽 지정 + sink 전용
map은 리파티션, mapValues는 아님
TopologyTestDriver는 브로커 불필요
절대적 표현("항상", "모든")이 든 선택지는 대개 오답
설정을 볼 때 소속을 먼저 확인 — 가장 흔한 함정입니다
남은 시간에 새로운 것을 넣지 마세요.
함정 사전 의 "혼동 시 결과" 행만 한 번 더 훑고,
플래시카드 의 설정 기본값 덱을 10분 돌리는 것이 가장 효율적입니다.
인쇄해서 보기
브라우저의 인쇄 기능(Ctrl +P 또는 Cmd +P )으로 그대로 출력할 수 있습니다.
인쇄 시 사이드바 · 목차 패널 · 버튼 · 복사 아이콘은 자동으로 숨겨지고,
표와 코드 블록은 페이지 중간에서 잘리지 않도록 처리됩니다.
표가 넓어 잘리는 경우에는 가로 방향(landscape) 으로 인쇄하세요.
링크는 인쇄본에서 URL이 함께 표시되므로, 종이에서도 어느 페이지를 참조하는지 알 수 있습니다.
공식 문서 출처
이 페이지의 모든 설정 기본값과 버전 사실은 Apache Kafka 4.3 공식 문서 에서 확인했습니다.
도메인 가중치는 공개 자료 기준 잠정치이며 Confluent 공식 Exam Guide로 재확인이 필요합니다
(CCDAK 개요 참조).