[DB] CDC - 개념 및 실습(KafkaConnect, Debezium) - 작성 중(수정 필요)

Donghee Kim·2026년 6월 1일

문득문득

목록 보기
16/16

오랜만에 velog를... 최근에는 notion 에 AWS-DEA 내용을 정리하고, 뭐 이래저래 바빴구만...

오늘은 CDC 개념 및 실습이다.

계속 Debezium 과 메시징 큐를 이용해서 이기종 DB 동기화를 어느정도 해보고 싶었어서, 개념 및 기술 플로우를 이해하고 자리잡으려고 시도해보았다.


CDC

url : https://youthful-yew-08d.notion.site/Change-Data-Capture-MSA-CQRS-1a45c4826975808294d2d70679f0e457?source=copy_link

해당 내용은 제가 정리했던 것인데, 해당 내용을 보시는게 더 도움 될 수 있습니다.

  • 필요 Infra
    • K8S
    • Heterogeneous Database (이기종 DB)
      -> 이번은 PG -> PG 로 실습해보고자 한다.
    • Kafka-Cluster
    • Custom Kafka Image (+Debezium Plugin)

설계

  • 실습 아키텍처

DB Component

- Source DB

CDC의 데이터 원천(Source) 역할 수행

  • 역할

    • 변경 데이터 생성
    • WAL 기록
    • Kafka로 전송할 CDC 이벤트 제공
  • 주요 설정

    • wal_level = logical 확인

      SHOW wal_level;
      1. replica : 물리 복제용 WAL에 기록
        INSERT INTO users ~~ => "디스크 블록 1234 변경"

      2. logical : 논리적 변경 이벤트 까지 WAL에 기록
        INSERT INTO users ~~ => "User 테이블 INSERT 발생"

      • Use Case : "외부 시스템이 데이터 변경 이벤트를 읽을 수 있게 하자"
    • Publication 설정

      Publication 은? 어떤 테이블어떤 변경 이벤트 를 외부에 공개할 것인지를 정의하는 객체

      publication_name은 kafka connector 생성 시, 값 option 설정이 다음과 같이 있음.

      publication.name : debezium_publication_v2
      • connector.yaml 생성 후, 확인 쿼리

        1. tb : pg_publication

          SELECT 
              pubname AS publication_name,
              puballtables AS all_tables,
              pubinsert AS insert_enabled,
              pubupdate AS update_enabled,
              pubdelete AS delete_enabled
          FROM pg_publication;

          • pubname : publication 이름
          • all_tables : 특정 테이블만인가? 모든 테이블인가?
          • insert_enabled : Insert 이벤트 수집 가능 여부
          • update_enabled : Update 이벤트 수집 가능 여부
          • delete_enabled : Delete 이벤트 수집 가능 여부
        2. tb : pg_publication_tables

          select * from pg_publication_tables;

          • pubname : publication 이름
          • schemaname : 테이블이 속한 스키마
          • tablename : CDC 대상 테이블
          • attnames : 복제 대상 컬럼 목록(Column Filtering)
          • rowfilter : 복제 대상 행 조건(Row Filtering)
    • Replication Slot 설정

      Debezium이 WAL을 어디까지 읽었는지 추적하는 객체

      Connector가 중지되었다가 재시작되더라도 마지막 읽은 위치(LSN)부터 다시 수집할 수 있도록 한다.

      slot_name은 connector 생성 시, 값 option 설정이 다음과 같이 있음.

      slot.name : debezium_slot_v2
      • connector.yaml 생성 후, 확인 쿼리

        1. tb : pg_replication_slots

          SELECT 
              slot_name,
              plugin,
              slot_type,
              active,
              restart_lsn,
              confirmed_flush_lsn
          FROM pg_replication_slots;

          • slot_name : Replication Slot 이름
          • plugin : WAL의 물리적 변경 내용을 논리적 변경 이벤트(INSERT, UPDATE, DELETE)로 변환
          • slot_type : Replication Slot 종류
          • active : Slot 사용 여부 (true => connector 연결 중 / false => connector 끊김)
          • restart_lsn : PostgreSQL이 WAL을 보관해야 하는 최소 위치
          • confirmed_flush_lsn : Debezium이 정상적으로 읽고 처리 완료한 마지막 WAL 위치

- Target DB

Kafka Topic으로 전달된 CDC 이벤트를 최종 저장하는 데이터 베이스

  • 역할
    • CDC 데이터 저장
    • 실시간 동기화 대상
    • UPSERT 및 DELETE 처리

Kafka Component

- Kafka Cluster

CDC 이벤트를 저장하고 전달하는 메시지 브로커 역할 수행

  • 역할
    • 이벤트 저장
    • 메시지 전달
    • Consumer Offset 관리

  • kafka-cluster.yaml
apiVersion: kafka.strimzi.io/v1
kind: Kafka
metadata:
  name: cdc-kafka
  namespace: kafka
spec:
  kafka:
    listeners:
      - name: plain
        port: 9092
        type: internal
        tls: false
    config:
      offsets.topic.replication.factor: 3
      transaction.state.log.replication.factor: 3
      transaction.state.log.min.isr: 2
      default.replication.factor: 3
      min.insync.replicas: 2
  entityOperator:
    topicOperator: {}
    userOperator: {}
  • kafka-ui.yaml
apiVersion: apps/v1
kind: Deployment
metadata:
  name: kafka-ui
  namespace: kafka

spec:
  replicas: 1

  selector:
    matchLabels:
      app: kafka-ui

  template:
    metadata:
      labels:
        app: kafka-ui

    spec:
      containers:
        - name: kafka-ui
          image: provectuslabs/kafka-ui:latest

          ports:
            - containerPort: 8080

          env:
            - name: KAFKA_CLUSTERS_0_NAME
              value: local

            - name: KAFKA_CLUSTERS_0_BOOTSTRAPSERVERS
              value: cdc-kafka-kafka-bootstrap:9092

            - name: KAFKA_CLUSTERS_0_KAFKACONNECT_0_NAME
              value: cdc-connect

            - name: KAFKA_CLUSTERS_0_KAFKACONNECT_0_ADDRESS
              value: http://cdc-connect-connect-api:8083
---
apiVersion: v1
kind: Service
metadata:
  name: kafka-ui
  namespace: kafka

spec:
  type: NodePort

  selector:
    app: kafka-ui

  ports:
    - port: 8080
      targetPort: 8080
      nodePort: 30080
  • pod

- Kafka Connect

Kafa와 외부 시스템을 연결하기 위한 데이터 통합 플랫폼
Connector를 실행하고 관리하는 Runtime 역할을 수행

  • 역할

    • Connector 실행
    • Offset 관리
    • Task 관리
    • 장애 복구
  • DockerFile (Custom Image + Debezium)

# Kafka \
FROM quay.io/strimzi/kafka:1.0.0-kafka-4.1.2

USER root

# plugin installation
RUN mkdir -p /opt/kafka/plugins/debezium-postgres \
    && mkdir -p /opt/kafka/plugins/debezium-jdbc-sink

# Debezium Postgres Source Connector 설치
ADD https://repo1.maven.org/maven2/io/debezium/debezium-connector-postgres/3.2.0.Final/debezium-connector-postgres-3.2.0.Final-plugin.tar.gz /tmp/debezium-postgres.tar.gz

# 1. Extract the tar.gz file to the plugins -C /directory
#    after Extract file directory structure will be like this:
#    /opt/kafka/plugins/debezium-postgres/
#    ├── debezium-connector-postgres-3.2.0.Final.jar
# 2. --strip-components=1 to remove the top-level directory from the archive
# 3. Remove the tar.gz file after extraction to save space
RUN tar -xzf /tmp/debezium-postgres.tar.gz \
    -C /opt/kafka/plugins/debezium-postgres \
    --strip-components=1 \
    && rm /tmp/debezium-postgres.tar.gz


# Debezium JDBC Sink Connector 설치
# 4. Download and extract the full Debezium JDBC Sink Connector package (tar.gz)
#    This package includes all built-in dialects (including PostgreSQL) to support UPSERT mode.
ADD https://repo1.maven.org/maven2/io/debezium/debezium-connector-jdbc/3.2.0.Final/debezium-connector-jdbc-3.2.0.Final-plugin.tar.gz /tmp/debezium-jdbc.tar.gz

RUN tar -xzf /tmp/debezium-jdbc.tar.gz \
    -C /opt/kafka/plugins/debezium-jdbc-sink \
    --strip-components=1 \
    && rm /tmp/debezium-jdbc.tar.gz

# PostgreSQL JDBC Driver 설치
ADD https://repo1.maven.org/maven2/org/postgresql/postgresql/42.7.3/postgresql-42.7.3.jar /opt/kafka/plugins/debezium-jdbc-sink/postgresql-42.7.3.jar


# 5. plugin directory like this:
# /opt/kafka/plugins
# ├── debezium-postgres
# │   ├── debezium-connector-postgres-3.2.0.Final.jar
# ├── debezium-jdbc-sink
# │   ├── debezium-connector-jdbc-3.2.0.Final.jar
# │   └── postgresql-42.7.3.jar

# 6. Change ownership of the plugins directory to user 1001 (the default user in the Strimzi Kafka image) to ensure that Kafka can access the plugin files without permission issues.
RUN chown -R 1001:0 /opt/kafka/plugins

USER 1001
  • kafka-connect.yaml
apiVersion: kafka.strimzi.io/v1
kind: KafkaConnect
metadata:
  name: cdc-connect
  namespace: kafka
  annotations:
    strimzi.io/use-connector-resources: "true"

spec:
  replicas: 1

  bootstrapServers: cdc-kafka-kafka-bootstrap:9092

  image: kimdonghee125/cdc-connect-debezium:custom-v6

  groupId: cdc-connect-cluster
  offsetStorageTopic: cdc-connect-offsets
  configStorageTopic: cdc-connect-configs
  statusStorageTopic: cdc-connect-status

  config:
    offset.storage.replication.factor: 3
    config.storage.replication.factor: 3
    status.storage.replication.factor: 3

    key.converter: org.apache.kafka.connect.json.JsonConverter
    value.converter: org.apache.kafka.connect.json.JsonConverter

    # [변경] 데이터 이동 시 스키마 메타데이터를 유지하도록 true로 켭니다.
    key.converter.schemas.enable: true
    value.converter.schemas.enable: true

  logging:
    type: external
    valueFrom:
      configMapKeyRef:
        name: connect-log4j
        key: log4j.properties
        
---
apiVersion: v1
kind: ConfigMap
metadata:
  name: connect-log4j
  namespace: kafka
data:
  log4j.properties: |
    # Log4j 2.x Properties Configuration
    status = error
    name = KafkaConnectLog4j2

    # 루트 로거 설정 (stdout 앱펜더 연결)
    rootLogger.level = info
    rootLogger.appenderRefs = stdout
    rootLogger.appenderRef.stdout.ref = STDOUT

    # 콘솔 앱펜더 정의 (Log4j 2.x 방식)
    appender.console.type = Console
    appender.console.name = STDOUT
    appender.console.layout.type = PatternLayout
    appender.console.layout.pattern = [%d] %p %m (%c)%n

    # 디버깅에 유용한 오픈소스 커넥터 로거 설정 (선택 사항)
    logger.connect.name = org.apache.kafka.connect
    logger.connect.level = info
    
    logger.debezium.name = io.debezium
    logger.debezium.level = info

- Kafka Connector

Kafka Connect 위에서 동작하는 플러그인 형태의 Component
Source 시스템의 데이터를 Kafka로 가져오거나, Kafka 데이터를 외부 시스템으로 전달하는 역할을 수행

1. Source Connector (Producer 역할)

Debezium PostgreSQL Connector를 사용
Source PostgreSQL의 Snapshot 및 WAL 변경사항을 읽어 Kafka Topic으로 전송

  • pg-source-connector.yaml
apiVersion: kafka.strimzi.io/v1
kind: KafkaConnector
metadata:
  # 1. 이름 변경: 기존 커넥터와 다른 이름으로 배포하여 스냅샷 트리거
  name: postgres-source-v2
  namespace: kafka
  labels:
    strimzi.io/cluster: cdc-connect
spec:
  class: io.debezium.connector.postgresql.PostgresConnector
  tasksMax: 1
  config:
    # 1. 대상 데이터 베이스
    database.hostname: sourceendpoint
    database.port: port
    database.user: app_user
    database.password: "1234"
    database.dbname: source_db

    # 2. Debezium이 모니터링할 스키마와 테이블 지정
    schema.include.list: "public,users"
    table.include.list: "users.org_user"

    # 3. Debezium이 생성하는 Kafka 토픽의 접두사
    topic.prefix: pg-source-initial

    # 4. Debezium CDC 및 스냅샷 핵심 설정
    plugin.name: pgoutput
    snapshot.mode: initial              # 🔥 스냅샷 강제 시작
    
    # 🔥 DB 내 슬롯과 퍼블리케이션 이름 변경 (새로 생성)
    publication.name: debezium_publication_v2
    slot.name: debezium_slot_v2
2. Sink Connector (Consumer 역할)

Debezium JDBC Sink Connector를 사용
Kafka Topic의 CDC 이벤트를 소비하여 Target PostgreSQL에 반영

  • pg-sink-connector.yaml
apiVersion: kafka.strimzi.io/v1
kind: KafkaConnector
metadata:
  name: postgres-sink-connector-v2
  namespace: kafka
  labels:
    strimzi.io/cluster: cdc-connect      # Strimzi KafkaConnect 클러스터 이름과 일치시킵니다.
spec:
  # 🔥 [핵심 교정] v6 이미지에 새로 내장한 Debezium 순정 JDBC 싱크 클래스를 조준합니다.
  class: io.debezium.connector.jdbc.JdbcSinkConnector
  tasksMax: 1
  config:
    # 1. 구독할 카프카 토픽 이름
    topics: pg-source-initial.users.org_user

    # 2. 대상 DB 연결 정보 (.user가 아니라 .username을 사용합니다)
    connection.url: jdbc:postgresql://targetendpoint:port/target_db
    connection.username: app_user
    connection.password: "1234"

    table.naming.strategy: io.debezium.connector.jdbc.naming.DefaultTableNamingStrategy
    table.name.format: copy_schema.cp2_users

    # 🔥 3. Debezium 전용 실시간 CUD 동기화 및 UPSERT 설정
    # Debezium 싱크는 주소를 분석해 PostgreSQL Dialect를 스스로 완벽하게 찾아냅니다. (dialect.name 불필요)
    insert.mode: upsert
    primary.key.mode: record_key         # Confluent의 pk.mode 대신 사용합니다.
    primary.key.fields: user_id          # Confluent의 pk.fields 대신 사용합니다.
    delete.enabled: true                 # 원본 DB에서 DELETE 발생 시 타겟 DB 데이터도 자동 삭제

    # basic 모드는 테이블이 없으면 자동으로 create하고, 새 컬럼이 들어오면 alter table로 추가해 줍니다.
    schema.evolution: basic

    # 5. 자동 Dialect 우회 및 스키마 일치 설정 (JSON w/ Schemas 규격 파싱)
    key.converter: org.apache.kafka.connect.json.JsonConverter
    key.converter.schemas.enable: true
    value.converter: org.apache.kafka.connect.json.JsonConverter
    value.converter.schemas.enable: true

    # 6. 오프셋 설정 (latest 또는 earliest 중 인프라 환경에 맞춰 사용)
    # latest : snapshot 이후의 변경 데이터만 처리 (운영 환경에서 일반적으로 권장)
    # earliest : snapshot 이후의 모든 변경 데이터 재처리 (개발/디버깅 시 유용)
    consumer.auto.offset.reset: earliest

    # 7. SMT 설정 - 토픽의 거대한 전체 JSON 내용 중 'after' 구조체만 깔끔하게 파내기
    transforms: ExtractAfter
    transforms.ExtractAfter.type: org.apache.kafka.connect.transforms.ExtractField$Value
    transforms.ExtractAfter.field: after

Kafka strimzi 설치

kubectl create namespace kafka

kubectl apply -f https://github.com/strimzi/strimzi-kafka-operator/releases/download/0.45.0/strimzi-cluster-operator-0.45.0.yaml -n kafka

위 리소스 설치

...

실습 결과

  • Kafka Connect/Connetor Check

  • source db insert query

  • topic check

  • topic-message event check

  • target db select query

=> 똑같이 들어감을 확인할 수 있다.

profile
WannaB.E/D.E

0개의 댓글