이 도메인의 학습 목표

이 도메인이 묻는 것

Kafka Connect는 9장 Kafka Connect의 범위입니다. Connect는 "커넥터 플러그인을 실행해 주는 워커 클러스터"이고, 개발자가 쓰는 것은 대부분 커넥터 설정 JSON과 REST API입니다. 따라서 출제도 설정 키 이름, REST 경로, 그리고 실패했을 때 무엇을 봐야 하는가에 몰립니다.

핵심 개념 (1) 구조와 실행 모드

Connect의 4개 구성 요소
요소역할단위
workerJVM 프로세스. 커넥터와 태스크를 실행하고 REST API를 제공프로세스
connector"무엇을 어디로 옮기는가"의 정의. 실제 데이터를 옮기지 않고 태스크로 쪼갠다논리 작업
task실제로 데이터를 읽고 쓰는 단위. 병렬성의 단위스레드
converterConnect 내부 표현 ↔ Kafka 바이트 직렬화 변환레코드
transform (SMT)레코드 하나 단위의 가벼운 변환. converter 앞/뒤에 체인으로 붙음레코드
standalone vs distributed
관점standalonedistributed
실행 스크립트connect-standalone.shconnect-distributed.sh
커넥터 설정 전달커맨드라인의 properties 파일REST API로만 (커맨드라인 불가)
오프셋 저장offset.storage.file.filename (로컬 파일)offset.storage.topic (Kafka 토픽)
설정·상태 저장로컬 (프로세스 재시작 시 초기화 위험)config.storage.topic / status.storage.topic
확장 · 내결함성없음 (프로세스 1개)워커 추가로 확장, 태스크 자동 재분배
용도개발·테스트, 단일 노드 로그 수집프로덕션

내부 토픽 3종

distributed 모드 내부 토픽 (Apache Kafka 4.3 워커 기본값)
설정 용도 파티션 권장 RF 기본값
config.storage.topic 커넥터·태스크 설정 반드시 1개 (순서가 중요) 3 (config.storage.replication.factor)
offset.storage.topic source 커넥터의 소스 시스템 오프셋 많이 (기본 offset.storage.partitions=25) 3
status.storage.topic 커넥터·태스크 상태 여러 개 (기본 status.storage.partitions=5) 3

세 토픽 모두 cleanup.policy=compact으로 만들어야 합니다. 키별 최신 값(최신 설정, 최신 오프셋, 최신 상태)만 필요하기 때문입니다. Connect가 자동 생성할 수도 있지만, 공식 문서는 파티션 수와 복제 계수를 직접 지정해 수동 생성하는 것을 권장합니다.

핵심 개념 (2) source와 sink의 오프셋 비대칭

이 도메인에서 가장 자주 나오는 단일 사실입니다.

오프셋을 어디에 저장하는가
관점Source 커넥터Sink 커넥터
무엇의 오프셋인가 소스 시스템의 위치 (파일 오프셋, DB binlog 위치, 테이블 커서 등) Kafka 토픽의 컨슈머 오프셋
저장 위치 (distributed) offset.storage.topic
(Connect 내부 토픽)
__consumer_offsets
(일반 컨슈머와 동일)
저장 위치 (standalone) offset.storage.file.filename
(로컬 파일)
__consumer_offsets (모드와 무관)
컨슈머 그룹 이름 없음 (컨슈머가 아님) connect-{커넥터 이름}
확인 방법 GET /connectors/{name}/offsets kafka-consumer-groups --describe --group connect-{name}
또는 같은 REST 엔드포인트
커밋 주기 offset.flush.interval.ms 기본 60000 (1분)

핵심 개념 (3) 컨버터와 SMT

컨버터 — 직렬화 형식을 결정합니다

worker 또는 커넥터 레벨에서 지정
key.converter=org.apache.kafka.connect.storage.StringConverter
value.converter=org.apache.kafka.connect.json.JsonConverter

# JsonConverter 는 기본적으로 Connect 스키마를 JSON 안에 함께 넣는다.
# 순수 JSON 을 원하면 끈다 (스키마 정보를 잃는 대가가 있다).
value.converter.schemas.enable=false

SMT 체인의 위치와 순서

SMT가 적용되는 지점
방향처리 순서
Source 소스 시스템 → SourceTask → SMT 체인 → converter → Kafka 토픽
Sink Kafka 토픽 → converter → SMT 체인 → SinkTask → 대상 시스템
SMT 체인 설정 — transforms에 나열한 순서대로 적용됩니다
transforms=unwrap,addSource,route

transforms.unwrap.type=org.apache.kafka.connect.transforms.ExtractField$Value
transforms.unwrap.field=payload

transforms.addSource.type=org.apache.kafka.connect.transforms.InsertField$Value
transforms.addSource.static.field=source_system
transforms.addSource.static.value=oracle-crm

transforms.route.type=org.apache.kafka.connect.transforms.RegexRouter
transforms.route.regex=(.*)
transforms.route.replacement=crm_$1
Apache Kafka에 포함된 대표 SMT (전부 org.apache.kafka.connect.transforms 패키지)
SMT용도주요 파라미터
InsertField필드 추가 (정적 값, 타임스탬프, 토픽·파티션·오프셋)static.field, timestamp.field, offset.field
ReplaceField필드 이름 변경·포함·제외renames, include, exclude
MaskField민감 필드를 기본값·지정값으로 마스킹fields, replacement
ValueToKey값의 필드를 키로 승격 (컴팩션·조인 준비에 필수)fields
ExtractField중첩 구조에서 한 필드만 꺼냄field
HoistField반대로, 스칼라를 구조체 한 필드로 감쌈field
Cast필드 타입 변환spec (예: amount:float64)
TimestampConverter문자열·유닉스시각·Date 상호 변환field, target.type, format
RegexRouter대상 토픽 이름을 정규식으로 변경regex, replacement
TimestampRouter토픽 이름에 시각을 붙임 (일별 인덱스 등)topic.format, timestamp.format
Filter조건에 맞는 레코드를 버림 (predicate와 함께 씀)predicate, negate
Flatten중첩 구조를 평탄화delimiter
SetSchemaMetadata스키마 이름·버전 지정schema.name, schema.version
HeaderFrom · InsertHeader · DropHeaders헤더 조작fields, headers

핵심 개념 (4) 에러 처리와 DLQ

Connect 에러 처리 설정 (Apache Kafka 4.3 기본값)
설정기본값의미
errors.tolerance none none이면 첫 에러에 태스크가 FAILED. all이면 건너뜀
errors.log.enable false 에러와 문제 레코드의 토픽·파티션·오프셋을 로그로 남김
errors.log.include.messages false 키·값·헤더까지 로그에 남김. 민감정보 노출 위험
errors.deadletterqueue.topic.name "" (빈 문자열) 비어 있으면 DLQ 없음. 토픽 이름을 지정해야 활성화. sink 전용
errors.deadletterqueue.topic.replication.factor 3 단일 브로커 개발 환경에서는 1로 낮춰야 함
errors.deadletterqueue.context.headers.enable false 실패 원인·원본 위치를 헤더에 기록. 사실상 켜야 쓸모가 있음
errors.retry.timeout 0 재시도에 쓰는 시간 예산. 0이면 재시도 없음. -1은 무한
errors.retry.delay.max.ms 60000 재시도 간 지수 백오프의 상한
실용적인 DLQ 설정 세트 (sink 커넥터)
{
  "name": "orders-jdbc-sink",
  "config": {
    "connector.class": "io.confluent.connect.jdbc.JdbcSinkConnector",
    "topics": "orders",
    "tasks.max": "3",

    "errors.tolerance": "all",
    "errors.log.enable": "true",
    "errors.deadletterqueue.topic.name": "orders-jdbc-sink-dlq",
    "errors.deadletterqueue.topic.replication.factor": "3",
    "errors.deadletterqueue.context.headers.enable": "true",
    "errors.retry.timeout": "300000",
    "errors.retry.delay.max.ms": "30000"
  }
}

핵심 개념 (5) REST API

Connect REST API 주요 엔드포인트 (Apache Kafka 4.3)
메서드 · 경로용도
GET /connectors커넥터 목록
POST /connectors커넥터 생성 ({"name":…, "config":{…}})
GET /connectors/{name}커넥터 정보
GET /connectors/{name}/config설정 조회
PUT /connectors/{name}/config설정 전체 교체 (없으면 생성)
PATCH /connectors/{name}/config설정 부분 수정
GET /connectors/{name}/status커넥터·태스크 상태와 실패 trace — 장애 시 첫 호출
GET /connectors/{name}/tasks태스크 목록과 설정
GET /connectors/{name}/tasks/{taskId}/status개별 태스크 상태
POST /connectors/{name}/restart커넥터 재시작 (?includeTasks=true&onlyFailed=true 지원)
POST /connectors/{name}/tasks/{taskId}/restart실패한 태스크만 재시작
PUT /connectors/{name}/pause · resume · stop일시정지 / 재개 / 정지
DELETE /connectors/{name}커넥터 삭제
GET /connectors/{name}/offsets오프셋 조회 (source·sink 모두)
PATCH /connectors/{name}/offsets · DELETE …/offsets오프셋 변경 / 초기화 (커넥터가 STOPPED여야 함)
GET /connectors/{name}/topics이 커넥터가 사용한 토픽 목록
GET /connector-plugins설치된 플러그인 목록
PUT /connector-plugins/{type}/config/validate배포 전 설정 검증
GET /admin/loggers · PUT /admin/loggers/{name}런타임 로그 레벨 조회·변경 (재시작 없이)

태스크 상태와 409 Conflict

상태 값과 대응
상태의미대응
RUNNING정상
PAUSED사람이 멈춤resume
FAILED예외로 중단. 자동 재시작하지 않음statustrace 확인 → 원인 수정 → restart
UNASSIGNED아직 워커에 배정되지 않음리밸런스 완료 대기
RESTARTING재시작 중대기

확장 전략 — tasks.max의 한계

tasks.max상한입니다. 커넥터가 그만큼의 병렬성을 만들 수 없으면 더 적은 태스크를 만듭니다.

반드시 외워야 할 설정값

Connect 워커·커넥터 필수 암기 설정 (Apache Kafka 4.3)
설정기본값시험 포인트
offset.flush.interval.ms60000source 오프셋 커밋 주기 (1분)
offset.flush.timeout.ms5000플러시 대기 상한
offset.storage.partitions25자동 생성 시 파티션 수
status.storage.partitions5
config.storage.replication.factor3offset·status도 동일하게 3
rebalance.timeout.ms60000워커 리밸런스 상한
plugin.pathnull비어 있으면 커넥터를 못 찾습니다 (기동 후 목록이 빔)
errors.tolerancenoneDLQ를 쓰려면 all
errors.retry.timeout0기본값은 재시도 없음
errors.deadletterqueue.topic.name""기본값은 DLQ 비활성
exactly.once.source.supportdisabledsource EOS(KIP-618)는 기본 비활성
transaction.boundary (source)pollinterval, connector 선택 가능
tasks.max(커넥터별)상한일 뿐, 실제 태스크 수는 더 적을 수 있음

자주 나오는 함정

함정 1 — source vs sink 오프셋 저장 위치

이 비대칭이 이 도메인 최다 출제 지점입니다
관점SourceSink
무엇인가외부 시스템 → KafkaKafka → 외부 시스템
어디 설정인가워커 offset.storage.topic컨슈머 (Connect가 내부 생성)
언제 발동offset.flush.interval.ms마다컨슈머 커밋 시점
혼동 시 결과복구 시 잘못된 곳을 초기화해 전량 재수집 또는 재적재
출제 형태"sink 커넥터의 진행 상황을 확인하는 명령은?" → kafka-consumer-groups

함정 2 — standalone vs distributed

"단일 노드 distributed"도 정상적인 프로덕션 구성입니다
관점standalonedistributed
무엇인가단일 프로세스, 로컬 상태클러스터, Kafka 토픽에 상태 저장
어디 설정인가워커 properties + 커넥터 properties워커 properties + REST
언제 발동기동 시 커넥터가 즉시 시작REST 호출 시 커넥터 시작
혼동 시 결과distributed에 커넥터 properties를 커맨드라인으로 주면 무시됩니다
출제 형태"워커를 재시작하니 커넥터가 사라졌다" → standalone의 로컬 상태

함정 3 — SMT vs Kafka Streams

"어디까지 SMT로 할 것인가"
관점SMTKafka Streams
무엇인가레코드 1건 단위의 stateless 변환상태를 갖는 스트림 처리 애플리케이션
어디 설정인가커넥터 설정 JSON별도 애플리케이션 (라이브러리)
가능한 것필드 추가·삭제·마스킹·타입 변환·라우팅·필터집계 · 조인 · 윈도우 · 상태 저장
혼동 시 결과SMT 체인이 10개를 넘어가면 유지보수 불가. 그때는 Streams로 옮길 신호입니다
출제 형태"두 토픽을 조인해 적재하려면?" → SMT로는 불가, Streams(또는 ksqlDB)

함정 4 — errors.tolerance vs errors.retry.timeout

둘 다 "에러"를 다루지만 하는 일이 다릅니다
관점errors.toleranceerrors.retry.timeout
무엇인가실패 레코드를 건너뛸지 결정실패를 다시 시도할 시간 예산
기본값none (첫 에러에 태스크 FAILED)0 (재시도 없음)
언제 발동변환·컨버터·sink 쓰기 실패 시재시도 가능한 실패 시
혼동 시 결과all만 켜면 일시적 장애로 데이터가 조용히 버려짐재시도만 켜면 영구 실패에 태스크가 계속 죽음
출제 형태"일시 장애는 재시도, 영구 불량은 DLQ로 보내려면?" → 둘을 함께 설정

함정 5 — 내부 토픽 3종의 파티션 수

세 토픽의 요구사항이 다릅니다
관점config.storage.topicoffset.storage.topicstatus.storage.topic
파티션1 (필수)많이 (기본 25)여러 개 (기본 5)
cleanup.policycompact — 세 개 모두
담는 것커넥터·태스크 설정source 오프셋커넥터·태스크 상태
혼동 시 결과파티션을 늘리면 설정 순서가 깨짐파티션이 적으면 오프셋 커밋 병목
출제 형태"내부 토픽 3종의 공통 정책은?" → compact

함정 6 — tasks.max vs 파티션 수 vs 워커 수

처리량을 늘리는 세 손잡이
관점tasks.max입력 토픽 파티션워커 수
무엇을 바꾸는가태스크 수 상한sink 태스크의 실질 상한태스크를 놓을 자리
어디 설정인가커넥터 설정토픽인프라
단독으로 효과파티션이 적으면 무효sink에서 효과적태스크가 1개면 무효
출제 형태"파티션 6개 토픽에 tasks.max=12. 실제 동작 태스크 수는?" → 6

설정 읽기 문제 대비

스니펫 1 — 이 커넥터에서 DLQ는 동작하는가

POST /connectors 본문
{
  "name": "es-sink",
  "config": {
    "connector.class": "io.confluent.connect.elasticsearch.ElasticsearchSinkConnector",
    "topics": "logs",
    "tasks.max": "4",
    "errors.deadletterqueue.topic.name": "es-sink-dlq",
    "errors.log.enable": "true"
  }
}

스니펫 2 — 이 SMT 체인의 결과

connector.properties
transforms=key,mask
transforms.key.type=org.apache.kafka.connect.transforms.ValueToKey
transforms.key.fields=userId
transforms.mask.type=org.apache.kafka.connect.transforms.MaskField$Value
transforms.mask.fields=userId,ssn

스니펫 3 — 이 요청이 409를 받는 이유

운영 중 실행한 명령과 응답
$ curl -s -o /dev/null -w "%{http_code}\n" \
    -X POST http://connect-1:8083/connectors/orders-sink/tasks/2/restart
409

$ curl -s http://connect-1:8083/connectors/orders-sink/status | head -20
# "connector": { "state": "RUNNING" }
# "tasks": [ { "id": 0, "state": "UNASSIGNED" }, { "id": 1, "state": "UNASSIGNED" } ]

도메인 미니 퀴즈

공식 문서 출처