학습 목표

Connect가 해결하는 문제

데이터베이스 변경을 Kafka로 보내고, Kafka의 이벤트를 검색 엔진과 데이터 웨어하우스로 내려보내는 코드는 어느 조직에나 있습니다. 그리고 그 코드는 거의 항상 같은 문제를 다시 풉니다 — 어디까지 읽었는지 기록하고, 실패하면 재시도하고, 프로세스가 죽으면 다른 노드에서 이어받고, 스키마를 변환합니다. 공식 문서는 Connect의 목적을 "Kafka와 다른 시스템 사이에서 데이터를 확장성 있고 신뢰성 있게 흘려보내는 도구"로 정의하며, 그 특징으로 커넥터 공통 프레임워크, standalone/distributed 두 모드, REST 인터페이스, 자동 오프셋 관리, 기본 제공되는 분산·확장성, 스트리밍/배치 통합을 듭니다.

Connect를 쓰는 판단 기준은 간단합니다. 데이터를 옮기기만 한다면 Connect, 옮기면서 계산해야 한다면 Kafka Streams입니다. 이 경계는 뒤에서 다루는 SMT의 제약에서 다시 등장합니다.

다섯 개의 명사

Kafka Connect의 구성 요소와 각각의 책임 범위
구성 요소 무엇인가 단위·수량
worker 커넥터와 태스크를 실제로 실행하는 JVM 프로세스. REST API를 노출하고, distributed 모드에서는 같은 group.id를 가진 워커끼리 하나의 Connect 클러스터를 이룹니다. 프로세스 단위. 늘리면 처리 용량과 내구성이 함께 올라갑니다.
connector "어디서 어디로 무엇을 옮길지"를 담은 논리적 작업 정의. 직접 데이터를 옮기지 않고, 태스크 설정 목록을 만들어 냅니다. 이름(name) 하나당 하나. 클러스터 전체에서 이름이 유일해야 합니다.
task 실제로 데이터를 읽고 쓰는 실행 단위(스레드). 워커에 분산 배치됩니다. tasks.max가 상한. 커넥터가 그만큼의 병렬성을 낼 수 없으면 더 적게 만듭니다.
converter Connect 내부 데이터 표현과 Kafka에 저장되는 바이트열 사이의 변환기. 커넥터와 독립이라 어떤 커넥터에도 어떤 직렬화 형식이든 붙일 수 있습니다. key / value / header 각각 하나.
transform (SMT) 레코드 하나씩 가볍게 고치거나 라우팅하는 변환. 여러 개를 체인으로 엮습니다. transforms 목록의 순서대로 적용.
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 클러스터에서 커넥터 정의가 태스크로 쪼개져 분산 배치되고, 리더 워커가 설정 변경과 리밸런스를 조정하는 구조

standalone과 distributed 모드

Connect는 두 가지 실행 모드를 지원합니다. standalone은 모든 작업이 단일 프로세스에서 수행되며 설정이 간단하지만, 공식 문서는 내고장성(fault tolerance) 같은 Connect의 일부 기능을 누리지 못한다고 명시합니다. distributed는 작업 자동 분배, 동적 스케일 인/아웃, 그리고 활성 태스크와 설정·오프셋 커밋 데이터 양쪽에 대한 내고장성을 제공합니다.

standalone과 distributed 모드 비교 (Apache Kafka 4.3 문서 기준)
관점 standalone distributed
기동 스크립트 connect-standalone.sh connect-distributed.sh
커넥터 등록 방법 명령줄에 properties/JSON 파일 전달, 또는 REST API REST API만 (명령줄로 커넥터 설정을 넘기지 않습니다)
source 오프셋 저장 offset.storage.file.filename로컬 파일 offset.storage.topicKafka 토픽
커넥터 설정 저장 로컬 파일 config.storage.topic
상태 저장 프로세스 메모리 status.storage.topic
확장 불가 (단일 프로세스) 워커 추가로 확장. 리밸런스로 태스크 재배치
내고장성 없음. 프로세스가 죽으면 멈춥니다 있음. 다른 워커가 태스크를 이어받습니다
source 커넥터 exactly-once 불가 — 공식 문서가 명시적으로 배제합니다 가능 (exactly.once.source.support=enabled)
권장 용도 로그 파일 수집처럼 워커가 하나여야 의미 있는 경우, 개발·테스트 프로덕션 전반

모든 워커가 요구하는 설정

connect-distributed.properties — distributed 모드의 최소 골격
# 모든 워커 공통 (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 기동 전에 직접 만들 것을 권장합니다. 자동 생성에 맡기면 브로커의 기본 파티션 수·복제 계수가 적용되어 용도에 맞지 않을 수 있기 때문입니다.

Connect 내부 토픽과 권장 구성. 기본값은 Apache Kafka 4.3 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가 보이는 것은 정상입니다.

커넥터·태스크가 가질 수 있는 상태 (Apache Kafka 4.3 Connect 문서 기준)
상태의미대응
UNASSIGNED아직 워커에 배정되지 않음리밸런스 완료를 기다립니다
RUNNING정상 실행 중
PAUSED관리자가 일시 정지. 태스크는 살아 있고 자원도 유지됩니다PUT /connectors/{name}/resume
STOPPED커넥터가 정지되어 태스크가 종료되고 자원이 해제됨. 태스크에는 이 상태가 없습니다오프셋 수정은 이 상태에서만 가능합니다
FAILED예외로 실패. trace에 스택트레이스가 담깁니다Connect는 실패한 태스크를 자동 재시작하지 않습니다. restart API를 직접 호출해야 합니다
RESTARTING재시작 중이거나 곧 재시작될 예정

source와 sink의 오프셋 저장 위치는 다릅니다

이 장에서 하나만 기억해야 한다면 이것입니다. 같은 "오프셋"이라는 단어를 쓰지만 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-{커넥터 이름}
source 와 sink 의 오프셋 저장 위치 — 비대칭입니다 소스 커넥터와 싱크 커넥터가 진행 위치를 어디에 저장하는지 좌우로 비교한 그림입니다. 왼쪽 소스 커넥터는 데이터베이스나 파일 같은 외부 시스템에서 읽어 Kafka 토픽에 씁니다. 이때 어디까지 읽었는지는 커넥터가 스스로 정의한 구조로 표현하며, Connect 의 내부 토픽인 offset.storage.topic 에 저장됩니다. 소스 쪽에는 컨슈머가 없으므로 __consumer_offsets 와는 아무 관계가 없습니다. 커밋 주기는 워커 설정 offset.flush.interval.ms 기본값 60000 밀리초입니다. 오른쪽 싱크 커넥터는 Kafka 토픽을 일반 KafkaConsumer 로 읽어 외부 시스템에 씁니다. 컨슈머 그룹 id 는 connect 하이픈 커넥터 이름 형식이며, 오프셋은 일반 컨슈머와 똑같이 __consumer_offsets 에 저장됩니다. 그래서 kafka-consumer-groups 명령으로 lag 을 확인할 수 있습니다. 싱크 커넥터는 offset.storage.topic 을 쓰지 않습니다. source vs sink 오프셋 저장 위치 — 같은 Connect 인데 저장소가 다릅니다 SOURCE 커넥터 — 외부 → Kafka 외부 시스템 (DB · 파일 · API) SourceTask.poll() 위치 표현은 커넥터가 정의: {file, position} Kafka 대상 토픽 (데이터) 오프셋 저장 위치 offset.storage.topic Connect 내부 토픽 · offset.flush.interval.ms 컨슈머가 없습니다 → __consumer_offsets 무관 SINK 커넥터 — Kafka → 외부 Kafka 소스 토픽 일반 KafkaConsumer 그대로 사용 group.id = connect-{커넥터 이름} SinkTask.put() → 외부 시스템 오프셋 저장 위치 __consumer_offsets 일반 컨슈머와 완전히 동일합니다 kafka-consumer-groups 로 lag 확인 가능 한 줄 요약: source → offset.storage.topic · sink → __consumer_offsets. 뒤집어 외우면 그대로 틀립니다. 싱크 오프셋을 사람이 손으로 offset.storage.topic 에서 찾으려 해도 없습니다. 반대도 마찬가지입니다.
source와 sink의 오프셋 저장 위치 비대칭 — source는 Connect의 offset.storage.topic으로, sink는 일반 컨슈머와 같은 __consumer_offsets로 진행 위치가 흘러가는 두 경로
source(파일)와 sink(Kafka)의 오프셋 응답 형식 차이 — 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 오프셋 커밋 주기

source 오프셋 플러시 관련 워커 설정 (Apache Kafka 4.3 기본값)
설정기본값설명튜닝 포인트
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으로도 쓸 수 있습니다.

Apache Kafka 배포판에 포함된 converter 클래스
클래스용도비고
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을 기대하면 파싱에 실패합니다.

worker.properties
value.converter=org.apache.kafka.connect.json.JsonConverter
# schemas.enable 을 지정하지 않음 → 기본값 true

스키마 없는 순수 JSON이 나갑니다. key/value를 각각 지정합니다.

worker.properties
key.converter=org.apache.kafka.connect.storage.StringConverter
value.converter=org.apache.kafka.connect.json.JsonConverter
value.converter.schemas.enable=false
converter 와 SMT 체인의 위치 — source 와 sink 에서 순서가 반대입니다 소스 커넥터와 싱크 커넥터에서 SMT 체인과 converter 가 놓이는 순서를 위아래로 비교한 그림입니다. 소스는 SourceTask 가 만든 레코드가 SMT 체인을 지나고 그 다음 converter 가 바이트로 직렬화해 Kafka 토픽에 씁니다. 싱크는 Kafka 토픽에서 읽은 바이트를 converter 가 먼저 Connect 데이터 모델로 역직렬화하고 그 다음 SMT 체인을 지나 SinkTask 로 전달됩니다. 즉 순서가 서로 반대입니다. 외우는 규칙은 하나입니다. converter 는 항상 Kafka 쪽에 붙고, SMT 는 항상 커넥터 쪽에 붙습니다. SMT 는 Connect 의 내부 데이터 모델인 Struct 와 Schema 위에서 동작하기 때문에 바이트 상태에서는 실행될 수 없습니다. 실패 단계 이름은 KEY_CONVERTER, VALUE_CONVERTER, TRANSFORMATION 처럼 DLQ 헤더의 stage 값으로 남습니다. converter 와 SMT — 순서가 서로 반대입니다 SOURCE — 외부 → Kafka SourceTask 외부 레코드 SMT 체인 transforms 순서대로 Converter Struct → 바이트 Kafka 토픽 직렬화된 레코드 SINK — Kafka → 외부 Kafka 토픽 직렬화된 레코드 Converter 바이트 → Struct SMT 체인 transforms 순서대로 SinkTask 외부로 쓸 레코드 규칙 하나로 외우세요 Converter 는 언제나 Kafka 쪽에, SMT 는 언제나 커넥터 쪽에 붙습니다. 그래서 순서가 뒤집혀 보입니다. SMT 는 Connect 내부 데이터 모델(Struct · Schema)에서 동작하므로 바이트 상태에서는 실행될 수 없습니다. 실패 단계는 KEY_CONVERTER · VALUE_CONVERTER · TRANSFORMATION 처럼 DLQ 헤더에 남습니다 (D-083).
converter와 SMT 체인의 적용 위치 — source는 SMT를 거친 뒤 converter가 직렬화해 Kafka에 쓰고, sink는 Kafka에서 읽어 converter가 역직렬화한 뒤 SMT를 거쳐 목적지에 씁니다

SMT — Single Message Transform

공식 문서는 transform을 "가벼운 메시지 단위(message-at-a-time) 수정"이라고 정의하며, 데이터 정리와 이벤트 라우팅에 편리하다고 설명합니다. 이 정의에 SMT의 능력과 한계가 모두 들어 있습니다.

기본 제공 SMT

Apache Kafka 4.3에 포함된 SMT. 클래스 목록으로 확인하세요 — 세는 기준(predicate 포함 여부 등)에 따라 개수가 달라집니다. 클래스는 모두 org.apache.kafka.connect.transforms 패키지입니다.
이름 용도 주요 파라미터
InsertField정적 값 또는 레코드 메타데이터로 필드를 추가static.field, static.value, topic.field, partition.field, offset.field, timestamp.field
ReplaceField필드를 걸러내거나 이름을 바꿈include, exclude, renames
MaskField필드를 타입의 null 값(0, 빈 문자열 등) 또는 지정한 값으로 치환fields, replacement
ValueToKeyvalue의 일부 필드로 레코드 키를 새로 만듦fields
ExtractFieldStruct/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는 키에 적용됩니다. 이 접미어를 빼면 클래스를 찾지 못해 커넥터 생성이 실패합니다.

connect-file-source.properties — 공식 문서의 SMT 체인 예시 (HoistField → InsertField)
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
변환 전(문자열)과 변환 후(JSON) 결과
변환 전
"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는 predicatenegate라는 암묵적 설정을 가지며, negate=true로 조건을 뒤집을 수 있습니다. Apache Kafka는 세 개의 predicate를 제공합니다.

기본 제공 predicate (org.apache.kafka.connect.transforms.predicates 패키지)
predicate매칭 조건파라미터
TopicNameMatches토픽 이름이 Java 정규식과 일치pattern
HasHeaderKey지정한 키의 헤더를 가진 레코드name
RecordIsTombstonevalue가 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"). 이 동작은 아래 네 설정의 기본값이 만들어 냅니다.

fail-fast 기본 동작과 동등한 설정 (공식 문서 인용)
# 실패 시 재시도하지 않음
errors.retry.timeout=0
# 오류와 그 문맥을 로그에 남기지 않음
errors.log.enable=false
# 오류 레코드를 DLQ 토픽에 기록하지 않음
errors.deadletterqueue.topic.name=
# 첫 오류에서 실패
errors.tolerance=none
Connect 에러 처리 설정 (Apache Kafka 4.3 커넥터 설정 기준). DLQ 관련 3개는 sink 커넥터 전용입니다.
설정기본값설명튜닝 포인트
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에 원인 없는 페이로드만 쌓입니다.
공식 문서 예시 — 재시도 + 로깅 + 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
Connect DLQ 흐름 — errors.tolerance 와 dead letter queue 싱크 커넥터의 처리 경로와 오류 처리 분기를 그린 그림입니다. 경로는 Kafka 소스 토픽, converter, SMT 체인, SinkTask, 외부 시스템 순서입니다. converter, SMT, SinkTask 세 단계 어디서든 오류가 날 수 있습니다. 오류가 나면 errors.tolerance 설정에 따라 갈립니다. 기본값 none 이면 첫 오류에서 태스크가 FAILED 상태로 멈추고 REST API 로 restart 해야 다시 돕니다. all 로 두면 오류 레코드를 건너뛰고 계속 처리합니다. 이때 errors.log.enable 을 true 로 하면 로그에 남고, errors.deadletterqueue.topic.name 을 지정하면 문제 레코드가 그 DLQ 토픽으로 전송됩니다. errors.deadletterqueue.context.headers.enable 을 true 로 하면 실패 원인이 헤더로 함께 들어갑니다. 헤더 이름은 모두 __connect.errors. 로 시작하며 topic, partition, offset, connector.name, task.id, stage, class.name, exception.class.name, exception.message, exception.stacktrace 가 있습니다. DLQ 설정은 싱크 커넥터에만 있습니다. 소스 커넥터에는 DLQ 설정이 없습니다. DLQ 흐름 — 기본값은 “멈춤”입니다 Kafka 소스 토픽 Converter 바이트 → Struct SMT 체인 transforms SinkTask put() 외부 시스템 ← 이 세 단계에서 오류 발생 errors.tolerance ? errors.tolerance = none 기본값 첫 오류에서 태스크가 FAILED 로 멈춥니다. Connect 는 실패한 태스크를 자동 재시작하지 않습니다. POST /connectors/{name}/tasks/{id}/restart errors.tolerance = all 오류 레코드를 건너뛰고 계속 처리합니다. errors.log.enable=true → 로그에 기록 errors.deadletterqueue.topic.name DLQ 토픽 — 원본 레코드 + 실패 원인 헤더 (…context.headers.enable=true) __connect.errors.topic .partition .offset .connector.name .task.id .stage .class.name .exception.class.name .exception.message .exception.stacktrace DLQ 토픽은 자동 생성됩니다 (RF 기본 3) DLQ 설정은 싱크 커넥터에만 있습니다. 소스 커넥터에는 errors.deadletterqueue.* 가 없습니다.
DLQ 흐름 — 변환·SMT·sink 쓰기 각 단계에서 발생한 오류가 errors.tolerance 판단을 거쳐 DLQ 토픽으로 가고, 실패 원인이 헤더에 실리는 경로

DLQ 헤더로 원인을 읽는 법

errors.deadletterqueue.context.headers.enable=true일 때 붙는 헤더입니다. 원본 레코드의 헤더와 충돌하지 않도록 모두 __connect.errors. 접두어를 갖습니다.

DLQ 오류 문맥 헤더 (Apache Kafka 4.3 DeadLetterQueueReporter 기준)
헤더 키담기는 값
__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입니다.

주요 REST 엔드포인트 (Apache Kafka 4.3 Connect 문서 기준)
메서드 · 경로용도
GET /워커 버전, git commit ID, 연결된 Kafka 클러스터 ID
GET /connectors활성 커넥터 목록
POST /connectors커넥터 생성. 본문은 nameconfig를 가진 JSON. initial_stateSTOPPED/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 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

태스크 재분배와 확장 전략

공식 문서에 따르면 리밸런스는 다음 상황에 발생합니다 — 커넥터가 클러스터에 처음 제출될 때, 커넥터가 요구하는 태스크 수가 늘거나 줄 때, 커넥터 설정이 변경될 때, 그리고 워커가 그룹에 추가되거나 제거될 때 (의도적 업그레이드든 장애든).

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 프로토콜은 리밸런스마다 모든 태스크를 멈췄습니다.
task 재분배 — 워커가 이탈했을 때 scheduled.rebalance.max.delay.ms만큼 기다린 뒤 남은 워커로 태스크가 증분 재할당되는 과정

증분 협력적 리밸런스

2.3.0 이전에는 리밸런스마다 클러스터의 모든 커넥터와 태스크를 전부 재배치했습니다. 2.3.0부터 기본값이 증분 협력적 리밸런스(incremental cooperative rebalancing)로 바뀌어, 새로 생기거나 없어지거나 이동해야 하는 태스크만 영향을 받습니다. 나머지 태스크는 멈추지도 재시작하지도 않습니다. connect.protocol의 4.3 기본값은 sessioned이며, eager로 두면 옛 전량 재배치 동작으로 돌아갑니다.

리밸런스·확장 관련 워커 설정 (Apache Kafka 4.3 기본값)
설정기본값설명튜닝 포인트
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의 병렬성은 세 단계로 제한됩니다.

  1. 커넥터가 만들 수 있는 태스크 수tasks.max는 상한일 뿐이고, 커넥터가 소스를 그만큼 쪼갤 수 없으면 태스크가 덜 생깁니다. sink 커넥터의 실질 상한은 구독 토픽의 총 파티션 수입니다.
  2. 워커 수 — 태스크는 워커에 배치되므로 워커가 적으면 한 프로세스에 태스크가 몰립니다.
  3. 목적지 시스템의 수용력 — 대부분의 실제 병목은 여기입니다.

Connect에서의 exactly-once

공식 문서는 sink 커넥터의 exactly-once는 0.11.0부터, source 커넥터는 3.3.0부터 지원한다고 밝히면서, 동시에 "exactly-once 지원은 실행하는 커넥터 종류에 크게 의존한다"고 경고합니다. 워커 설정을 모두 올바르게 맞추더라도, 커넥터 자체가 프레임워크의 기능을 활용하도록 설계되지 않았다면 exactly-once가 성립하지 않습니다.

Connect 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 메트릭 치트시트에 있습니다.

Connect 핵심 JMX 메트릭 (Apache Kafka 4.3 모니터링 문서 기준)
메트릭MBean무엇을 알려 주는가
statuskafka.connect:type=connector-task-metrics태스크 상태. failed가 하나라도 있으면 파이프라인이 부분적으로 멈춘 것입니다
running-ratiokafka.connect:type=connector-task-metrics태스크가 running 상태로 보낸 시간의 비율. 1에서 멀어지면 반복 재시작을 의심합니다
sink-record-lag-maxkafka.connect:type=sink-task-metricssink 태스크가 컨슈머 위치보다 얼마나 뒤처졌는지 (레코드 수)
put-batch-avg-time-mskafka.connect:type=sink-task-metrics목적지에 배치를 쓰는 평균 시간. 목적지 병목의 직접 지표입니다
poll-batch-avg-time-mskafka.connect:type=source-task-metrics소스에서 배치를 읽는 평균 시간
total-record-errors / total-record-failureskafka.connect:type=task-error-metrics레코드 처리 오류 / 실패 건수
total-records-skippedkafka.connect:type=task-error-metrics오류로 건너뛴 레코드 수. errors.tolerance=all일 때 여기가 조용히 오르면 데이터가 사라지고 있다는 뜻입니다
deadletterqueue-produce-requestskafka.connect:type=task-error-metricsDLQ 쓰기 시도 횟수
offset-commit-failure-percentagekafka.connect:type=connector-task-metrics오프셋 커밋 실패 비율. 오르면 offset.flush.timeout.ms를 검토합니다
rebalancing / rebalance-avg-time-mskafka.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 원인은 반드시 맞히고 넘어가세요.

공식 문서 출처

이 장의 설정 기본값·엔드포인트·클래스명은 모두 아래에서 확인했습니다 (Apache Kafka 4.3 문서 기준).