기본개념 · 9장
Kafka Connect
Kafka Connect는 외부 시스템과 Kafka 사이의 데이터 이동을 코드 없이 설정으로 처리하는 프레임워크입니다. 이 장에서는 worker · connector · task · converter · transform이라는 다섯 개의 어휘를 정확히 구분하고, 운영에서 가장 자주 물리는 두 지점 — source와 sink가 오프셋을 서로 다른 곳에 저장한다는 사실과 REST API가 409를 돌려주는 이유 — 를 원리부터 정리합니다. CCDAK에서 Kafka Connect는 15% 비중을 차지합니다.
학습 목표
- worker · connector · task · converter · transform의 역할과 경계를 구분할 수 있습니다.
- standalone과 distributed 모드의 차이, 그리고 distributed 모드에서만 가능한 것을 설명할 수 있습니다.
- source 커넥터는
offset.storage.topic에, sink 커넥터는__consumer_offsets에 진행 위치를 저장한다는 비대칭을 설명할 수 있습니다. - SMT가 stateless · 메시지 단위라는 제약과, 그래서 무엇을 할 수 없는지 판단할 수 있습니다.
errors.tolerance와 DLQ 설정 조합을 읽고 어떤 레코드가 어디로 가는지 예측할 수 있습니다.- REST API가
409 Conflict를 반환하는 세 가지 상황과 대응을 설명할 수 있습니다.
Connect가 해결하는 문제
데이터베이스 변경을 Kafka로 보내고, Kafka의 이벤트를 검색 엔진과 데이터 웨어하우스로 내려보내는 코드는 어느 조직에나 있습니다. 그리고 그 코드는 거의 항상 같은 문제를 다시 풉니다 — 어디까지 읽었는지 기록하고, 실패하면 재시도하고, 프로세스가 죽으면 다른 노드에서 이어받고, 스키마를 변환합니다. 공식 문서는 Connect의 목적을 "Kafka와 다른 시스템 사이에서 데이터를 확장성 있고 신뢰성 있게 흘려보내는 도구"로 정의하며, 그 특징으로 커넥터 공통 프레임워크, standalone/distributed 두 모드, REST 인터페이스, 자동 오프셋 관리, 기본 제공되는 분산·확장성, 스트리밍/배치 통합을 듭니다.
Connect를 쓰는 판단 기준은 간단합니다. 데이터를 옮기기만 한다면 Connect, 옮기면서 계산해야 한다면 Kafka Streams입니다. 이 경계는 뒤에서 다루는 SMT의 제약에서 다시 등장합니다.
다섯 개의 명사
| 구성 요소 | 무엇인가 | 단위·수량 |
|---|---|---|
| worker | 커넥터와 태스크를 실제로 실행하는 JVM 프로세스. REST API를 노출하고, distributed 모드에서는 같은 group.id를 가진 워커끼리 하나의 Connect 클러스터를 이룹니다. |
프로세스 단위. 늘리면 처리 용량과 내구성이 함께 올라갑니다. |
| connector | "어디서 어디로 무엇을 옮길지"를 담은 논리적 작업 정의. 직접 데이터를 옮기지 않고, 태스크 설정 목록을 만들어 냅니다. | 이름(name) 하나당 하나. 클러스터 전체에서 이름이 유일해야 합니다. |
| task | 실제로 데이터를 읽고 쓰는 실행 단위(스레드). 워커에 분산 배치됩니다. | tasks.max가 상한. 커넥터가 그만큼의 병렬성을 낼 수 없으면 더 적게 만듭니다. |
| converter | Connect 내부 데이터 표현과 Kafka에 저장되는 바이트열 사이의 변환기. 커넥터와 독립이라 어떤 커넥터에도 어떤 직렬화 형식이든 붙일 수 있습니다. | key / value / header 각각 하나. |
| transform (SMT) | 레코드 하나씩 가볍게 고치거나 라우팅하는 변환. 여러 개를 체인으로 엮습니다. | transforms 목록의 순서대로 적용. |
standalone과 distributed 모드
Connect는 두 가지 실행 모드를 지원합니다. standalone은 모든 작업이 단일 프로세스에서 수행되며 설정이 간단하지만, 공식 문서는 내고장성(fault tolerance) 같은 Connect의 일부 기능을 누리지 못한다고 명시합니다. distributed는 작업 자동 분배, 동적 스케일 인/아웃, 그리고 활성 태스크와 설정·오프셋 커밋 데이터 양쪽에 대한 내고장성을 제공합니다.
| 관점 | standalone | distributed |
|---|---|---|
| 기동 스크립트 | connect-standalone.sh |
connect-distributed.sh |
| 커넥터 등록 방법 | 명령줄에 properties/JSON 파일 전달, 또는 REST API | REST API만 (명령줄로 커넥터 설정을 넘기지 않습니다) |
| source 오프셋 저장 | offset.storage.file.filename — 로컬 파일 |
offset.storage.topic — Kafka 토픽 |
| 커넥터 설정 저장 | 로컬 파일 | config.storage.topic |
| 상태 저장 | 프로세스 메모리 | status.storage.topic |
| 확장 | 불가 (단일 프로세스) | 워커 추가로 확장. 리밸런스로 태스크 재배치 |
| 내고장성 | 없음. 프로세스가 죽으면 멈춥니다 | 있음. 다른 워커가 태스크를 이어받습니다 |
| source 커넥터 exactly-once | 불가 — 공식 문서가 명시적으로 배제합니다 | 가능 (exactly.once.source.support=enabled) |
| 권장 용도 | 로그 파일 수집처럼 워커가 하나여야 의미 있는 경우, 개발·테스트 | 프로덕션 전반 |
모든 워커가 요구하는 설정
# 모든 워커 공통 (standalone/distributed 동일)
bootstrap.servers=kafka-1:9092,kafka-2:9092,kafka-3:9092
key.converter=org.apache.kafka.connect.storage.StringConverter
value.converter=org.apache.kafka.connect.json.JsonConverter
value.converter.schemas.enable=false
plugin.path=/opt/kafka/connect-plugins
# distributed 모드 전용 — 이 클러스터를 식별하는 값. 컨슈머 그룹 ID와 겹치면 안 됩니다
group.id=connect-prod
# 내부 토픽 3종. 사전에 직접 만드는 것이 권장됩니다
config.storage.topic=connect-configs
offset.storage.topic=connect-offsets
status.storage.topic=connect-status
# 브로커 3대 이상이면 3을 씁니다 (기본값도 3)
config.storage.replication.factor=3
offset.storage.replication.factor=3
status.storage.replication.factor=3
offset.storage.partitions=25
status.storage.partitions=5
# REST 리스너 (기본값 http://:8083)
listeners=http://0.0.0.0:8083
내부 토픽 3종
distributed 모드의 Connect 클러스터는 상태를 워커 로컬에 두지 않고 Kafka 토픽 세 개에 둡니다. 그래서 워커가 죽어도 다른 워커가 이어받을 수 있습니다. 공식 문서는 이 세 토픽을 Connect 기동 전에 직접 만들 것을 권장합니다. 자동 생성에 맡기면 브로커의 기본 파티션 수·복제 계수가 적용되어 용도에 맞지 않을 수 있기 때문입니다.
| 설정 | 기본값 | 저장하는 것 | 권장 구성 |
|---|---|---|---|
config.storage.topic |
(필수·기본값 없음) | 커넥터와 태스크의 설정 | 파티션 1개, 복제, cleanup.policy=compact. 설정 순서가 전역으로 하나여야 하므로 단일 파티션입니다. |
offset.storage.topic |
(필수·기본값 없음) | source 커넥터가 소스 시스템에서 어디까지 읽었는지 | 파티션 여러 개(offset.storage.partitions 기본 25), 복제, cleanup.policy=compact |
status.storage.topic |
(필수·기본값 없음) | 커넥터·태스크의 상태와 커넥터가 사용 중인 토픽 목록 | 파티션 여러 개(status.storage.partitions 기본 5), 복제, cleanup.policy=compact |
*.storage.replication.factor |
3 | 세 토픽 각각의 복제 계수 | 브로커가 3대 미만인 개발 환경에서는 낮춰야 기동됩니다. -1은 브로커 기본값 사용. |
커넥터·태스크의 상태 변경은 status.storage.topic에 발행되고 모든 워커가 이 토픽을 구독합니다.
워커가 이 토픽을 비동기로 읽기 때문에 상태 변경이 REST API에 보이기까지 약간의 지연이 있습니다.
커넥터를 만든 직후 RUNNING이 아니라 UNASSIGNED가 보이는 것은 정상입니다.
| 상태 | 의미 | 대응 |
|---|---|---|
UNASSIGNED | 아직 워커에 배정되지 않음 | 리밸런스 완료를 기다립니다 |
RUNNING | 정상 실행 중 | — |
PAUSED | 관리자가 일시 정지. 태스크는 살아 있고 자원도 유지됩니다 | PUT /connectors/{name}/resume |
STOPPED | 커넥터가 정지되어 태스크가 종료되고 자원이 해제됨. 태스크에는 이 상태가 없습니다 | 오프셋 수정은 이 상태에서만 가능합니다 |
FAILED | 예외로 실패. trace에 스택트레이스가 담깁니다 | Connect는 실패한 태스크를 자동 재시작하지 않습니다. restart API를 직접 호출해야 합니다 |
RESTARTING | 재시작 중이거나 곧 재시작될 예정 | — |
source와 sink의 오프셋 저장 위치는 다릅니다
이 장에서 하나만 기억해야 한다면 이것입니다. 같은 "오프셋"이라는 단어를 쓰지만 source와 sink는 완전히 다른 것을 다른 곳에 저장합니다.
| source 커넥터 | sink 커넥터 | |
|---|---|---|
| 데이터 방향 | 외부 시스템 → Kafka | Kafka → 외부 시스템 |
| "오프셋"이 가리키는 것 | 소스 시스템 안의 위치 (파일 바이트 위치, DB의 binlog 좌표, 테이블의 증가 컬럼 값 등). 커넥터가 스스로 정의한 임의 구조입니다. | Kafka 토픽-파티션의 컨슈머 오프셋. 형식이 모든 sink 커넥터에 공통입니다. |
| 저장 위치 (distributed) | offset.storage.topic — Connect의 내부 토픽 |
__consumer_offsets — 일반 컨슈머와 똑같은 경로. Connect 내부 토픽이 아닙니다. |
| 저장 위치 (standalone) | offset.storage.file.filename — 로컬 파일 |
__consumer_offsets (모드와 무관하게 동일) |
| 진행 위치 확인 방법 | GET /connectors/{name}/offsets |
GET /connectors/{name}/offsets 또는 kafka-consumer-groups.sh --describe |
| 컨슈머 그룹 이름 | 해당 없음 (컨슈머가 아닙니다) | 기본 connect-{커넥터 이름} |
offset.storage.topic으로,
sink는 일반 컨슈머와 같은 __consumer_offsets로 진행 위치가 흘러가는 두 경로
GET /connectors/{name}/offsets// FileStreamSourceConnector — partition/offset 구조를 커넥터가 정의합니다
{
"offsets": [
{ "partition": { "filename": "test.txt" },
"offset": { "position": 30 } }
]
}
// FileStreamSinkConnector — 모든 sink 커넥터가 이 공통 형식을 씁니다
{
"offsets": [
{ "partition": { "kafka_topic": "test", "kafka_partition": 0 },
"offset": { "kafka_offset": 5 } },
{ "partition": { "kafka_topic": "test", "kafka_partition": 1 },
"offset": null }
]
}
offset 필드가 null이면 해당 파티션의 오프셋을 초기화합니다.
공식 문서는 source 커넥터의 요청 본문 형식은 커넥터 구현에 따라 다르고,
sink 커넥터는 공통 형식이라는 점을 명시합니다 — 이것도 위 비대칭의 결과입니다.
source 오프셋 커밋 주기
| 설정 | 기본값 | 설명 | 튜닝 포인트 |
|---|---|---|---|
offset.flush.interval.ms |
60000 (1분) |
태스크의 오프셋 커밋을 시도하는 주기 | 줄이면 장애 시 재처리 구간이 짧아지지만 내부 토픽 쓰기가 늘어납니다. |
offset.flush.timeout.ms |
5000 (5초) |
레코드 플러시와 오프셋 커밋을 기다리는 최대 시간. 초과하면 취소하고 다음 시도로 넘깁니다 | 느린 sink에서 커밋 실패 로그가 반복되면 이 값을 올립니다. exactly-once source에는 영향이 없습니다. |
task.shutdown.graceful.timeout.ms |
5000 (5초) |
태스크 전체가 정상 종료되기를 기다리는 총합 시간 (태스크당이 아닙니다) | 태스크 수가 많으면 늘려야 강제 종료를 피할 수 있습니다. |
Converter — Kafka에 무슨 바이트가 쓰이는가
커넥터는 데이터를 Connect의 내부 표현(Struct/Map + Schema)으로 다룹니다.
그것을 Kafka에 저장될 바이트열로 바꾸는 것이 converter입니다.
converter는 커넥터와 독립이라서, 같은 JDBC source 커넥터를 Avro로도 JSON으로도 쓸 수 있습니다.
| 클래스 | 용도 | 비고 |
|---|---|---|
org.apache.kafka.connect.json.JsonConverter |
JSON 직렬화. 스키마를 페이로드에 함께 담을 수 있습니다 | schemas.enable 기본값 true |
org.apache.kafka.connect.storage.StringConverter |
문자열로 변환. 키에 자주 씁니다 | 스키마 개념이 없습니다 |
org.apache.kafka.connect.converters.ByteArrayConverter |
변환하지 않고 바이트를 그대로 통과 | 바이너리 페이로드 파이프스루 |
...converters.{Boolean,Double,Float,Integer,Long,Short}Converter |
원시 타입 단일 값 변환 | 키가 숫자인 파이프라인에서 유용합니다 |
org.apache.kafka.connect.storage.SimpleHeaderConverter |
헤더 값 변환. header.converter의 기본값 |
— |
스키마 봉투가 붙어 나갑니다. 목적지가 순수 JSON을 기대하면 파싱에 실패합니다.
value.converter=org.apache.kafka.connect.json.JsonConverter
# schemas.enable 을 지정하지 않음 → 기본값 true
스키마 없는 순수 JSON이 나갑니다. key/value를 각각 지정합니다.
key.converter=org.apache.kafka.connect.storage.StringConverter
value.converter=org.apache.kafka.connect.json.JsonConverter
value.converter.schemas.enable=false
SMT — Single Message Transform
공식 문서는 transform을 "가벼운 메시지 단위(message-at-a-time) 수정"이라고 정의하며, 데이터 정리와 이벤트 라우팅에 편리하다고 설명합니다. 이 정의에 SMT의 능력과 한계가 모두 들어 있습니다.
기본 제공 SMT
| 이름 | 용도 | 주요 파라미터 |
|---|---|---|
InsertField | 정적 값 또는 레코드 메타데이터로 필드를 추가 | static.field, static.value, topic.field, partition.field, offset.field, timestamp.field |
ReplaceField | 필드를 걸러내거나 이름을 바꿈 | include, exclude, renames |
MaskField | 필드를 타입의 null 값(0, 빈 문자열 등) 또는 지정한 값으로 치환 | fields, replacement |
ValueToKey | value의 일부 필드로 레코드 키를 새로 만듦 | fields |
ExtractField | Struct/Map에서 특정 필드만 뽑아 결과로 삼음 | field, field.syntax.version |
HoistField | 이벤트 전체를 Struct/Map의 단일 필드로 감쌈 | field |
Cast | 필드 또는 키/값 전체를 특정 타입으로 캐스팅 | spec |
TimestampConverter | 타임스탬프 형식 변환 (문자열 ↔ Unix ↔ Date 등) | target.type, field, format, unix.precision |
RegexRouter | 정규식으로 대상 토픽 이름을 변경 | regex, replacement |
TimestampRouter | 원래 토픽명과 타임스탬프로 토픽명을 재구성 (일자별 인덱스 등) | topic.format, timestamp.format |
Filter | 레코드를 이후 처리에서 제거. predicate와 함께 써야 의미가 있습니다 | (없음 — predicate로 조건 지정) |
Flatten | 중첩 구조를 평탄화 | delimiter |
HeaderFrom | 키/값의 필드를 헤더로 복사 또는 이동 | fields, headers, operation |
InsertHeader | 정적 값으로 헤더 추가 | header, value.literal |
DropHeaders | 이름으로 헤더 제거 | headers |
SetSchemaMetadata | 스키마 이름 또는 버전 변경 | schema.name, schema.version |
SMT 클래스 이름에 $Key / $Value 접미어를 붙여 키에 적용할지 값에 적용할지를
고릅니다. 예를 들어 InsertField$Value는 값에, ExtractField$Key는 키에 적용됩니다.
이 접미어를 빼면 클래스를 찾지 못해 커넥터 생성이 실패합니다.
name=local-file-source
connector.class=FileStreamSource
tasks.max=1
file=test.txt
topic=connect-test
# 적용 순서가 곧 이 목록의 순서입니다
transforms=MakeMap, InsertSource
transforms.MakeMap.type=org.apache.kafka.connect.transforms.HoistField$Value
transforms.MakeMap.field=line
transforms.InsertSource.type=org.apache.kafka.connect.transforms.InsertField$Value
transforms.InsertSource.static.field=data_source
transforms.InsertSource.static.value=test-file-source
변환 전
"foo"
"bar"
"hello world"
변환 후
{"line":"foo","data_source":"test-file-source"}
{"line":"bar","data_source":"test-file-source"}
{"line":"hello world","data_source":"test-file-source"}
Predicate — 조건부 SMT
SMT를 조건을 만족하는 레코드에만 적용하려면 predicate를 씁니다.
모든 SMT는 predicate와 negate라는 암묵적 설정을 가지며,
negate=true로 조건을 뒤집을 수 있습니다.
Apache Kafka는 세 개의 predicate를 제공합니다.
| predicate | 매칭 조건 | 파라미터 |
|---|---|---|
TopicNameMatches | 토픽 이름이 Java 정규식과 일치 | pattern |
HasHeaderKey | 지정한 키의 헤더를 가진 레코드 | name |
RecordIsTombstone | value가 null인 tombstone 레코드 | (없음) |
foo 토픽은 전부 버리고, bar가 아닌 토픽에만 ExtractField 적용transforms=Filter,Extract
transforms.Filter.type=org.apache.kafka.connect.transforms.Filter
transforms.Filter.predicate=IsFoo
transforms.Extract.type=org.apache.kafka.connect.transforms.ExtractField$Key
transforms.Extract.field=other_field
transforms.Extract.predicate=IsBar
transforms.Extract.negate=true
predicates=IsFoo,IsBar
predicates.IsFoo.type=org.apache.kafka.connect.transforms.predicates.TopicNameMatches
predicates.IsFoo.pattern=foo
predicates.IsBar.type=org.apache.kafka.connect.transforms.predicates.TopicNameMatches
predicates.IsBar.pattern=bar
에러 처리와 DLQ
공식 문서의 기본 동작 서술은 명확합니다 — 변환(conversion)이나 transform 도중 발생한 오류는 기본적으로 커넥터를 실패시킵니다("fail fast"). 이 동작은 아래 네 설정의 기본값이 만들어 냅니다.
# 실패 시 재시도하지 않음
errors.retry.timeout=0
# 오류와 그 문맥을 로그에 남기지 않음
errors.log.enable=false
# 오류 레코드를 DLQ 토픽에 기록하지 않음
errors.deadletterqueue.topic.name=
# 첫 오류에서 실패
errors.tolerance=none
| 설정 | 기본값 | 설명 | 튜닝 포인트 |
|---|---|---|---|
errors.tolerance |
none |
none이면 어떤 오류든 즉시 태스크 실패. all이면 문제 레코드를 건너뜁니다 |
값은 이 둘뿐입니다. 숫자나 비율을 넣을 수 없습니다. |
errors.retry.timeout |
0 | 실패한 작업을 재시도하는 최대 시간(ms). 0은 재시도 없음, -1은 무한 재시도 |
일시적 네트워크 오류가 잦은 sink에 10분(600000) 정도가 흔한 출발점입니다. |
errors.retry.delay.max.ms |
60000 (1분) |
연속 재시도 사이의 최대 간격. 이 상한에 도달하면 jitter를 더해 thundering herd를 방지합니다 | — |
errors.log.enable |
false | 각 오류와 실패한 작업, 문제 레코드의 토픽·파티션·오프셋을 Connect 애플리케이션 로그에 기록 | 운영에서는 켜 두는 편이 낫습니다. 원인 추적의 출발점입니다. |
errors.log.include.messages |
false | 실패한 레코드의 키·값·헤더까지 로그에 기록 | 민감정보가 로그로 흘러갑니다. 디버깅 기간에만 한시적으로 켜세요. |
errors.deadletterqueue.topic.name |
"" (빈 값) |
DLQ로 쓸 토픽 이름. 비어 있으면 DLQ 기능 자체가 꺼진 것입니다 | 커넥터별로 나누면 원인 분석이 쉽습니다. |
errors.deadletterqueue.topic.replication.factor |
3 | DLQ 토픽이 없을 때 생성할 복제 계수 | 브로커 3대 미만 환경에서는 낮춰야 생성됩니다. |
errors.deadletterqueue.context.headers.enable |
false | 실패 문맥을 헤더로 DLQ 메시지에 붙입니다. 헤더 키는 모두 __connect.errors.로 시작합니다 |
사실상 켜야 합니다. 끄면 DLQ에 원인 없는 페이로드만 쌓입니다. |
# 최대 10분간 재시도, 연속 재시도 간격은 최대 30초
errors.retry.timeout=600000
errors.retry.delay.max.ms=30000
# 오류 문맥을 로그에 남기되 레코드 내용은 제외
errors.log.enable=true
errors.log.include.messages=false
# 오류 문맥을 Kafka 토픽으로 보냄
errors.deadletterqueue.topic.name=my-connector-errors
# 모든 오류를 허용 (실패시키지 않고 보고만 함)
errors.tolerance=all
errors.tolerance 판단을 거쳐 DLQ 토픽으로 가고, 실패 원인이 헤더에 실리는 경로
DLQ 헤더로 원인을 읽는 법
errors.deadletterqueue.context.headers.enable=true일 때 붙는 헤더입니다.
원본 레코드의 헤더와 충돌하지 않도록 모두 __connect.errors. 접두어를 갖습니다.
| 헤더 키 | 담기는 값 |
|---|---|
__connect.errors.topic | 원본 레코드의 토픽 |
__connect.errors.partition | 원본 파티션 |
__connect.errors.offset | 원본 오프셋 |
__connect.errors.connector.name | 커넥터 이름 |
__connect.errors.task.id | 태스크 ID |
__connect.errors.stage | 실패한 단계 (변환 / transform / put 등) |
__connect.errors.class.name | 실패 시점에 실행 중이던 클래스 |
__connect.errors.exception.class.name | 예외 클래스 |
__connect.errors.exception.message | 예외 메시지 |
__connect.errors.exception.stacktrace | 스택트레이스 |
REST API
Connect는 서비스로 실행되는 것을 전제로 하기 때문에, 커넥터 관리를 REST API로 합니다.
이 API는 standalone과 distributed 양쪽에서 모두 제공됩니다.
listeners를 지정하지 않으면 기본값은 http://:8083입니다.
| 메서드 · 경로 | 용도 |
|---|---|
GET / | 워커 버전, git commit ID, 연결된 Kafka 클러스터 ID |
GET /connectors | 활성 커넥터 목록 |
POST /connectors | 커넥터 생성. 본문은 name과 config를 가진 JSON. initial_state로 STOPPED/PAUSED/RUNNING(기본) 지정 가능 |
GET /connectors/{name} | 특정 커넥터 정보 |
GET /connectors/{name}/config | 커넥터 설정 조회 |
PUT /connectors/{name}/config | 설정 전체 교체. 없으면 생성됩니다(멱등 배포에 유용) |
PATCH /connectors/{name}/config | 설정 부분 변경. 본문의 null 값은 해당 키 제거를 뜻합니다 |
GET /connectors/{name}/status | 커넥터·전체 태스크 상태, 배정된 워커 ID, 실패 시 trace |
GET /connectors/{name}/tasks | 실행 중인 태스크와 그 설정 목록 |
GET /connectors/{name}/tasks/{taskid}/status | 개별 태스크 상태 |
PUT /connectors/{name}/pause | 일시 정지 (자원 유지) |
PUT /connectors/{name}/stop | 정지 (태스크 종료·자원 해제). 오프셋 수정의 전제 조건 |
PUT /connectors/{name}/resume | 일시 정지·정지 상태에서 재개 |
POST /connectors/{name}/restart?includeTasks=&onlyFailed= | 커넥터·태스크 재시작. 두 파라미터의 기본값은 모두 false(이전 버전과 동일한 동작) |
POST /connectors/{name}/tasks/{taskId}/restart | 개별 태스크 재시작 — 실패한 태스크 복구의 기본 수단 |
DELETE /connectors/{name} | 커넥터 삭제 (모든 태스크 중지 + 설정 삭제) |
GET /connectors/{name}/topics | 커넥터가 사용 중인 토픽 집합 |
PUT /connectors/{name}/topics/reset | 사용 토픽 집합 초기화 |
GET /connectors/{name}/offsets | 현재 오프셋 조회 (KIP-875) |
DELETE /connectors/{name}/offsets | 오프셋 초기화. STOPPED 상태 필수 |
PATCH /connectors/{name}/offsets | 오프셋을 특정 값으로 변경. STOPPED 상태 필수 |
GET /connector-plugins | 설치된 플러그인 목록. 요청을 처리한 워커만 확인하므로 롤링 업그레이드 중에는 결과가 워커마다 다를 수 있습니다 |
PUT /connector-plugins/{type}/config/validate | 설정 검증. 배포 전 CI에서 유용합니다 |
GET /admin/loggers | 레벨이 명시적으로 지정된 로거 목록 |
PUT /admin/loggers/{name} | 런타임에 로그 레벨 변경 (재시작 불필요, KIP-495) |
409 Conflict가 나는 이유
distributed 모드에서 REST API는 사용자 인터페이스이면서
워커 사이의 클러스터 내부 통신 채널이기도 합니다.
팔로워 워커가 받은 일부 요청은 리더 워커로 전달(forward)됩니다.
409 Conflict는 이 구조와 리밸런스가 만나는 지점에서 나옵니다.
| 원인 | 무슨 상황인가 | 대응 |
|---|---|---|
| 리밸런스 진행 중 | 워커 합류·이탈, 커넥터 추가·설정 변경으로 클러스터가 재구성되는 동안 쓰기성 요청이 들어옴. 공식 문서는 리밸런스 중 태스크 재시작을 시도하면 409를 반환한다고 명시합니다 | 리밸런스 완료 후 재시도. 다만 리밸런스 자체가 커넥터·태스크를 사실상 재시작하므로 재시도가 불필요할 수도 있습니다 |
| 동명 커넥터가 이미 존재 | POST /connectors에 이미 등록된 name을 다시 보냄. 공식 문서는 같은 이름으로 재등록하면 실패한다고 명시합니다 |
GET /connectors로 확인. 설정을 바꾸려는 것이면 PUT /connectors/{name}/config를 쓰세요 |
| 리더가 아닌 워커로 간 요청의 포워딩 실패 | 쓰기성 요청은 리더가 처리해야 하는데, 리더가 아직 정해지지 않았거나 포워딩 대상에 도달할 수 없음. rest.advertised.*가 잘못돼 워커끼리 서로 못 찾는 경우가 대표적입니다 |
rest.advertised.host.name·rest.advertised.port·rest.advertised.listener가 다른 워커에서 실제로 도달 가능한 주소인지 확인 |
CONNECT=http://connect-1:8083
# 1. 커넥터와 태스크 상태를 함께 본다. 커넥터만 보면 안 됩니다
curl -s $CONNECT/connectors/pg-source/status
# 2. FAILED 태스크의 trace 를 읽는다 (실패 원인과 어느 워커였는지)
curl -s $CONNECT/connectors/pg-source/status \
| grep -o '"trace":"[^"]*"' | head -1
# 3. 실패한 태스크만 재시작
curl -s -X POST "$CONNECT/connectors/pg-source/restart?includeTasks=true&onlyFailed=true"
# 4. 원인 추적이 더 필요하면 런타임에 로그 레벨을 올린다 (재시작 불필요)
curl -s -X PUT -H "Content-Type: application/json" \
--data '{"level":"DEBUG"}' \
$CONNECT/admin/loggers/org.apache.kafka.connect.runtime.WorkerSinkTask
# 5. 오프셋을 되돌려야 하면 먼저 STOPPED 로
curl -s -X PUT $CONNECT/connectors/pg-source/stop
curl -s -X DELETE $CONNECT/connectors/pg-source/offsets
curl -s -X PUT $CONNECT/connectors/pg-source/resume
태스크 재분배와 확장 전략
공식 문서에 따르면 리밸런스는 다음 상황에 발생합니다 — 커넥터가 클러스터에 처음 제출될 때, 커넥터가 요구하는 태스크 수가 늘거나 줄 때, 커넥터 설정이 변경될 때, 그리고 워커가 그룹에 추가되거나 제거될 때 (의도적 업그레이드든 장애든).
scheduled.rebalance.max.delay.ms만큼 기다린 뒤
남은 워커로 태스크가 증분 재할당되는 과정
증분 협력적 리밸런스
2.3.0 이전에는 리밸런스마다 클러스터의 모든 커넥터와 태스크를 전부 재배치했습니다.
2.3.0부터 기본값이 증분 협력적 리밸런스(incremental cooperative rebalancing)로 바뀌어,
새로 생기거나 없어지거나 이동해야 하는 태스크만 영향을 받습니다.
나머지 태스크는 멈추지도 재시작하지도 않습니다.
connect.protocol의 4.3 기본값은 sessioned이며,
eager로 두면 옛 전량 재배치 동작으로 돌아갑니다.
| 설정 | 기본값 | 설명 | 튜닝 포인트 |
|---|---|---|---|
connect.protocol |
sessioned |
Connect 프로토콜 호환 모드. eager / compatible / sessioned |
eager는 전량 재배치(2.3 이전 동작)로 되돌립니다. 특별한 이유 없이 바꾸지 마세요. |
scheduled.rebalance.max.delay.ms |
300000 (5분) |
워커가 이탈했을 때 돌아오기를 기다리는 최대 지연. 이 기간에 그 워커의 태스크는 미배정 상태로 남습니다 | 지연 안에 복귀하면 태스크를 그대로 되돌려받습니다. 롤링 재시작 중 불필요한 재배치를 막는 장치지만, 그만큼 데이터가 흐르지 않습니다. 짧은 배포라면 줄이는 것을 검토하세요. |
tasks.max |
1 | 커넥터가 만들 태스크 수의 상한 (커넥터 설정) | 커넥터가 그만큼의 병렬성을 낼 수 없으면 더 적게 만듭니다. |
plugin.discovery |
hybrid_warn |
플러그인 탐색 전략. 워커 기동 시간에 큰 영향을 줍니다 | service_load가 가장 빠르지만 모든 플러그인이 호환되는지 먼저 검증해야 합니다. hybrid_fail로 테스트 환경에서 확인하세요. |
connector.client.config.override.policy |
All |
커넥터가 오버라이드할 수 있는 클라이언트 설정 범위 | Kafka 4.2부터 Allowlist 권장이며 5.0에서 기본값이 됩니다. All이면 커넥터가 sasl.jaas.config 같은 값까지 덮어쓸 수 있습니다. |
확장의 실제 상한
Connect의 병렬성은 세 단계로 제한됩니다.
- 커넥터가 만들 수 있는 태스크 수 —
tasks.max는 상한일 뿐이고, 커넥터가 소스를 그만큼 쪼갤 수 없으면 태스크가 덜 생깁니다. sink 커넥터의 실질 상한은 구독 토픽의 총 파티션 수입니다. - 워커 수 — 태스크는 워커에 배치되므로 워커가 적으면 한 프로세스에 태스크가 몰립니다.
- 목적지 시스템의 수용력 — 대부분의 실제 병목은 여기입니다.
Connect에서의 exactly-once
공식 문서는 sink 커넥터의 exactly-once는 0.11.0부터, source 커넥터는 3.3.0부터 지원한다고 밝히면서, 동시에 "exactly-once 지원은 실행하는 커넥터 종류에 크게 의존한다"고 경고합니다. 워커 설정을 모두 올바르게 맞추더라도, 커넥터 자체가 프레임워크의 기능을 활용하도록 설계되지 않았다면 exactly-once가 성립하지 않습니다.
| 설정 | 위치·기본값 | 설명 |
|---|---|---|
consumer.isolation.level |
워커 · (컨슈머 기본값 read_uncommitted) |
sink의 exactly-once 전제. read_committed로 두어 abort된 트랜잭션 레코드를 무시하게 합니다. 커넥터별로는 consumer.override.isolation.level |
exactly.once.source.support |
워커 · disabled |
source의 프레임워크 레벨 지원. disabled / preparing / enabled. distributed 모드 전용 |
exactly.once.support |
source 커넥터 · requested |
required로 두면 생성 시점에 사전 검사를 강제해, 커넥터가 EOS를 제공할 수 없으면 생성이 실패합니다 |
transaction.boundary |
source 커넥터 · poll |
트랜잭션 경계. poll(배치마다) / interval(transaction.boundary.interval.ms마다) / connector(커넥터가 스스로 정의) |
무엇을 모니터링해야 하는가
Connect는 JMX로 메트릭을 노출합니다. 장애를 조기에 잡는 데 실제로 쓰이는 것들만 골랐습니다. 전체 메트릭 목록과 알림 임계값 정리는 JMX 메트릭 치트시트에 있습니다.
| 메트릭 | MBean | 무엇을 알려 주는가 |
|---|---|---|
status | kafka.connect:type=connector-task-metrics | 태스크 상태. failed가 하나라도 있으면 파이프라인이 부분적으로 멈춘 것입니다 |
running-ratio | kafka.connect:type=connector-task-metrics | 태스크가 running 상태로 보낸 시간의 비율. 1에서 멀어지면 반복 재시작을 의심합니다 |
sink-record-lag-max | kafka.connect:type=sink-task-metrics | sink 태스크가 컨슈머 위치보다 얼마나 뒤처졌는지 (레코드 수) |
put-batch-avg-time-ms | kafka.connect:type=sink-task-metrics | 목적지에 배치를 쓰는 평균 시간. 목적지 병목의 직접 지표입니다 |
poll-batch-avg-time-ms | kafka.connect:type=source-task-metrics | 소스에서 배치를 읽는 평균 시간 |
total-record-errors / total-record-failures | kafka.connect:type=task-error-metrics | 레코드 처리 오류 / 실패 건수 |
total-records-skipped | kafka.connect:type=task-error-metrics | 오류로 건너뛴 레코드 수. errors.tolerance=all일 때 여기가 조용히 오르면 데이터가 사라지고 있다는 뜻입니다 |
deadletterqueue-produce-requests | kafka.connect:type=task-error-metrics | DLQ 쓰기 시도 횟수 |
offset-commit-failure-percentage | kafka.connect:type=connector-task-metrics | 오프셋 커밋 실패 비율. 오르면 offset.flush.timeout.ms를 검토합니다 |
rebalancing / rebalance-avg-time-ms | kafka.connect:type=connect-worker-rebalance-metrics | 리밸런스 중인지, 평균 소요 시간. 계속 1이면 리밸런스 스톰입니다 |
보안에서 놓치기 쉬운 것
Connect 워커의 principal은 최소한 다음 ACL이 필요합니다 —
Connect 클러스터 group.id에 대한 Group Read,
내부 토픽 3종에 대한 Topic Read·Write
(토픽이 아직 없으면 Create까지).
커넥터별로는 source가 대상 토픽에 Write,
sink가 connect-{커넥터명} 그룹에 Read와 구독 토픽에 Read,
그리고 DLQ 토픽에 Write가 필요합니다.
DLQ 권한을 빼먹어 DLQ 쓰기가 실패하는 것이 흔한 초기 설정 실수입니다.
시험 포인트 정리
확인 문제
네 가지 문항 유형(단일 선택 · 복수 선택 · 연결형 · 순서 배열)이 섞여 있습니다. 오프셋 저장 위치와 409 원인은 반드시 맞히고 넘어가세요.
이어서 볼 곳
- 10장 · Kafka Streams와 ksqlDB SMT로 못 하는 집계·조인·윈도우를 어디서 처리할지.
- 예제 8 · Connect CDC 파이프라인 동작하는 source→sink 구성과 오프셋 확인 절차.
- 예제 6 · DLQ + 재시도 패턴 errors.* 조합을 실제로 돌려 보고 DLQ 헤더를 읽습니다.
- Connect 치트시트 REST 엔드포인트 전체 표, SMT 목록, 상태 코드 대응.
- CCDAK · Kafka Connect 도메인 15% 도메인의 압축 정리와 자주 나오는 함정.
- 트러블슈팅 치트시트 증상 → 확인 명령 → 원인 후보 결정 트리.
공식 문서 출처
이 장의 설정 기본값·엔드포인트·클래스명은 모두 아래에서 확인했습니다 (Apache Kafka 4.3 문서 기준).
- Kafka Connect — Overview — Connect의 목적과 특징 6가지
- Running Kafka Connect — standalone/distributed 차이, 워커 필수 설정, 내부 토픽 권장 구성,
producer./consumer.접두어 규칙 - Transformations — SMT 목록과 각 파라미터, predicate 3종,
negate - Connect REST API — 엔드포인트 전체,
listeners기본값http://:8083,admin.listeners, 오프셋 관리 엔드포인트(KIP-875) - Error Reporting in Connect — fail-fast 기본값 4개, DLQ 설정,
errors.log.include.messages의 위험 - Exactly-once support — sink 0.11.0 / source 3.3.0, distributed 전용,
preparing→enabled2단계 롤링 - Connect Administration — 리밸런스 발생 조건, 증분 협력적 리밸런스(2.3.0),
scheduled.rebalance.max.delay.ms5분, 상태 6종, 409 Conflict, pause/stop 차이 - Connect Security — REST API 무인증 경고,
connector.client.config.override.policy, 워커·커넥터별 ACL 표 - Kafka Connect Configs —
offset.storage.partitions25,status.storage.partitions5,connect.protocolsessioned,plugin.discoveryhybrid_warn,offset.flush.* - Sink Connector Configs —
errors.*기본값, DLQ 3종이 sink 전용임 - Source Connector Configs —
exactly.once.support,transaction.boundary,offsets.storage.topic - Connect Monitoring —
sink-record-lag-max,task-error-metrics,rebalancing등 MBean - DeadLetterQueueReporter (Apache Kafka 4.3 소스) —
__connect.errors.*헤더 키 10개