CCAAK · 섹션 3
Kafka Connect
개발자 시험은 커넥터 설정을 묻고, 운영자 시험은 커넥터 클러스터를 묻습니다.
워커를 몇 대 두고 어떻게 늘리는가, 내부 토픽 3종이 어떤 성질을 가져야 하는가,
워커가 죽었을 때 로그를 어디서 찾는가, 그리고 409 Conflict가 왜 나는가입니다.
이 페이지는 그 네 가지에 집중합니다.
학습 목표
- 분산 모드 워커 그룹의 구성 요소와
group.id의 제약을 설명할 수 있습니다. - 내부 토픽 3종의 파티션·복제·정리 정책 요구사항을 구분할 수 있습니다.
- 커넥터·태스크 상태 6종을 구분하고 FAILED에서 자동 복구되지 않는다는 점을 말할 수 있습니다.
- 워커 장애 시 로그를 찾아가는 세 경로를 순서대로 실행할 수 있습니다.
409 Conflict가 나오는 두 상황을 구분할 수 있습니다.
이 도메인이 묻는 것
| 질문 형태 | 실제로 확인하는 것 |
|---|---|
| "워커를 3대에서 5대로 늘리면 무엇이 일어나는가?" | 리밸런스 동작과 tasks.max의 상한 역할 |
"config.storage.topic의 파티션 수는?" |
내부 토픽 3종의 요구사항 차이 |
| "태스크가 FAILED다. 어떻게 조치하는가?" | 자동 재시작이 없다는 사실과 restart API |
| "워커 1대가 죽었다. 그 태스크의 예외를 어디서 보는가?" | 로그 위치 · REST의 trace · 런타임 로그 레벨 변경 |
| "REST 호출이 409를 반환했다" | 리밸런스 진행 중 재시작 요청, 또는 리더가 아닌 워커로의 요청 상황 판단 |
핵심 개념 요약 — 운영 관점
워커 그룹의 모양
분산 모드의 워커들은 같은 group.id로 하나의 Connect 클러스터를 이룹니다.
커넥터는 논리적 작업 정의이고 실제 일은 태스크가 합니다.
태스크는 워커들에 분배되며, 커넥터가 요구하는 태스크 수는 tasks.max가 상한입니다.
리밸런스 — 왜 워커를 죽여도 바로 재분배되지 않는가
Connect는 기본적으로 증분 협력 리밸런스(incremental cooperative rebalancing)를 씁니다.
새로 추가·제거·이동이 필요한 태스크만 건드리고 나머지는 멈추지 않습니다.
connect.protocol의 기본값은 sessioned이며,
eager로 되돌리면 예전처럼 전체 태스크를 회수·재분배합니다.
워커가 그룹을 떠나면 Connect는 즉시 재분배하지 않고
scheduled.rebalance.max.delay.ms(기본 300000ms = 5분)을 기다립니다.
그 안에 워커가 돌아오면 이전 태스크를 그대로 되돌려받습니다.
돌아오지 않으면 남은 워커들에 재할당됩니다.
내부 토픽 3종 — 요구사항이 서로 다릅니다
분산 모드 워커는 상태를 로컬에 두지 않고 Kafka 토픽에 둡니다. 세 토픽 모두 복제되어야 하고 compaction이 걸려야 하지만, 파티션 요구사항이 다릅니다. 이 차이가 matching 문항으로 자주 나옵니다.
| 토픽 설정 | 담는 것 | 파티션 요구 | 기본 파티션 | 기본 RF |
|---|---|---|---|---|
config.storage.topic |
커넥터·태스크 설정 | 반드시 파티션 1개 — 설정 변경 순서가 전역으로 정해져야 합니다 | 1 (고정) | 3 |
offset.storage.topic |
source 커넥터의 오프셋 | 파티션 여러 개 (처리량) | 25 | 3 |
status.storage.topic |
커넥터·태스크 상태, 사용 토픽 목록 | 파티션 여러 개 | 5 | 3 |
커넥터·태스크 상태 6종
| 상태 | 의미 | 운영 조치 |
|---|---|---|
UNASSIGNED | 아직 워커에 할당되지 않음 | 리밸런스 대기 중일 수 있습니다. 잠시 기다립니다. |
RUNNING | 정상 동작 | — |
PAUSED | 관리자가 일시 정지. 자원은 계속 점유 | 재개는 빠릅니다. 원격 시스템 점검 시 사용. |
STOPPED | 태스크를 종료하고 자원 반납. 커넥터에만 적용 | 오프셋 관리 API는 이 상태에서만 동작합니다. |
FAILED | 예외로 실패. 원인은 상태 출력의 trace에 | 자동 재시작되지 않습니다. restart API로 수동 조치. |
RESTARTING | 재시작 중이거나 곧 재시작 예정 | — |
필수 명령어 — Connect는 CLI 대신 REST입니다
# 커넥터 목록
curl -s http://worker-1:8083/connectors
# 커넥터와 모든 태스크의 상태. 실패 시 각 태스크에 "trace" 필드가 붙습니다
curl -s http://worker-1:8083/connectors/pg-source/status
# 특정 태스크만
curl -s http://worker-1:8083/connectors/pg-source/tasks/0/status
# 이 커넥터가 실제로 쓰는 토픽 목록 (status.storage.topic 이 추적)
curl -s http://worker-1:8083/connectors/pg-source/topics
# 워커 버전과 연결된 클러스터 ID
curl -s http://worker-1:8083/
# 실패한 인스턴스만 재시작 (커넥터 + 태스크)
curl -s -X POST 'http://worker-1:8083/connectors/pg-source/restart?includeTasks=true&onlyFailed=true'
# 개별 태스크 재시작
curl -s -X POST http://worker-1:8083/connectors/pg-source/tasks/0/restart
# 일시 정지 / 재개 (자원 유지)
curl -s -X PUT http://worker-1:8083/connectors/pg-source/pause
curl -s -X PUT http://worker-1:8083/connectors/pg-source/resume
# 완전 정지 (자원 반납) — 오프셋 관리 전 필수 단계
curl -s -X PUT http://worker-1:8083/connectors/pg-source/stop
# 현재 오프셋 조회
curl -s http://worker-1:8083/connectors/pg-source/offsets
# 리셋 (커넥터가 STOPPED 여야 함)
curl -s -X DELETE http://worker-1:8083/connectors/pg-source/offsets
# 특정 위치로 변경 (커넥터가 STOPPED 여야 함)
curl -s -X PATCH http://worker-1:8083/connectors/pg-source/offsets \
-H 'Content-Type: application/json' \
-d '{"offsets":[{"partition":{"filename":"test.txt"},"offset":{"position":30}}]}'
# 현재 명시적으로 지정된 로거와 레벨
curl -s http://worker-1:8083/admin/loggers
# 특정 로거 레벨 조회
curl -s http://worker-1:8083/admin/loggers/org.apache.kafka.connect.runtime.WorkerSinkTask
# 레벨 변경 — 워커 재시작 없이 즉시 적용
curl -s -X PUT http://worker-1:8083/admin/loggers/org.apache.kafka.connect.runtime.WorkerSinkTask \
-H 'Content-Type: application/json' -d '{"level":"DEBUG"}'
반드시 외워야 할 설정값
| 설정 | 기본값 | 의미 | 운영 포인트 |
|---|---|---|---|
listenersworker |
http://:8083 |
REST API 리스너 | 워커 간 요청 전달도 이 경로를 씁니다. |
connect.protocolworker |
sessioned |
Connect 그룹 프로토콜 (eager/compatible/sessioned) |
eager는 전체 정지 재분배로 되돌립니다. 특별한 이유 없이 바꾸지 않습니다. |
scheduled.rebalance.max.delay.msworker |
300000 (5분) |
워커 이탈 후 재분배까지 대기 | 이 시간 동안 태스크가 미할당 상태로 남습니다. 롤링 업그레이드를 위한 값입니다. |
session.timeout.msworker |
10000 | 워커 장애 감지 시간 | 브로커의 그룹 최소·최대 세션 타임아웃 범위 안이어야 합니다. |
rebalance.timeout.msworker |
60000 | 리밸런스 시작 후 각 워커의 합류 제한 시간 | 초과하면 워커가 그룹에서 빠지고 오프셋 커밋이 실패합니다. |
offset.flush.interval.msworker |
60000 | source 오프셋 커밋 주기 | 워커가 죽으면 최대 이 주기만큼 재처리될 수 있습니다(at-least-once). |
offset.flush.timeout.msworker |
5000 | 오프셋 커밋 제한 시간 | 느린 sink에서 이 값이 짧으면 커밋 실패 로그가 반복됩니다. |
config.storage.replication.factorworker |
3 | 내부 토픽 복제 계수 (offset·status도 동일 기본값) | 브로커가 3대 미만인 실습 환경에서는 워커 기동이 실패합니다. |
offset.storage.partitionsworker |
25 | source 오프셋 토픽 파티션 수 | status.storage.partitions는 5, config는 1 고정입니다. |
exactly.once.source.supportworker |
disabled |
source 커넥터 EOS | 기존 클러스터는 preparing → enabled 2단계 롤링이 필요합니다. |
tasks.maxconnector |
1 | 태스크 수 상한 | 기본값이 1이라 워커를 늘려도 병렬성이 안 늘어납니다. 스케일링 사고 1위. |
errors.toleranceconnector |
none |
오류 허용 여부 | 기본은 첫 오류에 태스크 FAILED. all로 바꿔야 건너뜁니다. |
errors.log.enableconnector |
false |
오류를 애플리케이션 로그에 기록 | 켜야 토픽·파티션·오프셋이 로그에 남습니다. 진단의 출발점입니다. |
errors.deadletterqueue.topic.namesink |
빈 문자열 | DLQ 토픽 | 비어 있으면 DLQ가 없습니다. RF 기본값은 3, 헤더 기록은 기본 false입니다. |
errors.retry.timeoutconnector |
0 | 재시도 총 예산(ms) | 기본값 0은 재시도하지 않음입니다. |
장애 시나리오와 대응
시나리오 1 — 워커 1대가 죽었다. 그 태스크의 예외를 어디서 보는가
이것이 이 섹션의 대표 문항입니다. 죽은 노드의 로그는 REST로 가져올 수 없으므로 세 경로를 순서대로 씁니다.
- 죽은 노드의 로컬 로그 파일을 직접 본다. Connect 워커의 예외 스택은 그 워커의 프로세스 로그에 남습니다. 중앙 로그 수집이 있다면 그쪽에서 해당 워커의 호스트로 필터링합니다. 이 경로가 유일하게 완전한 정보를 줍니다.
-
살아 있는 워커에서
GET /connectors/{name}/status를 호출해trace를 읽는다. 상태는status.storage.topic으로 공유되므로, 죽은 워커가 남긴 실패 원인도 다른 워커에서 조회됩니다. 단 스택 트레이스 요약 수준이라 상세 컨텍스트는 부족할 수 있습니다. -
재현을 위해
PUT /admin/loggers/{logger}로 런타임 로그 레벨을 올린다. 워커를 재시작하지 않고 특정 로거만DEBUG로 바꿔 같은 실패를 다시 유도한 뒤 상세 로그를 수집합니다. 조사 후에는 반드시 원래 레벨로 되돌립니다 — 디스크가 빠르게 찹니다.
시나리오 2 — REST 호출이 409 Conflict를 반환한다
| 상황 | 왜 409인가 | 조치 |
|---|---|---|
| 리밸런스 진행 중에 restart 요청 | 태스크 할당이 확정되지 않은 상태에서 재시작을 받아들일 수 없습니다. | 리밸런스가 끝난 뒤 재시도합니다. 다만 리밸런스 자체가 태스크를 다시 시작하므로 재시도가 불필요할 수 있습니다. |
| 설정 변경·커넥터 생성 요청이 리밸런스와 겹침 | 설정 쓰기는 리더 워커를 거쳐야 하고, 리밸런스 중에는 리더가 확정되지 않습니다. | 잠시 후 재시도합니다. 자동화 스크립트라면 409를 재시도 대상으로 다뤄야 합니다. |
| 워커가 리더로 요청을 전달할 수 없음 | follower 워커가 받은 쓰기 요청은 리더 REST로 전달됩니다. 그 경로가 막히면 실패합니다. | rest.advertised.host.name / rest.advertised.port / rest.advertised.listener가 워커 간에 실제 도달 가능한 값인지 확인합니다. |
시나리오 3 — 워커를 늘렸는데 처리량이 그대로다
GET /connectors/{name}/tasks로 실제 태스크 수를 확인합니다.tasks.max가 기본값 1이면 워커가 몇 대든 태스크는 1개입니다. 이것이 1순위 원인입니다.tasks.max를 올려도 늘지 않으면 커넥터 자체의 상한을 봅니다. source 커넥터는 나눌 수 있는 단위(테이블·파일·파티션) 수를 넘겨 태스크를 만들지 못합니다.- sink 커넥터의 태스크 수는 입력 토픽의 파티션 수가 실질 상한입니다. 파티션 4개에 태스크 8개를 줘도 4개만 일합니다.
- 태스크가 충분한데도 느리면
errors.retry.timeout·타임아웃 설정과 대상 시스템 쪽을 봅니다.
시나리오 4 — 태스크가 실패했는데 아무 일도 일어나지 않는다
- Connect는 실패한 태스크를 자동으로 재시작하지 않습니다. 이것이 정상 동작입니다.
GET /connectors/{name}/status의trace로 원인을 먼저 확인합니다. 원인을 모른 채 재시작하면 같은 실패가 반복됩니다.- 데이터 문제(역직렬화 실패, 스키마 불일치)라면
errors.tolerance=all과 DLQ를 붙여 나쁜 레코드를 격리합니다. DLQ 토픽 이름을 지정하지 않으면 DLQ가 만들어지지 않습니다. - 일시적 문제(대상 시스템 재시작)라면
POST /connectors/{name}/restart?includeTasks=true&onlyFailed=true로 실패분만 되살립니다. - 재발 방지를 위해
errors.log.enable=true와errors.deadletterqueue.context.headers.enable=true를 켜 두면 실패 레코드의 원래 토픽·파티션·오프셋이 남습니다.
자주 나오는 함정
관련 케이스 스터디
- 케이스 10 · 큰 메시지가 무한 재시도로 쌓였다 — 재시도 예산과 DLQ 설계
- 케이스 9 · 스키마 배포 후 전체 컨슈머가 죽었다 — 역직렬화 실패와
errors.tolerance
커넥터 설정·SMT 문법은 Connect REST · SMT 치트시트, 구조와 개발 관점은 9장 Kafka Connect, 동작하는 CDC 파이프라인은 예제 8에 있습니다.
미니 퀴즈
내부 토픽 요구사항과 장애 조사 순서가 중심입니다.
공식 문서 출처
- Connect Administration — 증분 협력 리밸런스,
scheduled.rebalance.max.delay.ms, 상태 6종,409 Conflict, 상태 전파 지연 - Running Kafka Connect — 분산 모드
group.id제약, 내부 토픽 3종의 파티션·복제·compaction 요구 - Connect REST Interface — 상태·재시작·pause/stop/resume, 오프셋 관리 엔드포인트,
/admin/loggers,rest.advertised.* - Error Reporting in Connect —
errors.log.enable,errors.tolerance, DLQ - Kafka Connect Configs — 워커 설정 기본값
- Sink Connector Configs — DLQ·오류 처리 기본값,
tasks.max