학습 목표

시나리오

주문 서비스의 MySQL orders 테이블을 분석 시스템과 검색 인덱스가 함께 필요로 합니다. 기존에는 두 팀이 각각 야간 배치로 전체를 덤프해 가져갔습니다. 데이터가 최대 24시간 낡고, 배치가 돌 때마다 DB에 부하가 걸렸습니다.

CDC(Change Data Capture)로 바꿉니다. Debezium이 MySQL binlog를 읽어 INSERT/UPDATE/DELETE를 이벤트로 발행하면, 하류는 그 스트림만 소비합니다. DB에 추가 쿼리가 발생하지 않고 지연은 초 단위가 됩니다.

애플리케이션 코드를 고칠 수 없다는 제약이 있어 outbox 패턴 대신 테이블 직접 캡처를 선택했습니다.

아키텍처

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 아키텍처 — worker / connector / task 분산과 리더 워커의 역할
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
distributed 모드의 내부 토픽 3개
설정 저장하는 것 권장 파티션 cleanup.policy
config.storage.topic 커넥터·태스크 설정 반드시 1 (순서가 의미를 갖습니다) compact
offset.storage.topic source 커넥터의 진행 위치 (binlog 파일·포지션 등) 기본 25 compact
status.storage.topic 커넥터·태스크 상태, 사용 중인 토픽 목록 기본 5 compact

사전 요구사항

검증 환경
항목버전비고
Apache Kafka4.3.1예제 1의 클러스터. Connect도 같은 apache/kafka:4.3.1 이미지로 띄웁니다
Debezium MySQL 커넥터3.6.0.FinalMaven Central의 io.debezium:debezium-connector-mysql 플러그인 tarball을 이미지 빌드 시 내려받습니다
MySQLmysql:8.4binlog를 ROW 포맷으로 켭니다
Docker Composev2

전체 코드

디렉터리 구조
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 이미지

connect-cdc/Dockerfile.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 설정

connect-cdc/mysql/my.cnf
[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
connect-cdc/mysql/init.sql
-- 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대

connect-cdc/docker-compose.connect.yml
# 예제 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 소스 커넥터

connect-cdc/connectors/mysql-source.json
{
  "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 커넥터 — 바로 동작하는 버전

connect-cdc/connectors/file-sink.json — Apache Kafka 내장 커넥터
{
  "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"
  }
}
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.* 가 없습니다.
Connect DLQ 흐름 — errors.tolerance와 DLQ 토픽, 헤더에 담기는 실패 원인

S3 Sink — 참고 설정

connect-cdc/connectors/s3-sink.json (참고)
{
  "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
MySQL에 변경을 만들어 CDC를 트리거
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의 오프셋이 다른 곳에 있는지 확인

source 커넥터 — Connect REST로 조회
curl -s http://localhost:8083/connectors/shop-orders-source/offsets | jq
기대 출력 — binlog 파일과 포지션이 보입니다
{
  "offsets": [
    {
      "partition": { "server": "shopdb" },
      "offset": {
        "file": "mysql-bin.000003",
        "pos": 4127,
        "ts_sec": 1785000000,
        "snapshot_completed": true
      }
    }
  ]
}
sink 커넥터 — 컨슈머 그룹으로 조회됩니다
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. 워커 이탈 시 태스크가 재분배되는가

워커 1대를 정지
# 태스크가 있는 워커를 죽입니다.
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
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 재분배 — 워커 이탈 시 남은 워커로 태스크가 옮겨가는 과정

5. DLQ가 동작하는가

sink가 처리할 수 없는 레코드를 직접 넣습니다
# 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

Kafka 4.3 Connect REST API — 목적별
하고 싶은 것요청비고
커넥터 목록GET /connectors
배포POST /connectors본문에 nameconfig. initial_stateSTOPPED/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}/offsetssource·sink 모두 조회 가능
오프셋 리셋DELETE /connectors/{name}/offsetsstop 상태여야 합니다. source 커넥터를 처음부터 다시 읽히는 방법
오프셋 특정 위치로PATCH /connectors/{name}/offsetsstop 상태 필요. 본문은 offsets 배열
사용 중인 토픽GET /connectors/{name}/topicsstatus.storage.topic에 기록된 추적 정보
삭제DELETE /connectors/{name}태스크를 멈추고 설정을 지웁니다. 오프셋은 지워지지 않습니다
설치된 플러그인GET /connector-pluginsplugin.path 문제를 진단할 때 첫 번째로 봅니다
source 커넥터를 처음부터 다시 읽히기 — 순서가 중요합니다
# 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)

자주 하는 실수

이어서 볼 곳

공식 문서 출처