학습 목표

이 도메인이 묻는 것

이 섹션의 전형적인 질문 형태
질문 형태실제로 확인하는 것
"워커를 3대에서 5대로 늘리면 무엇이 일어나는가?" 리밸런스 동작과 tasks.max의 상한 역할
"config.storage.topic의 파티션 수는?" 내부 토픽 3종의 요구사항 차이
"태스크가 FAILED다. 어떻게 조치하는가?" 자동 재시작이 없다는 사실과 restart API
"워커 1대가 죽었다. 그 태스크의 예외를 어디서 보는가?" 로그 위치 · REST의 trace · 런타임 로그 레벨 변경
"REST 호출이 409를 반환했다" 리밸런스 진행 중 재시작 요청, 또는 리더가 아닌 워커로의 요청 상황 판단

핵심 개념 요약 — 운영 관점

워커 그룹의 모양

분산 모드의 워커들은 같은 group.id로 하나의 Connect 클러스터를 이룹니다. 커넥터는 논리적 작업 정의이고 실제 일은 태스크가 합니다. 태스크는 워커들에 분배되며, 커넥터가 요구하는 태스크 수는 tasks.max가 상한입니다.

Kafka Connect 아키텍처 — 워커 · 커넥터 · 태스크 · 내부 토픽 3종 분산 모드 Connect 클러스터를 그린 그림입니다. 같은 group.id 를 가진 워커 세 대가 하나의 클러스터를 이루고, 그중 한 대가 그룹 리더가 되어 태스크 할당을 계산합니다. 커넥터 A 는 소스 커넥터로 태스크 세 개, 커넥터 B 는 싱크 커넥터로 태스크 세 개를 요구하며 합계 여섯 개의 태스크가 세 워커에 두 개씩 분산됩니다. 커넥터는 설정일 뿐이고 실제로 데이터를 옮기는 것은 태스크입니다. 모든 워커가 REST API 를 제공하므로 어느 워커에 요청해도 됩니다. 클러스터 상태는 Kafka 내부 토픽 세 개에 저장됩니다. config.storage.topic 은 커넥터와 태스크 설정을 담고 파티션 하나에 compact 정책을 씁니다. offset.storage.topic 은 소스 커넥터의 오프셋을 담고 기본 25 파티션에 compact 정책을 씁니다. status.storage.topic 은 커넥터와 태스크 상태를 담고 compact 정책을 씁니다. 워커가 자기 상태를 이 토픽들에 저장하기 때문에 워커 자체는 상태를 갖지 않고 교체할 수 있습니다. Kafka Connect 아키텍처 — 상태는 워커가 아니라 Kafka 에 있습니다 Connect 클러스터 — group.id=connect-cluster (컨슈머 그룹 id 와 겹치면 안 됩니다) worker-1 (그룹 리더) REST :8083 A · task 0 B · task 0 태스크 할당 계산 설정 변경 감지 → 리밸런스 worker-2 REST :8083 A · task 1 B · task 1 리더가 준 할당대로 실행 worker-3 REST :8083 A · task 2 B · task 2 리더가 준 할당대로 실행 config.storage.topic 커넥터 · 태스크 설정 파티션 1개 · compact 필수 offset.storage.topic 소스 커넥터 오프셋 기본 25 파티션 · compact status.storage.topic 커넥터 · 태스크 상태 compact · 상태 API 가 읽음 커넥터는 설정, 태스크는 일꾼입니다. tasks.max (기본 1) 가 태스크 수의 상한이고 실제 수는 커넥터가 정합니다. 싱크 커넥터의 오프셋은 이 세 토픽이 아니라 __consumer_offsets 에 저장됩니다 — 비대칭은 D-081 을 보세요.
Connect 아키텍처 — 워커 3대에 커넥터와 태스크가 분배되는 구조, 리더 워커의 역할

리밸런스 — 왜 워커를 죽여도 바로 재분배되지 않는가

Connect는 기본적으로 증분 협력 리밸런스(incremental cooperative rebalancing)를 씁니다. 새로 추가·제거·이동이 필요한 태스크만 건드리고 나머지는 멈추지 않습니다. connect.protocol의 기본값은 sessioned이며, eager로 되돌리면 예전처럼 전체 태스크를 회수·재분배합니다.

워커가 그룹을 떠나면 Connect는 즉시 재분배하지 않고 scheduled.rebalance.max.delay.ms(기본 300000ms = 5분)을 기다립니다. 그 안에 워커가 돌아오면 이전 태스크를 그대로 되돌려받습니다. 돌아오지 않으면 남은 워커들에 재할당됩니다.

Connect 태스크 재분배 — 워커 이탈과 지연 리밸런스 워커 세 대에 태스크 여섯 개가 나뉘어 있던 Connect 클러스터에서 워커 한 대가 빠졌을 때 무슨 일이 일어나는지 세 단계로 보여줍니다. 첫 단계는 정상 상태로 worker-1 이 A0 과 B0, worker-2 가 A1 과 B1, worker-3 이 A2 와 B2 를 실행합니다. 두 번째 단계는 worker-3 이 빠진 직후입니다. Connect 는 즉시 재할당하지 않고 scheduled.rebalance.max.delay.ms 기본값 300000 밀리초, 즉 5분을 기다립니다. 이 시간 안에 worker-3 이 돌아오면 원래 태스크를 그대로 되돌려 받습니다. 대신 그동안 A2 와 B2 는 아무도 실행하지 않는 상태로 남습니다. 세 번째 단계는 대기 시간이 지나도 돌아오지 않은 경우입니다. A2 와 B2 가 남은 두 워커로 재할당됩니다. 증분 협력 리밸런싱이 기본이므로 이미 잘 돌고 있는 A0, B0, A1, B1 은 멈추지 않습니다. 태스크 재분배 — 워커가 빠지면 바로 옮기지 않습니다 ① 정상 — 태스크 6개 / 워커 3대 worker-1 A0 B0 worker-2 A1 B1 worker-3 A2 B2 ② worker-3 이탈 → 대기 worker-1 그대로 실행 A0 B0 worker-2 그대로 실행 A1 B1 worker-3 이탈 A2 ? B2 ? ③ 대기 시간 경과 → 재할당 worker-1 A0 B0 A2 worker-2 A1 B1 B2 worker-3 없음 새 워커가 들어오면 다시 리밸런스 왜 바로 옮기지 않는가 — scheduled.rebalance.max.delay.ms 기본 300000 (5분) 이 시간 안에 worker-3 이 돌아오면 원래 태스크를 그대로 되돌려 받습니다 (재시작·재복구 비용 절약). 대신 그동안 A2 · B2 는 아무도 실행하지 않습니다 — 롤링 재시작 때 지연이 보이는 이유입니다. 기본은 증분 협력 리밸런싱입니다. 영향받지 않는 A0·B0·A1·B1 은 멈추지 않습니다 (A = 소스 커넥터, B = 싱크 커넥터). 2.3 이전의 eager 프로토콜은 리밸런스마다 모든 태스크를 멈췄습니다.
워커 이탈 시 태스크 재분배 — 지연 대기 후 남은 워커로 재할당되는 흐름

내부 토픽 3종 — 요구사항이 서로 다릅니다

분산 모드 워커는 상태를 로컬에 두지 않고 Kafka 토픽에 둡니다. 세 토픽 모두 복제되어야 하고 compaction이 걸려야 하지만, 파티션 요구사항이 다릅니다. 이 차이가 matching 문항으로 자주 나옵니다.

Connect 내부 토픽 3종
토픽 설정 담는 것 파티션 요구 기본 파티션 기본 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/
복구 조작 — FAILED는 수동으로 되살립니다
# 실패한 인스턴스만 재시작 (커넥터 + 태스크)
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
오프셋 관리 — STOPPED 상태에서만 됩니다
# 현재 오프셋 조회
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"}'

반드시 외워야 할 설정값

Kafka 4.3 Connect 워커·커넥터 기본값
설정 기본값 의미 운영 포인트
listeners
worker
http://:8083 REST API 리스너 워커 간 요청 전달도 이 경로를 씁니다.
connect.protocol
worker
sessioned Connect 그룹 프로토콜 (eager/compatible/sessioned) eager는 전체 정지 재분배로 되돌립니다. 특별한 이유 없이 바꾸지 않습니다.
scheduled.rebalance.max.delay.ms
worker
300000
(5분)
워커 이탈 후 재분배까지 대기 이 시간 동안 태스크가 미할당 상태로 남습니다. 롤링 업그레이드를 위한 값입니다.
session.timeout.ms
worker
10000 워커 장애 감지 시간 브로커의 그룹 최소·최대 세션 타임아웃 범위 안이어야 합니다.
rebalance.timeout.ms
worker
60000 리밸런스 시작 후 각 워커의 합류 제한 시간 초과하면 워커가 그룹에서 빠지고 오프셋 커밋이 실패합니다.
offset.flush.interval.ms
worker
60000 source 오프셋 커밋 주기 워커가 죽으면 최대 이 주기만큼 재처리될 수 있습니다(at-least-once).
offset.flush.timeout.ms
worker
5000 오프셋 커밋 제한 시간 느린 sink에서 이 값이 짧으면 커밋 실패 로그가 반복됩니다.
config.storage.replication.factor
worker
3 내부 토픽 복제 계수 (offset·status도 동일 기본값) 브로커가 3대 미만인 실습 환경에서는 워커 기동이 실패합니다.
offset.storage.partitions
worker
25 source 오프셋 토픽 파티션 수 status.storage.partitions5, config는 1 고정입니다.
exactly.once.source.support
worker
disabled source 커넥터 EOS 기존 클러스터는 preparingenabled 2단계 롤링이 필요합니다.
tasks.max
connector
1 태스크 수 상한 기본값이 1이라 워커를 늘려도 병렬성이 안 늘어납니다. 스케일링 사고 1위.
errors.tolerance
connector
none 오류 허용 여부 기본은 첫 오류에 태스크 FAILED. all로 바꿔야 건너뜁니다.
errors.log.enable
connector
false 오류를 애플리케이션 로그에 기록 켜야 토픽·파티션·오프셋이 로그에 남습니다. 진단의 출발점입니다.
errors.deadletterqueue.topic.name
sink
빈 문자열 DLQ 토픽 비어 있으면 DLQ가 없습니다. RF 기본값은 3, 헤더 기록은 기본 false입니다.
errors.retry.timeout
connector
0 재시도 총 예산(ms) 기본값 0은 재시도하지 않음입니다.

장애 시나리오와 대응

시나리오 1 — 워커 1대가 죽었다. 그 태스크의 예외를 어디서 보는가

이것이 이 섹션의 대표 문항입니다. 죽은 노드의 로그는 REST로 가져올 수 없으므로 세 경로를 순서대로 씁니다.

  1. 죽은 노드의 로컬 로그 파일을 직접 본다. Connect 워커의 예외 스택은 그 워커의 프로세스 로그에 남습니다. 중앙 로그 수집이 있다면 그쪽에서 해당 워커의 호스트로 필터링합니다. 이 경로가 유일하게 완전한 정보를 줍니다.
  2. 살아 있는 워커에서 GET /connectors/{name}/status를 호출해 trace를 읽는다. 상태는 status.storage.topic으로 공유되므로, 죽은 워커가 남긴 실패 원인도 다른 워커에서 조회됩니다. 단 스택 트레이스 요약 수준이라 상세 컨텍스트는 부족할 수 있습니다.
  3. 재현을 위해 PUT /admin/loggers/{logger}로 런타임 로그 레벨을 올린다. 워커를 재시작하지 않고 특정 로거만 DEBUG로 바꿔 같은 실패를 다시 유도한 뒤 상세 로그를 수집합니다. 조사 후에는 반드시 원래 레벨로 되돌립니다 — 디스크가 빠르게 찹니다.

시나리오 2 — REST 호출이 409 Conflict를 반환한다

409 Conflict의 원인과 조치
상황왜 409인가조치
리밸런스 진행 중에 restart 요청 태스크 할당이 확정되지 않은 상태에서 재시작을 받아들일 수 없습니다. 리밸런스가 끝난 뒤 재시도합니다. 다만 리밸런스 자체가 태스크를 다시 시작하므로 재시도가 불필요할 수 있습니다.
설정 변경·커넥터 생성 요청이 리밸런스와 겹침 설정 쓰기는 리더 워커를 거쳐야 하고, 리밸런스 중에는 리더가 확정되지 않습니다. 잠시 후 재시도합니다. 자동화 스크립트라면 409를 재시도 대상으로 다뤄야 합니다.
워커가 리더로 요청을 전달할 수 없음 follower 워커가 받은 쓰기 요청은 리더 REST로 전달됩니다. 그 경로가 막히면 실패합니다. rest.advertised.host.name / rest.advertised.port / rest.advertised.listener가 워커 간에 실제 도달 가능한 값인지 확인합니다.

시나리오 3 — 워커를 늘렸는데 처리량이 그대로다

  1. GET /connectors/{name}/tasks실제 태스크 수를 확인합니다.
  2. tasks.max가 기본값 1이면 워커가 몇 대든 태스크는 1개입니다. 이것이 1순위 원인입니다.
  3. tasks.max를 올려도 늘지 않으면 커넥터 자체의 상한을 봅니다. source 커넥터는 나눌 수 있는 단위(테이블·파일·파티션) 수를 넘겨 태스크를 만들지 못합니다.
  4. sink 커넥터의 태스크 수는 입력 토픽의 파티션 수가 실질 상한입니다. 파티션 4개에 태스크 8개를 줘도 4개만 일합니다.
  5. 태스크가 충분한데도 느리면 errors.retry.timeout·타임아웃 설정과 대상 시스템 쪽을 봅니다.

시나리오 4 — 태스크가 실패했는데 아무 일도 일어나지 않는다

  1. Connect는 실패한 태스크를 자동으로 재시작하지 않습니다. 이것이 정상 동작입니다.
  2. GET /connectors/{name}/statustrace로 원인을 먼저 확인합니다. 원인을 모른 채 재시작하면 같은 실패가 반복됩니다.
  3. 데이터 문제(역직렬화 실패, 스키마 불일치)라면 errors.tolerance=all과 DLQ를 붙여 나쁜 레코드를 격리합니다. DLQ 토픽 이름을 지정하지 않으면 DLQ가 만들어지지 않습니다.
  4. 일시적 문제(대상 시스템 재시작)라면 POST /connectors/{name}/restart?includeTasks=true&onlyFailed=true로 실패분만 되살립니다.
  5. 재발 방지를 위해 errors.log.enable=trueerrors.deadletterqueue.context.headers.enable=true를 켜 두면 실패 레코드의 원래 토픽·파티션·오프셋이 남습니다.

자주 나오는 함정

관련 케이스 스터디

커넥터 설정·SMT 문법은 Connect REST · SMT 치트시트, 구조와 개발 관점은 9장 Kafka Connect, 동작하는 CDC 파이프라인은 예제 8에 있습니다.

미니 퀴즈

내부 토픽 요구사항과 장애 조사 순서가 중심입니다.

공식 문서 출처