실무 예제 · 8
Kafka Connect CDC 파이프라인
MySQL의 binlog를 Debezium이 읽어 Kafka 토픽으로 흘리고, sink 커넥터가 그것을 외부 저장소로 내립니다. 코드를 한 줄도 쓰지 않고 JSON 설정만으로 파이프라인이 만들어지는 대신, 내부 토픽 3개·오프셋 저장 위치의 비대칭·에러 처리를 모르면 장애 때 손을 쓸 수 없습니다. 그 부분을 중심으로 다룹니다.
학습 목표
- Connect distributed 모드의 내부 토픽 3개가 무엇을 저장하고 왜
compact여야 하는지 설명할 수 있습니다. - source 커넥터와 sink 커넥터의 오프셋 저장 위치가 다르다는 것을 확인할 수 있습니다.
errors.*설정으로 DLQ를 구성하고 실패 원인 헤더를 읽을 수 있습니다.- REST API로 커넥터를 배포·상태 확인·재시작·오프셋 조회할 수 있습니다.
시나리오
주문 서비스의 MySQL orders 테이블을 분석 시스템과 검색 인덱스가 함께 필요로 합니다.
기존에는 두 팀이 각각 야간 배치로 전체를 덤프해 가져갔습니다.
데이터가 최대 24시간 낡고, 배치가 돌 때마다 DB에 부하가 걸렸습니다.
CDC(Change Data Capture)로 바꿉니다. Debezium이 MySQL binlog를 읽어 INSERT/UPDATE/DELETE를 이벤트로 발행하면, 하류는 그 스트림만 소비합니다. DB에 추가 쿼리가 발생하지 않고 지연은 초 단위가 됩니다.
애플리케이션 코드를 고칠 수 없다는 제약이 있어 outbox 패턴 대신 테이블 직접 캡처를 선택했습니다.
아키텍처
offset.storage.topic,
sink는 __consumer_offsets
| 설정 | 저장하는 것 | 권장 파티션 | cleanup.policy |
|---|---|---|---|
config.storage.topic |
커넥터·태스크 설정 | 반드시 1 (순서가 의미를 갖습니다) | compact |
offset.storage.topic |
source 커넥터의 진행 위치 (binlog 파일·포지션 등) | 기본 25 |
compact |
status.storage.topic |
커넥터·태스크 상태, 사용 중인 토픽 목록 | 기본 5 |
compact |
사전 요구사항
| 항목 | 버전 | 비고 |
|---|---|---|
| Apache Kafka | 4.3.1 | 예제 1의 클러스터. Connect도 같은 apache/kafka:4.3.1 이미지로 띄웁니다 |
| Debezium MySQL 커넥터 | 3.6.0.Final | Maven Central의 io.debezium:debezium-connector-mysql 플러그인 tarball을 이미지 빌드 시 내려받습니다 |
| MySQL | mysql:8.4 | binlog를 ROW 포맷으로 켭니다 |
| Docker Compose | v2 |
전체 코드
connect-cdc/
├── docker-compose.connect.yml # MySQL + Connect 워커 2대
├── Dockerfile.connect # apache/kafka 4.3.1 + Debezium 플러그인
├── mysql/
│ ├── my.cnf # binlog 설정
│ └── init.sql # 스키마 + 사용자 + 샘플 데이터
└── connectors/
├── mysql-source.json # Debezium MySQL 소스
├── file-sink.json # FileStreamSink (Apache Kafka 내장, 바로 동작)
└── s3-sink.json # S3 Sink 참고 설정 (플러그인 별도 설치 필요)
Connect 이미지
# Apache Kafka 4.3.1 이미지에 Debezium MySQL 플러그인을 얹습니다.
# 브로커와 Connect 워커의 버전을 일치시키기 위해 같은 베이스를 씁니다.
FROM apache/kafka:4.3.1
USER root
# 플러그인은 반드시 "커넥터별 하위 디렉터리" 로 격리해야 합니다.
# 한 디렉터리에 여러 커넥터의 jar 를 섞으면 클래스로더 충돌이 발생합니다.
ARG DEBEZIUM_VERSION=3.6.0.Final
ARG PLUGIN_DIR=/opt/connect-plugins
RUN mkdir -p ${PLUGIN_DIR}/debezium-connector-mysql \
&& curl -fsSL \
"https://repo1.maven.org/maven2/io/debezium/debezium-connector-mysql/${DEBEZIUM_VERSION}/debezium-connector-mysql-${DEBEZIUM_VERSION}-plugin.tar.gz" \
-o /tmp/dbz.tar.gz \
&& tar -xzf /tmp/dbz.tar.gz -C ${PLUGIN_DIR}/debezium-connector-mysql --strip-components=1 \
&& rm /tmp/dbz.tar.gz \
&& chown -R appuser:appuser ${PLUGIN_DIR}
USER appuser
MySQL 설정
[mysqld]
# Debezium 은 binlog 를 읽습니다. ROW 포맷이어야 변경된 행의 값 전체를 얻습니다.
# STATEMENT 나 MIXED 로는 CDC 가 성립하지 않습니다.
binlog_format = ROW
# 변경 전/후 이미지를 모두 기록합니다. UPDATE 이벤트의 before 를 얻기 위해 필요합니다.
binlog_row_image = FULL
# 복제 소스로 동작하기 위한 서버 ID. Debezium 의 database.server.id 와는 다른 값입니다.
server_id = 1
log_bin = mysql-bin
# binlog 보관 기간. 이 기간보다 오래 Debezium 이 멈춰 있으면
# 필요한 binlog 가 사라져 스냅샷을 다시 떠야 합니다(운영 중 가장 위험한 지점).
binlog_expire_logs_seconds = 604800
# GTID 를 켜면 소스 전환 시 위치 추적이 쉬워집니다.
gtid_mode = ON
enforce_gtid_consistency = ON
-- CDC 전용 사용자. 애플리케이션 계정을 재사용하지 마세요.
CREATE USER 'debezium'@'%' IDENTIFIED BY 'dbz-secret';
-- Debezium 이 필요한 최소 권한입니다.
-- SELECT : 초기 스냅샷
-- RELOAD : 스냅샷 시 글로벌 읽기 잠금
-- SHOW DATABASES : 대상 탐색
-- REPLICATION SLAVE : binlog 스트림 읽기
-- REPLICATION CLIENT : binlog 위치 조회
GRANT SELECT, RELOAD, SHOW DATABASES, REPLICATION SLAVE, REPLICATION CLIENT
ON *.* TO 'debezium'@'%';
FLUSH PRIVILEGES;
CREATE DATABASE IF NOT EXISTS shop;
USE shop;
CREATE TABLE orders (
id BIGINT NOT NULL AUTO_INCREMENT,
order_no VARCHAR(40) NOT NULL,
customer_id BIGINT NOT NULL,
amount DECIMAL(12,2) NOT NULL,
status VARCHAR(20) NOT NULL,
created_at DATETIME(3) NOT NULL DEFAULT CURRENT_TIMESTAMP(3),
updated_at DATETIME(3) NOT NULL DEFAULT CURRENT_TIMESTAMP(3)
ON UPDATE CURRENT_TIMESTAMP(3),
PRIMARY KEY (id),
UNIQUE KEY uk_order_no (order_no)
);
-- 초기 스냅샷 대상이 될 데이터
INSERT INTO orders (order_no, customer_id, amount, status) VALUES
('ORD-1001', 501, 32000.00, 'CREATED'),
('ORD-1002', 502, 15500.00, 'PAID'),
('ORD-1003', 501, 9800.00, 'CREATED');
Compose — Connect 워커 2대
# 예제 1의 kafka-lab 네트워크에 MySQL 과 Connect 워커 2대를 붙입니다.
# 워커를 2대로 두는 이유: 태스크 재분배와 워커 이탈 복구를 실습하기 위해서입니다.
---
name: kafka-lab-connect
x-connect-common: &connect-common
build:
context: .
dockerfile: Dockerfile.connect
restart: unless-stopped
# 워커를 "connect-distributed.sh" 로 띄웁니다.
# 설정은 환경변수가 아니라 프로퍼티 파일이 필요하므로
# 셸에서 파일을 생성해 넘깁니다(이미지가 Connect 용 변환을 하지 않습니다).
command:
- /bin/bash
- -c
- |
cat > /tmp/connect.properties <<'EOF'
bootstrap.servers=kafka-1:19092,kafka-2:19092,kafka-3:19092
# 이 값이 같은 워커들이 하나의 Connect 클러스터를 이룹니다.
# 컨슈머 그룹 ID 와 같은 이름공간을 쓰므로 다른 애플리케이션과 겹치면 안 됩니다.
group.id=cdc-connect-cluster
# REST 리스너. 4.x 에서는 rest.port 가 없고 listeners 를 씁니다(기본 http://:8083).
listeners=http://0.0.0.0:8083
# 워커끼리 서로를 찾을 때 쓰는 주소. 컨테이너 이름을 광고합니다.
rest.advertised.host.name=${CONNECT_HOSTNAME}
# --- 컨버터 ---------------------------------------------------------
# 이 예제는 스키마 레지스트리 없이 JSON 을 씁니다.
# schemas.enable=true 면 메시지에 스키마가 함께 들어가 크기가 커집니다.
# 운영에서는 Avro + Schema Registry 를 권장합니다(예제 7).
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
# --- 내부 토픽 3개 ---------------------------------------------------
# config 는 파티션이 반드시 1이어야 합니다(설정 변경 순서가 의미를 가집니다).
config.storage.topic=connect-cdc-configs
config.storage.replication.factor=3
# source 커넥터의 진행 위치가 여기 저장됩니다. 기본 파티션 25.
offset.storage.topic=connect-cdc-offsets
offset.storage.replication.factor=3
offset.storage.partitions=25
# 커넥터/태스크 상태. 기본 파티션 5.
status.storage.topic=connect-cdc-status
status.storage.replication.factor=3
status.storage.partitions=5
# source 오프셋을 커밋하는 주기. 기본값 60000(1분).
# 짧게 두면 재시작 시 재처리 구간이 줄지만 쓰기가 늘어납니다.
offset.flush.interval.ms=10000
# 오프셋 커밋 타임아웃. 기본값 5000.
offset.flush.timeout.ms=5000
# --- 플러그인 -------------------------------------------------------
# 커넥터별 하위 디렉터리로 격리해 두었습니다.
plugin.path=/opt/connect-plugins
# 기본값 hybrid_warn. 서비스 로더 방식으로 마이그레이션되지 않은 플러그인에
# 경고를 남깁니다. service_load 로 올리면 기동이 빨라집니다.
plugin.discovery=hybrid_warn
EOF
exec /opt/kafka/bin/connect-distributed.sh /tmp/connect.properties
networks:
- kafka-lab_default
healthcheck:
test: ['CMD-SHELL', 'curl -sf http://localhost:8083/ || exit 1']
interval: 10s
timeout: 5s
retries: 18
start_period: 40s
services:
mysql:
image: mysql:8.4
hostname: mysql
container_name: mysql
ports:
- '3306:3306'
environment:
MYSQL_ROOT_PASSWORD: root-secret
volumes:
- ./mysql/my.cnf:/etc/mysql/conf.d/binlog.cnf:ro
- ./mysql/init.sql:/docker-entrypoint-initdb.d/init.sql:ro
- mysql-data:/var/lib/mysql
networks:
- kafka-lab_default
healthcheck:
test: ['CMD-SHELL', 'mysqladmin ping -h 127.0.0.1 -uroot -proot-secret || exit 1']
interval: 10s
timeout: 5s
retries: 12
start_period: 30s
connect-1:
<<: *connect-common
hostname: connect-1
container_name: connect-1
ports:
- '8083:8083'
environment:
CONNECT_HOSTNAME: connect-1
depends_on:
mysql:
condition: service_healthy
connect-2:
<<: *connect-common
hostname: connect-2
container_name: connect-2
ports:
- '8084:8083'
environment:
CONNECT_HOSTNAME: connect-2
depends_on:
mysql:
condition: service_healthy
volumes:
mysql-data:
networks:
kafka-lab_default:
external: true
Debezium MySQL 소스 커넥터
{
"name": "shop-orders-source",
"config": {
"connector.class": "io.debezium.connector.mysql.MySqlConnector",
"//": "MySQL binlog 는 단일 스트림이므로 소스 태스크는 항상 1개입니다.",
"//": "tasks.max 를 올려도 태스크가 늘지 않습니다 — 흔한 오해입니다.",
"tasks.max": "1",
"database.hostname": "mysql",
"database.port": "3306",
"database.user": "debezium",
"database.password": "dbz-secret",
"//": "MySQL 복제 클러스터 안에서 유일해야 하는 값입니다.",
"//": "MySQL 서버의 server_id 와 겹치면 기존 복제가 끊어집니다.",
"database.server.id": "184054",
"//": "생성되는 모든 토픽 이름의 접두어입니다. 토픽명은 {prefix}.{db}.{table} 이 됩니다.",
"//": "이 값을 나중에 바꾸면 토픽이 전부 새로 생기고 하류가 끊깁니다. 처음에 잘 정하세요.",
"topic.prefix": "shopdb",
"database.include.list": "shop",
"table.include.list": "shop.orders",
"//": "DDL 이력을 저장하는 내부 토픽. 커넥터마다 별도 토픽이어야 합니다.",
"schema.history.internal.kafka.bootstrap.servers": "kafka-1:19092,kafka-2:19092,kafka-3:19092",
"schema.history.internal.kafka.topic": "schema-changes.shop",
"//": "initial: 최초 1회 전체 스냅샷 후 binlog 스트리밍. 기본 동작입니다.",
"//": "never 로 두면 스냅샷 없이 지금부터의 변경만 받습니다.",
"snapshot.mode": "initial",
"//": "Debezium 이벤트는 before/after/source/op 를 감싼 봉투(envelope) 구조입니다.",
"//": "하류가 '변경 후의 행' 만 원하면 ExtractNewRecordState SMT 로 펼칩니다.",
"transforms": "unwrap,route",
"transforms.unwrap.type": "io.debezium.transforms.ExtractNewRecordState",
"//": "DELETE 는 tombstone(null value)이 됩니다. 컴팩션 토픽에서는 이것이 삭제 마커입니다.",
"transforms.unwrap.delete.handling.mode": "rewrite",
"//": "op(c/u/d/r) 와 타임스탬프를 헤더가 아닌 필드로 남깁니다. 하류 디버깅에 유용합니다.",
"transforms.unwrap.add.fields": "op,source.ts_ms,source.db,source.table",
"//": "토픽 이름을 shopdb.shop.orders → cdc.orders 로 단순화합니다.",
"transforms.route.type": "org.apache.kafka.connect.transforms.RegexRouter",
"transforms.route.regex": "shopdb\\.shop\\.(.*)",
"transforms.route.replacement": "cdc.$1",
"//": "--- 에러 처리 / DLQ (Connect 가 설정으로 제공합니다) ---",
"//": "errors.tolerance 기본값은 none 이며, 첫 실패에서 태스크가 FAILED 로 죽습니다.",
"errors.tolerance": "all",
"//": "errors.log.enable 기본값 false. 켜지 않으면 무엇이 실패했는지 알 수 없습니다.",
"errors.log.enable": "true",
"errors.log.include.messages": "true",
"//": "errors.retry.timeout 기본값 0 = 재시도 없음. -1 은 무제한입니다.",
"errors.retry.timeout": "60000",
"//": "errors.retry.delay.max.ms 기본값 60000.",
"errors.retry.delay.max.ms": "30000",
"//": "source 커넥터는 DLQ 를 지원하지 않습니다 — DLQ 는 sink 전용 기능입니다.",
"//": "위 errors.* 중 deadletterqueue.* 는 sink 커넥터에만 유효합니다.",
"//": "생성되는 토픽의 기본 스펙. 미리 만들어 두는 편이 더 안전합니다.",
"topic.creation.default.replication.factor": "3",
"topic.creation.default.partitions": "3",
"topic.creation.default.cleanup.policy": "compact",
"topic.creation.groups": ""
}
}
Sink 커넥터 — 바로 동작하는 버전
{
"name": "orders-file-sink",
"config": {
"connector.class": "org.apache.kafka.connect.file.FileStreamSinkConnector",
"tasks.max": "1",
"topics": "cdc.orders",
"file": "/tmp/cdc-orders.out",
"key.converter": "org.apache.kafka.connect.storage.StringConverter",
"value.converter": "org.apache.kafka.connect.json.JsonConverter",
"value.converter.schemas.enable": "false",
"errors.tolerance": "all",
"errors.log.enable": "true",
"errors.log.include.messages": "true",
"errors.deadletterqueue.topic.name": "cdc.orders.DLQ",
"errors.deadletterqueue.topic.replication.factor": "3",
"errors.deadletterqueue.context.headers.enable": "true",
"errors.retry.timeout": "60000",
"errors.retry.delay.max.ms": "30000"
}
}
errors.tolerance와 DLQ 토픽, 헤더에 담기는 실패 원인
S3 Sink — 참고 설정
{
"name": "orders-s3-sink",
"config": {
"connector.class": "io.confluent.connect.s3.S3SinkConnector",
"tasks.max": "3",
"topics": "cdc.orders",
"s3.bucket.name": "analytics-raw",
"s3.region": "ap-northeast-2",
"storage.class": "io.confluent.connect.s3.storage.S3Storage",
"format.class": "io.confluent.connect.s3.format.parquet.ParquetFormat",
"flush.size": "10000",
"rotate.interval.ms": "300000",
"partitioner.class": "io.confluent.connect.storage.partitioner.TimeBasedPartitioner",
"path.format": "'dt'=yyyy-MM-dd/'hour'=HH",
"partition.duration.ms": "3600000",
"locale": "ko-KR",
"timezone": "Asia/Seoul",
"key.converter": "org.apache.kafka.connect.storage.StringConverter",
"value.converter": "org.apache.kafka.connect.json.JsonConverter",
"value.converter.schemas.enable": "false",
"errors.tolerance": "all",
"errors.log.enable": "true",
"errors.deadletterqueue.topic.name": "cdc.orders.DLQ",
"errors.deadletterqueue.topic.replication.factor": "3",
"errors.deadletterqueue.context.headers.enable": "true"
}
}
실행 방법
# 0. 예제 1의 클러스터
cd kafka-lab && docker compose ps
# 1. config 내부 토픽은 파티션이 1이어야 하므로 미리 만듭니다.
# Connect 가 자동 생성하게 두어도 되지만, 명시하는 편이 안전합니다.
./kcli kafka-topics.sh --create --if-not-exists --topic connect-cdc-configs \
--partitions 1 --replication-factor 3 --config cleanup.policy=compact
./kcli kafka-topics.sh --create --if-not-exists --topic connect-cdc-offsets \
--partitions 25 --replication-factor 3 --config cleanup.policy=compact
./kcli kafka-topics.sh --create --if-not-exists --topic connect-cdc-status \
--partitions 5 --replication-factor 3 --config cleanup.policy=compact
# 2. MySQL + Connect 워커 2대 기동 (이미지 빌드 포함)
cd ../connect-cdc
docker compose -f docker-compose.connect.yml up -d --build
# 3. 워커가 준비되고 플러그인이 로드되었는지 확인
curl -s http://localhost:8083/ | jq
curl -s http://localhost:8083/connector-plugins | jq -r '.[].class' | grep -i mysql
# → io.debezium.connector.mysql.MySqlConnector
# 4. 소스 커넥터 배포
curl -s -X POST http://localhost:8083/connectors \
-H 'Content-Type: application/json' \
-d @connectors/mysql-source.json | jq
# 5. sink 커넥터 배포
curl -s -X POST http://localhost:8083/connectors \
-H 'Content-Type: application/json' \
-d @connectors/file-sink.json | jq
# 6. 상태 확인 — 둘 다 RUNNING 이어야 합니다.
curl -s http://localhost:8083/connectors/shop-orders-source/status | jq
curl -s http://localhost:8083/connectors/orders-file-sink/status | jq
docker exec -i mysql mysql -uroot -proot-secret shop <<'SQL'
INSERT INTO orders (order_no, customer_id, amount, status)
VALUES ('ORD-2001', 503, 47000.00, 'CREATED');
UPDATE orders SET status = 'PAID', amount = 47000.00 WHERE order_no = 'ORD-2001';
DELETE FROM orders WHERE order_no = 'ORD-1003';
SQL
검증 방법
1. CDC 토픽이 만들어지고 이벤트가 들어왔는가
cd ../kafka-lab
./kcli kafka-topics.sh --list | grep -E 'cdc\.|schema-changes'
# → cdc.orders
# schema-changes.shop
./kcli kafka-console-consumer.sh --topic cdc.orders --from-beginning \
--timeout-ms 10000 --property print.key=true 2>/dev/null | head -10
op 필드로 변경 종류를 구분합니다{"id":1} {"id":1,"order_no":"ORD-1001","customer_id":501,"amount":"...","status":"CREATED","__op":"r","__source_ts_ms":...}
{"id":4} {"id":4,"order_no":"ORD-2001","customer_id":503,"amount":"...","status":"CREATED","__op":"c","__source_ts_ms":...}
{"id":4} {"id":4,"order_no":"ORD-2001","customer_id":503,"amount":"...","status":"PAID","__op":"u","__source_ts_ms":...}
{"id":3} {"id":3,"order_no":"ORD-1003","__deleted":"true","__op":"d","__source_ts_ms":...}
__op 값이 핵심입니다.
r은 초기 스냅샷(read), c는 INSERT, u는 UPDATE,
d는 DELETE입니다.
키가 테이블의 프라이머리 키({"id":4})로 설정되므로
같은 행의 변경이 항상 같은 파티션으로 가 순서가 보장됩니다.
이것이 CDC에서 키 설계가 자동으로 해결되는 이유입니다.
2. source와 sink의 오프셋이 다른 곳에 있는지 확인
curl -s http://localhost:8083/connectors/shop-orders-source/offsets | jq
{
"offsets": [
{
"partition": { "server": "shopdb" },
"offset": {
"file": "mysql-bin.000003",
"pos": 4127,
"ts_sec": 1785000000,
"snapshot_completed": true
}
}
]
}
cd ../kafka-lab
# sink 커넥터는 connect-{커넥터명} 이라는 컨슈머 그룹을 씁니다.
./kcli kafka-consumer-groups.sh --list | grep connect-
# → connect-orders-file-sink
./kcli kafka-consumer-groups.sh --describe --group connect-orders-file-sink
source는 REST로만 보이고 sink는 컨슈머 그룹으로 보입니다. 이 비대칭이 실무에서 가장 자주 헷갈리는 지점이고 시험에도 나옵니다.
3. sink 파일에 실제로 적재되었는가
# 태스크가 어느 워커에 배정됐는지 먼저 확인합니다.
curl -s http://localhost:8083/connectors/orders-file-sink/status \
| jq -r '.tasks[] | "task \(.id) → \(.worker_id) (\(.state))"'
# → task 0 → connect-2:8083 (RUNNING)
# 그 워커의 컨테이너에서 파일을 봅니다.
docker exec connect-2 tail -5 /tmp/cdc-orders.out
4. 워커 이탈 시 태스크가 재분배되는가
# 태스크가 있는 워커를 죽입니다.
docker compose -f ../connect-cdc/docker-compose.connect.yml stop connect-2
# 잠시 뒤 남은 워커로 태스크가 옮겨갑니다.
sleep 20
curl -s http://localhost:8083/connectors/orders-file-sink/status \
| jq -r '.tasks[] | "task \(.id) → \(.worker_id) (\(.state))"'
# → task 0 → connect-1:8083 (RUNNING)
docker compose -f ../connect-cdc/docker-compose.connect.yml start connect-2
5. DLQ가 동작하는가
# value.converter 가 JsonConverter 인 토픽에 JSON 이 아닌 값을 넣습니다.
echo 'not-a-json' | ./kcli kafka-console-producer.sh --topic cdc.orders
# DLQ 토픽에 헤더와 함께 나타나야 합니다.
./kcli kafka-console-consumer.sh --topic cdc.orders.DLQ --from-beginning \
--timeout-ms 10000 --property print.headers=true 2>/dev/null
errors.deadletterqueue.context.headers.enable=true일 때만 붙습니다__connect.errors.topic:cdc.orders,__connect.errors.partition:0,__connect.errors.offset:5,
__connect.errors.connector.name:orders-file-sink,__connect.errors.task.id:0,
__connect.errors.stage:VALUE_CONVERTER,__connect.errors.class.name:org.apache.kafka.connect.json.JsonConverter,
__connect.errors.exception.class.name:org.apache.kafka.connect.errors.DataException,...
__connect.errors.stage가 어느 단계에서 실패했는지를 알려 줍니다 —
VALUE_CONVERTER(역직렬화), TRANSFORMATION(SMT),
TASK_PUT(sink 쓰기) 등입니다.
이 헤더가 없으면 DLQ는 원인을 알 수 없는 레코드 더미가 됩니다.
운영에 필요한 REST API
| 하고 싶은 것 | 요청 | 비고 |
|---|---|---|
| 커넥터 목록 | GET /connectors | |
| 배포 | POST /connectors | 본문에 name과 config. initial_state로 STOPPED/PAUSED/RUNNING(기본)을 지정할 수 있습니다 |
| 설정 변경 | PUT /connectors/{name}/config | 본문은 config 객체만. 없으면 생성됩니다 |
| 상태 확인 | GET /connectors/{name}/status | 커넥터·태스크 상태, 배정된 워커, 실패 시 스택트레이스 |
| 실패한 태스크만 재시작 | POST /connectors/{name}/restart?includeTasks=true&onlyFailed=true | 가장 자주 쓰는 복구 명령입니다 |
| 개별 태스크 재시작 | POST /connectors/{name}/tasks/{taskId}/restart | |
| 일시 정지 | PUT /connectors/{name}/pause | 자원은 유지. 재개가 빠릅니다 |
| 정지(자원 해제) | PUT /connectors/{name}/stop | 오프셋을 수정하려면 이 상태여야 합니다 |
| 재개 | PUT /connectors/{name}/resume | |
| 오프셋 조회 | GET /connectors/{name}/offsets | source·sink 모두 조회 가능 |
| 오프셋 리셋 | DELETE /connectors/{name}/offsets | stop 상태여야 합니다. source 커넥터를 처음부터 다시 읽히는 방법 |
| 오프셋 특정 위치로 | PATCH /connectors/{name}/offsets | stop 상태 필요. 본문은 offsets 배열 |
| 사용 중인 토픽 | GET /connectors/{name}/topics | status.storage.topic에 기록된 추적 정보 |
| 삭제 | DELETE /connectors/{name} | 태스크를 멈추고 설정을 지웁니다. 오프셋은 지워지지 않습니다 |
| 설치된 플러그인 | GET /connector-plugins | plugin.path 문제를 진단할 때 첫 번째로 봅니다 |
# 1) 정지 (pause 가 아니라 stop 이어야 합니다)
curl -s -X PUT http://localhost:8083/connectors/shop-orders-source/stop
# 2) 상태가 STOPPED 가 되었는지 확인
curl -s http://localhost:8083/connectors/shop-orders-source/status | jq '.connector.state'
# 3) 오프셋 리셋
curl -s -X DELETE http://localhost:8083/connectors/shop-orders-source/offsets | jq
# 4) 재개 → snapshot.mode=initial 이므로 스냅샷부터 다시 시작합니다
curl -s -X PUT http://localhost:8083/connectors/shop-orders-source/resume
프로덕션 고려사항
| 항목 | 이 예제 | 프로덕션 |
|---|---|---|
| 컨버터 | JSON, schemas.enable=false |
Avro + Schema Registry. 스키마 없는 JSON은 DDL 변경 시 하류가 조용히 깨집니다(예제 7) |
| binlog 보관 | 7일 | Debezium이 멈춘 채로 이 기간을 넘기면 필요한 binlog가 사라져 전체 스냅샷을 다시 떠야 합니다. lag 알림이 필수입니다 |
| 초기 스냅샷 | 테이블 3행 | 수천만 행이면 스냅샷이 수 시간 걸리고 DB에 큰 부하가 갑니다. 읽기 전용 복제본을 캡처 대상으로 두거나 incremental snapshot을 씁니다 |
| 자격증명 | 설정 파일에 평문 | config.providers로 외부 시크릿에서 주입합니다. Connect REST는 GET /connectors/{name}/config로 설정을 그대로 반환하므로 평문 비밀번호가 노출됩니다 |
| REST 보안 | 인증 없음 | Connect REST는 인증이 없으면 누구나 커넥터를 삭제할 수 있습니다. 리버스 프록시로 인증을 걸고 네트워크를 제한합니다 |
워커 수와 tasks.max |
워커 2, task 1 | MySQL binlog는 단일 스트림이라 소스 태스크가 1개를 넘지 못합니다. 워커를 늘려도 소스 처리량은 늘지 않습니다. sink는 파티션 수까지 늘릴 수 있습니다 |
| 내부 토픽 | 미리 생성 | RF 3 + compact + config는 파티션 1. 여러 Connect 클러스터가 같은 내부 토픽을 공유하면 서로를 망칩니다 — group.id와 토픽 이름을 함께 분리하세요 |
| 모니터링 | REST 수동 확인 | 커넥터·태스크 상태를 주기적으로 폴링해 FAILED에 알림. sink lag은 connect-{name} 컨슈머 그룹으로 측정합니다(예제 10) |
자주 하는 실수
관련 케이스 스터디
이어서 볼 곳
공식 문서 출처
- Kafka Connect — distributed 모드, 내부 토픽 3개, worker/connector/task 모델
- Connect REST Interface — 엔드포인트 목록,
initial_state,restart?includeTasks&onlyFailed,pause/stop/resume, 오프셋 관리 엔드포인트가stop상태를 요구한다는 서술 - Connect Worker Configs —
listeners=http://:8083,plugin.path=null,plugin.discovery=hybrid_warn,offset.flush.interval.ms,offset.flush.timeout.ms=5000,offset.storage.partitions=25,status.storage.partitions=5,*.storage.replication.factor=3,header.converter=SimpleHeaderConverter - Connect Error Reporting and DLQ —
errors.tolerance=none,errors.log.enable=false,errors.log.include.messages=false,errors.retry.timeout=0,errors.retry.delay.max.ms=60000,errors.deadletterqueue.topic.name=빈 문자열,errors.deadletterqueue.topic.replication.factor=3,errors.deadletterqueue.context.headers.enable=false - Connect Administration —
GET /connectors/{name}/status출력 형식, 토픽 추적(topic.tracking.enable), pause/resume 의미 - Debezium MySQL Connector —
io.debezium.connector.mysql.MySqlConnector,topic.prefix,database.server.id,database.include.list,table.include.list,schema.history.internal.kafka.*,snapshot.mode, binlogROW포맷 요구, 필요한 MySQL 권한 - FileStreamSinkConnector (Apache Kafka 4.3) —
file설정 하나로 동작하는 내장 sink 커넥터