initial_state로 STOPPED / PAUSED / RUNNING(기본)을 지정할 수 있습니다
GET
/connectors/{connector}
커넥터 상세 (설정 + 태스크 목록)
GET
/connectors/{connector}/config
커넥터 설정만 조회
백업·비교에 유용합니다
PUT
/connectors/{connector}/config
커넥터 생성 또는 전체 재설정
멱등합니다. GitOps 배포에 적합. 단 설정 전체를 보내야 하며 빠진 키는 기본값으로 돌아갑니다
PATCH
/connectors/{connector}/config
커넥터 설정을 부분 변경
일부 키만 바꿀 때. PUT과 달리 나머지를 보존합니다
DELETE
/connectors/{connector}
커넥터 삭제
태스크가 정지되고 설정이 제거됩니다. 활성 토픽 정보도 함께 삭제됩니다
상태 · 생애주기
GET
/connectors/{connector}/status
커넥터와 각 태스크의 상태·담당 워커·trace
장애 대응의 출발점.trace에 스택트레이스가 들어 있습니다
POST
/connectors/{connector}/restart
커넥터 재시작
?includeTasks=true&onlyFailed=true가 실무 기본값입니다. 파라미터를 안 주면 커넥터만 재시작합니다
GET
/connectors/{connector}/tasks
태스크 목록과 각 태스크의 설정
파티션이 어떻게 분배됐는지 확인
GET
/connectors/{connector}/tasks/{task}/status
개별 태스크 상태
POST
/connectors/{connector}/tasks/{task}/restart
개별 태스크 재시작
리밸런스 중이면 409가 반환됩니다
PUT
/connectors/{connector}/pause
일시 정지 — 태스크는 유휴 상태로 남습니다
자원을 계속 점유합니다. 재개가 빠릅니다. 이미 paused면 그대로 성공합니다
PUT
/connectors/{connector}/resume
일시 정지 해제
pause 상태는 영속이라 클러스터를 재시작해도 유지됩니다
PUT
/connectors/{connector}/stop
정지 — 태스크를 완전히 종료하고 자원을 반납
3.5.0 도입. 오프셋을 수정하려면 이 상태여야 합니다. 재개는 pause보다 느립니다
오프셋 관리
GET
/connectors/{connector}/offsets
현재 오프셋 조회
source는 커넥터 고유 파티션, sink는 Kafka 토픽-파티션 형태로 나옵니다
PATCH
/connectors/{connector}/offsets
오프셋 변경
커넥터가 STOPPED 상태여야 합니다
DELETE
/connectors/{connector}/offsets
오프셋 초기화 — 처음부터 다시 읽기
커넥터가 STOPPED 상태여야 합니다. 되돌릴 수 없습니다
토픽 추적
GET
/connectors/{connector}/topics
이 커넥터가 실제로 사용 중인 토픽 목록
status.storage.topic에 기록된 정보를 읽습니다. topic.tracking.enable이 필요
PUT
/connectors/{connector}/topics/reset
활성 토픽 목록 초기화
topic.tracking.allow.reset=false면 거부됩니다
가장 자주 쓰는 호출 — 그대로 붙여 쓰세요
export CW=http://localhost:8083
# 1) 전체 커넥터의 상태를 한 번에 (가장 유용한 한 줄)
curl -s "$CW/connectors?expand=status" | jq -r '
to_entries[] | "\(.value.status.connector.state)\t\(.key)\t태스크: " +
([.value.status.tasks[].state] | join(","))'
# 출력 예
# RUNNING mysql-source 태스크: RUNNING,RUNNING
# FAILED s3-sink 태스크: FAILED,RUNNING
# 2) 실패한 커넥터만 골라내기
curl -s "$CW/connectors?expand=status" | jq -r '
to_entries[]
| select(.value.status.connector.state == "FAILED"
or any(.value.status.tasks[]; .state == "FAILED"))
| .key'
# 3) 실패 원인 확인 — trace 에 스택트레이스가 들어 있습니다
curl -s "$CW/connectors/s3-sink/status" | jq -r '.tasks[] | select(.state=="FAILED") | .trace'
# 4) 실패한 태스크만 재시작 (Connect 는 자동 재시작하지 않습니다)
curl -s -X POST "$CW/connectors/s3-sink/restart?includeTasks=true&onlyFailed=true"
커넥터 생성 · 설정 검증 · 변경
# 배포 전 검증 — 실제로 만들지 않고 오류만 확인합니다
curl -s -X PUT "$CW/connector-plugins/FileStreamSinkConnector/config/validate" \
-H "Content-Type: application/json" \
-d '{
"connector.class": "org.apache.kafka.connect.file.FileStreamSinkConnector",
"tasks.max": "1",
"topics": "orders",
"file": "/tmp/orders.sink.txt"
}' | jq '.error_count, [.configs[] | select(.value.errors | length > 0) | {name: .value.name, errors: .value.errors}]'
# 생성 (POST) — name 과 config 를 감싼 형태
curl -s -X POST "$CW/connectors" \
-H "Content-Type: application/json" \
-d '{
"name": "orders-file-sink",
"config": {
"connector.class": "org.apache.kafka.connect.file.FileStreamSinkConnector",
"tasks.max": "1",
"topics": "orders",
"file": "/tmp/orders.sink.txt"
}
}'
# 멱등 생성/갱신 (PUT) — config 만 보냅니다. GitOps 에 적합합니다.
# 주의: 설정 "전체"를 보내야 합니다. 빠진 키는 기본값으로 돌아갑니다.
curl -s -X PUT "$CW/connectors/orders-file-sink/config" \
-H "Content-Type: application/json" \
-d '{
"connector.class": "org.apache.kafka.connect.file.FileStreamSinkConnector",
"tasks.max": "2",
"topics": "orders",
"file": "/tmp/orders.sink.txt"
}'
# 일부만 바꾸기 (PATCH) — 나머지는 보존됩니다
curl -s -X PATCH "$CW/connectors/orders-file-sink/config" \
-H "Content-Type: application/json" \
-d '{ "tasks.max": "4" }'
오프셋 초기화 — STOPPED 상태가 전제입니다
# 1) 반드시 stop (pause 로는 안 됩니다)
curl -s -X PUT "$CW/connectors/mysql-source/stop"
# 2) STOPPED 확인
curl -s "$CW/connectors/mysql-source/status" | jq -r '.connector.state'
# 3) 현재 오프셋 확인 — 되돌릴 수 있게 저장해 두세요
curl -s "$CW/connectors/mysql-source/offsets" | tee mysql-source-offsets.bak.json | jq
# 4) 초기화 (되돌릴 수 없습니다)
curl -s -X DELETE "$CW/connectors/mysql-source/offsets"
# 5) 재개
curl -s -X PUT "$CW/connectors/mysql-source/resume"
재시작 없이 디버그 로그 켜기
curl -s -X PUT "$CW/admin/loggers/org.apache.kafka.connect.runtime.WorkerSinkTask" \
-H "Content-Type: application/json" -d '{"level": "DEBUG"}'
# 현재 명시적으로 설정된 로거 확인
curl -s "$CW/admin/loggers" | jq
# 조사 후 되돌리기
curl -s -X PUT "$CW/admin/loggers/org.apache.kafka.connect.runtime.WorkerSinkTask" \
-H "Content-Type: application/json" -d '{"level": "INFO"}'
커넥터 상태와 대응
커넥터·태스크 상태 — 정의는 Apache Kafka 4.3.1 Connect Administration 문서 기준
상태
의미
태스크에도 있는가
대응
UNASSIGNED
아직 워커에 배정되지 않았습니다
O
리밸런스가 끝나면 해소됩니다. 오래 지속되면 워커 수·tasks.max·워커 로그 확인
RUNNING
정상 동작 중
O
—
PAUSED
관리자가 일시 정지시켰습니다. 태스크는 유휴 상태로 남아 자원을 점유합니다
O
PUT /resume. 태스크가 처리 중이던 작업을 끝내야 하므로 전이에 시간이 걸릴 수 있습니다. 실패한 태스크는 재시작하기 전까지 PAUSED로 전이하지 않습니다
STOPPED
커넥터가 정지되어 태스크가 완전히 종료되고 자원이 반납되었습니다
X
태스크가 아예 없어지므로 status API에 태스크가 보이지 않습니다. 오프셋 수정은 이 상태에서만 가능합니다
FAILED
예외가 발생해 실패했습니다. 예외는 status 출력의 trace에 담깁니다
O
Connect는 실패한 태스크를 자동으로 재시작하지 않습니다.trace를 읽고 원인을 고친 뒤 POST /restart?includeTasks=true&onlyFailed=true
RESTARTING
재시작 중이거나 곧 재시작될 예정입니다
O
기다립니다
pause vs stop
두 정지 방식의 차이
관점
pause
stop (3.5.0+)
태스크
유휴 상태로 남습니다
완전히 종료됩니다
자원 점유
계속 점유 (커넥션·스레드)
반납
재개 속도
빠릅니다
느립니다 (태스크 재생성)
오프셋 수정
불가
가능
status의 태스크
보입니다 (PAUSED)
보이지 않습니다
지속성
영속 — 클러스터 재시작에도 유지
영속
언제 쓰는가
대상 시스템 점검처럼 곧 재개할 때
장기 정지, 오프셋 조작, 자원 회수
409 Conflict의 원인
409 Conflict의 원인 3가지
원인
어떤 요청에서
확인
대응
① 리밸런스 진행 중
태스크·커넥터 재시작
워커 로그의 rebalance 메시지. GET /connectors?expand=status로 담당 워커 변화 확인
완료를 기다린 뒤 상태를 재확인하세요. 리밸런스 자체가 태스크를 재시작하므로 재시도가 불필요할 수 있습니다(공식 문서 명시)
② 같은 이름의 커넥터가 이미 존재
POST /connectors
GET /connectors로 이름 중복 확인
멱등하게 만들려면 POST 대신 PUT /connectors/{name}/config를 쓰세요. 배포 스크립트에서 특히 중요합니다
③ 리더 워커로의 요청 전달 실패
쓰기 성격의 모든 요청
rest.advertised.host.name · rest.advertised.port가 다른 워커에서 도달 가능한지
Connect는 쓰기 요청을 리더 워커로 전달(forward)합니다. 광고 주소가 잘못되면 전달이 실패합니다 — 컨테이너 환경의 흔한 원인
409를 견디는 재시작 스크립트
#!/usr/bin/env bash
set -euo pipefail
CW="${CW:-http://localhost:8083}"
CONNECTOR="$1"
for attempt in 1 2 3 4 5; do
code=$(curl -s -o /dev/null -w '%{http_code}' -X POST \
"$CW/connectors/$CONNECTOR/restart?includeTasks=true&onlyFailed=true")
if [ "$code" = "200" ] || [ "$code" = "204" ]; then
echo "재시작 요청 성공"
break
fi
if [ "$code" = "409" ]; then
# 리밸런스 중입니다. 리밸런스가 이미 태스크를 재시작했을 수 있으므로
# 재시도 전에 상태를 먼저 확인합니다.
state=$(curl -s "$CW/connectors/$CONNECTOR/status" | jq -r '.connector.state')
if [ "$state" = "RUNNING" ]; then
echo "리밸런스가 이미 복구했습니다 — 재시작 불필요"
exit 0
fi
echo "409 (리밸런스 중) — ${attempt}회, 대기 후 재시도"
sleep $((attempt * 10))
continue
fi
echo "예상치 못한 응답: $code" >&2
exit 1
done
# 최종 상태 확인
curl -s "$CW/connectors/$CONNECTOR/status" | jq '{connector: .connector.state, tasks: [.tasks[].state]}'
SMT (Single Message Transform)
Apache Kafka에 내장된 SMT 목록입니다.
대부분 $Key와 $Value 두 변형을 가지며,
레코드의 키에 적용할지 값에 적용할지를 클래스명으로 구분합니다
(예: org.apache.kafka.connect.transforms.Cast$Value).
내장 SMT — Apache Kafka 4.3.1 생성 문서(connect_transforms) 기준
이름
용도
주요 파라미터
메모
Cast
필드 또는 전체 키·값의 타입을 변환
spec, replace.null.with.default
정수·부동소수·불리언·문자열 간 변환. binary → string(base64)만 가능. $Key/$Value
{
"transforms": "dropTombstones",
"transforms.dropTombstones.type": "org.apache.kafka.connect.transforms.Filter",
"transforms.dropTombstones.predicate": "isTombstone",
"predicates": "isTombstone",
"predicates.isTombstone.type": "org.apache.kafka.connect.transforms.predicates.RecordIsTombstone",
"//": "negate 를 true 로 두면 반대로 tombstone 만 남깁니다",
"//negate": "transforms.dropTombstones.negate: true"
}
내부 토픽 3개
분산 모드 워커는 세 개의 내부 토픽에 상태를 저장합니다.
이 토픽들의 설정을 잘못 두는 것이 Connect 운영 사고의 가장 흔한 원인입니다.
Connect 내부 토픽 — 기본값은 Apache Kafka 4.3.1 Connect 설정 문서 기준
설정
파티션 기본값
RF 기본값
필수 조건
config.storage.topic
반드시 1개 (설정 항목 없음)
config.storage.replication.factor = 3
파티션 1개 + cleanup.policy=compact. 순서가 의미를 가지므로 파티션을 늘리면 설정이 깨집니다
offset.storage.topic
offset.storage.partitions = 25
offset.storage.replication.factor = 3
cleanup.policy=compact. source 커넥터 전용입니다
status.storage.topic
status.storage.partitions = 5
status.storage.replication.factor = 3
cleanup.policy=compact. 활성 토픽 추적 정보도 여기 들어갑니다
내부 토픽을 직접 만들기 — 자동 생성에 맡기지 않는 방법
export BS=localhost:9092
# config: 파티션 1개 필수, compact
bin/kafka-topics.sh --bootstrap-server $BS --create --if-not-exists \
--topic connect-configs --partitions 1 --replication-factor 3 \
--config cleanup.policy=compact
# offsets: source 커넥터 오프셋
bin/kafka-topics.sh --bootstrap-server $BS --create --if-not-exists \
--topic connect-offsets --partitions 25 --replication-factor 3 \
--config cleanup.policy=compact
# status: 커넥터·태스크 상태
bin/kafka-topics.sh --bootstrap-server $BS --create --if-not-exists \
--topic connect-status --partitions 5 --replication-factor 3 \
--config cleanup.policy=compact
# 확인 — cleanup.policy 가 compact 인지 반드시 봅니다
for t in connect-configs connect-offsets connect-status; do
echo "--- $t"
bin/kafka-topics.sh --bootstrap-server $BS --describe --topic "$t"
done
connect-distributed.properties — 최소 구성
bootstrap.servers=localhost:9092
# 워커 그룹 ID. 컨슈머 그룹 ID 와 겹치지 않게 하세요.
group.id=connect-cluster
# 컨버터 — 스키마를 쓰지 않으면 schemas.enable 를 false 로
key.converter=org.apache.kafka.connect.json.JsonConverter
value.converter=org.apache.kafka.connect.json.JsonConverter
key.converter.schemas.enable=false
value.converter.schemas.enable=false
# 내부 토픽
config.storage.topic=connect-configs
config.storage.replication.factor=3
offset.storage.topic=connect-offsets
offset.storage.partitions=25
offset.storage.replication.factor=3
status.storage.topic=connect-status
status.storage.partitions=5
status.storage.replication.factor=3
# 플러그인 격리 — 여기에 없으면 커넥터가 보이지 않습니다
plugin.path=/opt/connectors
# REST
listeners=http://:8083
# 컨테이너에서는 다른 워커가 도달할 수 있는 주소를 알려야 합니다
#rest.advertised.host.name=connect-1.internal
#rest.advertised.port=8083
# 롤링 재시작 중 불필요한 태스크 재배치를 미룹니다
scheduled.rebalance.max.delay.ms=300000
에러 처리와 DLQ
errors.* 설정 — 기본값은 Apache Kafka 4.3.1 커넥터 설정 문서 기준
설정
기본값
역할
errors.tolerance
none
none이면 첫 오류에서 태스크가 FAILED가 됩니다. all이면 오류 레코드를 건너뜁니다
errors.retry.timeout
0
0은 재시도하지 않음입니다. -1은 무한 재시도
errors.retry.delay.max.ms
60000 (1분)
재시도 사이 최대 지연
errors.log.enable
false
오류를 워커 로그에 기록합니다. DLQ를 쓸 때도 켜 두면 추적이 쉬워집니다
errors.log.include.messages
false
로그에 레코드 내용을 포함합니다. 개인정보가 로그로 새어 나갈 수 있습니다
errors.deadletterqueue.topic.name
"" (빈 값)
sink 커넥터 전용. 비워 두면 DLQ가 동작하지 않습니다
errors.deadletterqueue.topic.replication.factor
3
브로커가 3대 미만이면 낮춰야 합니다
errors.deadletterqueue.context.headers.enable
false
원본 토픽·파티션·오프셋·예외를 헤더로 남깁니다.반드시 켜세요 — 끄면 DLQ 레코드가 어디서 왔는지 알 수 없습니다
실무 DLQ 설정 — 네 줄이 한 세트입니다
{
"connector.class": "io.example.MySinkConnector",
"topics": "orders",
"//": "이 네 개는 함께 설정해야 의미가 있습니다",
"errors.tolerance": "all",
"errors.deadletterqueue.topic.name": "orders-sink-dlq",
"errors.deadletterqueue.topic.replication.factor": "3",
"errors.deadletterqueue.context.headers.enable": "true",
"//log": "DLQ 와 별개로 로그도 남깁니다. 다만 레코드 내용은 개인정보 위험이 있습니다",
"errors.log.enable": "true",
"errors.log.include.messages": "false",
"//retry": "일시적 오류는 재시도로 흡수합니다 (기본값 0 은 재시도 없음)",
"errors.retry.timeout": "60000",
"errors.retry.delay.max.ms": "10000"
}
DLQ 컨텍스트 헤더 10개
errors.deadletterqueue.context.headers.enable=true일 때
DLQ 레코드에 붙는 헤더입니다. 접두어는 __connect.errors.이며
헤더 이름은 DeadLetterQueueReporter.java에서 확인했습니다.
REST 엔드포인트 목록은 Apache Kafka 4.3.1 배포판에 포함된
Connect OpenAPI 정의(connect_rest.yaml)에서 추출했습니다.
SMT·Predicate의 이름과 파라미터는 생성된 변환 문서에서,
상태 정의와 409 동작은 Connect Administration 문서에서,
설정 기본값은 Connect 설정 문서에서 확인했습니다.
curl 예시와 재시작 스크립트는 이 가이드의 구성입니다.