장애 대응 시 3단계

REST API 전체 표

$CWhttp://localhost:8083으로 읽으세요. 워커의 REST 리스너 기본값은 http://:8083입니다.

Kafka Connect REST API — Apache Kafka 4.3.1 배포판의 OpenAPI 정의(connect_rest.yaml)에서 추출
메서드 경로 용도 메모
워커
GET / 이 워커의 정보와 연결된 Kafka 클러스터 ID 워커 버전 확인에 씁니다
GET /health 워커의 readiness · liveness 확인 K8s probe에 쓰기 좋습니다
GET /admin/loggers 레벨이 명시적으로 설정된 로거 목록과 레벨 admin.listeners를 따로 두면 이 리스너로만 접근됩니다
GET /admin/loggers/{logger} 특정 로거의 레벨 조회
PUT /admin/loggers/{logger} 특정 로거의 레벨 변경 재시작 없이 디버그 로그를 켤 수 있습니다
플러그인
GET /connector-plugins 설치된 커넥터 플러그인 목록 여기 없으면 plugin.path를 확인하세요
GET /connector-plugins/{pluginName}/config 플러그인의 설정 정의 조회 어떤 설정을 받는지 모를 때
PUT /connector-plugins/{pluginName}/config/validate 설정을 실제로 만들지 않고 검증 배포 파이프라인에 넣기 가장 좋은 엔드포인트입니다. 오류를 필드별로 알려 줍니다
커넥터 생성 · 조회
GET /connectors 활성 커넥터 목록 ?expand=status&expand=info로 상세를 한 번에 받을 수 있습니다
POST /connectors 커넥터 생성. 본문은 name + config 객체 initial_stateSTOPPED / 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

두 정지 방식의 차이
관점pausestop (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
DropHeaders 지정한 헤더를 제거 headers
ExtractField Struct·Map에서 특정 필드만 꺼내 전체를 대체 field, field.syntax.version, replace.null.with.default null 값은 그대로 통과합니다. $Key/$Value
Filter 레코드를 전부 드롭 없음 Predicate와 조건부로 써야 의미가 있습니다. 단독으로 쓰면 모든 레코드가 사라집니다
Flatten 중첩 구조를 평탄화. 필드명을 구분자로 이어 붙입니다 delimiter RDB sink처럼 평탄한 스키마가 필요할 때. $Key/$Value
HeaderFrom 키·값의 필드를 헤더로 이동 또는 복사 fields, headers, operation, replace.null.with.default operationmove인지 copy인지가 핵심. $Key/$Value
HoistField 값을 지정한 필드명으로 한 겹 감쌉니다 field ExtractField의 반대. 스키마 없는 값을 구조화할 때. $Key/$Value
InsertField 레코드 메타데이터 또는 정적 값을 필드로 추가 topic.field, partition.field, offset.field, timestamp.field, static.field, static.value sink에 원본 위치를 남기는 데 매우 유용합니다. $Key/$Value
InsertHeader 모든 레코드에 헤더 추가 header, value.literal 리터럴 값만 넣을 수 있습니다
MaskField 지정 필드를 타입별 null 값(0, false, 빈 문자열 등)으로 마스킹 fields, replacement, replace.null.with.default 개인정보 제거에 자주 씁니다. $Key/$Value
RegexRouter 정규식으로 대상 토픽명을 변경 regex, replacement 토픽 접두어 제거·추가에 가장 많이 쓰이는 SMT
ReplaceField 필드를 필터링하거나 이름을 변경 include, exclude, renames, replace.null.with.default 불필요한 필드 제거 + 이름 정규화. $Key/$Value
SetSchemaMetadata 스키마 이름·버전을 설정 schema.name, schema.version, replace.null.with.default Schema Registry 연동 시 subject 제어. $Key/$Value
TimestampConverter Unix epoch · 문자열 · Connect Date/Timestamp 간 변환 target.type, field, format, unix.precision, replace.null.with.default 개별 필드 또는 전체 값에 적용. $Key/$Value
TimestampRouter 원래 토픽명 + 레코드 타임스탬프로 토픽명을 재구성 timestamp.format, topic.format 날짜별 토픽·인덱스 라우팅 (예: Elasticsearch sink)
ValueToKey 값의 일부 필드로 키를 새로 만듭니다 fields, replace.null.with.default source가 키를 안 주는 경우 파티셔닝을 위해 필수. ExtractField$Key와 조합하는 패턴이 흔합니다

Predicate

SMT를 조건부로 적용하려면 Predicate를 씁니다. negate로 조건을 반전할 수 있습니다.

내장 Predicate — Apache Kafka 4.3.1 생성 문서 기준
이름참이 되는 조건파라미터
HasHeaderKey지정한 이름의 헤더가 하나 이상 있는 레코드name
RecordIsTombstonetombstone 레코드 (값이 null)없음
TopicNameMatches토픽명이 정규식에 매칭되는 레코드pattern
SMT 체인 — 별칭 순서대로 적용됩니다
{
  "connector.class": "io.example.MySinkConnector",
  "tasks.max": "2",
  "topics.regex": "prod\\.orders\\..*",

  "//": "transforms 에 나열한 별칭 순서가 곧 적용 순서입니다",
  "transforms": "extractKey,maskPii,addSource,route",

  "//1": "1) 값의 orderId 필드로 키를 만듭니다 (파티셔닝 근거 확보)",
  "transforms.extractKey.type": "org.apache.kafka.connect.transforms.ValueToKey",
  "transforms.extractKey.fields": "orderId",

  "//2": "2) 개인정보 필드를 마스킹합니다",
  "transforms.maskPii.type": "org.apache.kafka.connect.transforms.MaskField$Value",
  "transforms.maskPii.fields": "email,phone,ssn",
  "transforms.maskPii.replacement": "REDACTED",

  "//3": "3) 원본 위치를 필드로 남깁니다 (사후 추적용)",
  "transforms.addSource.type": "org.apache.kafka.connect.transforms.InsertField$Value",
  "transforms.addSource.topic.field": "_src_topic",
  "transforms.addSource.partition.field": "_src_partition",
  "transforms.addSource.offset.field": "_src_offset",

  "//4": "4) prod. 접두어를 떼어 대상 토픽명을 바꿉니다",
  "transforms.route.type": "org.apache.kafka.connect.transforms.RegexRouter",
  "transforms.route.regex": "prod\\.(.*)",
  "transforms.route.replacement": "$1"
}
Predicate로 조건부 적용 — tombstone만 드롭
{
  "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에서 확인했습니다.

DLQ 컨텍스트 헤더 — 이것이 없으면 DLQ 레코드는 출처를 알 수 없습니다
헤더내용
__connect.errors.topic원본 토픽 이름
__connect.errors.partition원본 파티션
__connect.errors.offset원본 오프셋 — 이 세 개로 원본 레코드를 정확히 다시 찾을 수 있습니다
__connect.errors.connector.name실패한 커넥터 이름
__connect.errors.task.id실패한 태스크 ID
__connect.errors.stage실패 단계 (변환 · 컨버터 · put 등) — 어디서 깨졌는지를 알려 줍니다
__connect.errors.class.name실패 시점에 실행 중이던 클래스
__connect.errors.exception.class.name예외 클래스명
__connect.errors.exception.message예외 메시지
__connect.errors.exception.stacktrace스택트레이스
DLQ 레코드의 원인 헤더 읽기
bin/kafka-console-consumer.sh --bootstrap-server localhost:9092 \
  --topic orders-sink-dlq --from-beginning --max-messages 5 \
  --property print.headers=true \
  --property print.key=true \
  --property print.offset=true

# 실패 단계와 예외만 빠르게 집계하려면 (stage 로 원인 분류)
bin/kafka-console-consumer.sh --bootstrap-server localhost:9092 \
  --topic orders-sink-dlq --from-beginning --timeout-ms 10000 \
  --property print.headers=true --property print.value=false \
  | grep -oE '__connect\.errors\.(stage|exception\.class\.name):[^,]*' \
  | sort | uniq -c | sort -rn

source vs sink

두 커넥터 종류의 차이 — 오프셋 저장 위치가 핵심입니다
관점SourceSink
방향 외부 시스템 → Kafka Kafka → 외부 시스템
오프셋 저장 위치 offset.storage.topic (Connect 내부 토픽) __consumer_offsets (일반 컨슈머와 동일)
오프셋의 의미 커넥터가 정의하는 임의의 위치 (파일 오프셋, DB binlog 위치 등) Kafka 토픽-파티션의 오프셋
병렬성 상한 커넥터가 정의하는 분할 단위 (테이블 수, 파일 수 등) 입력 토픽의 파티션 수
DLQ 없음 있음 (errors.deadletterqueue.*)
토픽 지정 커넥터 고유 설정 (또는 topic.creation.groups) topics 또는 topics.regex (둘 중 하나만)
EOS 워커 exactly.once.source.support + 커넥터 transaction.boundary 커넥터 구현에 따름 (외부 시스템의 멱등성 필요)
sink 커넥터의 lag 확인 — 결국 컨슈머 그룹입니다
# sink 커넥터는 __consumer_offsets 를 쓰므로 일반 컨슈머 도구가 통합니다
bin/kafka-consumer-groups.sh --bootstrap-server localhost:9092 --list \
  | grep '^connect-'

bin/kafka-consumer-groups.sh --bootstrap-server localhost:9092 \
  --describe --group connect-orders-file-sink

# source 커넥터는 위 방법으로 볼 수 없습니다 — REST 를 씁니다
curl -s "$CW/connectors/mysql-source/offsets" | jq

standalone vs distributed

두 실행 모드
관점standalonedistributed
기동connect-standalone.shconnect-distributed.sh
커넥터 설정파일로 인자 전달REST API로만
오프셋·설정·상태 저장로컬 파일Kafka 내부 토픽 3개
확장불가 (단일 프로세스)워커 추가로 확장
장애 복구없음태스크가 다른 워커로 재배치
용도개발·테스트, 단일 노드에서만 가능한 소스(로컬 파일 등)운영

자주 걸리는 지점

공식 문서 출처

REST 엔드포인트 목록은 Apache Kafka 4.3.1 배포판에 포함된 Connect OpenAPI 정의(connect_rest.yaml)에서 추출했습니다. SMT·Predicate의 이름과 파라미터는 생성된 변환 문서에서, 상태 정의와 409 동작은 Connect Administration 문서에서, 설정 기본값은 Connect 설정 문서에서 확인했습니다. curl 예시와 재시작 스크립트는 이 가이드의 구성입니다.