세 번째로 큰 도메인입니다. 코드를 쓰지 않고 설정(JSON)으로 데이터를 옮기는 프레임워크라서,
출제 지점도 "무엇을 어디에 설정하는가"에 집중됩니다.
특히 source와 sink의 오프셋 저장 위치가 다르다는 비대칭,
SMT의 능력 범위, DLQ 설정 조합이 반복해서 나옵니다.
이 도메인의 학습 목표
worker · connector · task · converter · transform의 역할을 구분할 수 있습니다.
standalone과 distributed 모드의 차이, 내부 토픽 3종의 용도를 말할 수 있습니다.
source와 sink가 오프셋을 어디에 저장하는지 정확히 구분할 수 있습니다.
SMT로 할 수 있는 일과 할 수 없는 일의 경계를 판단할 수 있습니다.
DLQ를 켜는 설정 조합과 errors.tolerance의 의미를 설명할 수 있습니다.
이 도메인이 묻는 것
Kafka Connect는 9장 Kafka Connect의 범위입니다.
Connect는 "커넥터 플러그인을 실행해 주는 워커 클러스터"이고,
개발자가 쓰는 것은 대부분 커넥터 설정 JSON과 REST API입니다.
따라서 출제도 설정 키 이름, REST 경로, 그리고 실패했을 때 무엇을 봐야 하는가에 몰립니다.
핵심 개념 (1) 구조와 실행 모드
Connect의 4개 구성 요소
요소
역할
단위
worker
JVM 프로세스. 커넥터와 태스크를 실행하고 REST API를 제공
프로세스
connector
"무엇을 어디로 옮기는가"의 정의. 실제 데이터를 옮기지 않고 태스크로 쪼갠다
논리 작업
task
실제로 데이터를 읽고 쓰는 단위. 병렬성의 단위
스레드
converter
Connect 내부 표현 ↔ Kafka 바이트 직렬화 변환
레코드
transform (SMT)
레코드 하나 단위의 가벼운 변환. converter 앞/뒤에 체인으로 붙음
레코드
standalone vs distributed
관점
standalone
distributed
실행 스크립트
connect-standalone.sh
connect-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
StringConverter · ByteArrayConverter — 스키마 없음. 가장 단순
JsonConverter — schemas.enable로 스키마 포함 여부를 결정. Apache Kafka에 포함