오랜만에 velog를... 최근에는 notion 에 AWS-DEA 내용을 정리하고, 뭐 이래저래 바빴구만...
오늘은 CDC 개념 및 실습이다.
계속 Debezium 과 메시징 큐를 이용해서 이기종 DB 동기화를 어느정도 해보고 싶었어서, 개념 및 기술 플로우를 이해하고 자리잡으려고 시도해보았다.
해당 내용은 제가 정리했던 것인데, 해당 내용을 보시는게 더 도움 될 수 있습니다.


CDC의 데이터 원천(Source) 역할 수행
역할
주요 설정
wal_level = logical 확인
SHOW wal_level;
replica: 물리 복제용 WAL에 기록
INSERT INTO users ~~=> "디스크 블록 1234 변경"
logical: 논리적 변경 이벤트 까지 WAL에 기록
INSERT INTO users ~~=> "User 테이블 INSERT 발생"
Publication 설정
Publication 은?
어떤 테이블의어떤 변경 이벤트를 외부에 공개할 것인지를 정의하는 객체
publication_name은 kafka connector 생성 시, 값 option 설정이 다음과 같이 있음.
publication.name : debezium_publication_v2
connector.yaml 생성 후, 확인 쿼리
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 이벤트 수집 가능 여부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 생성 후, 확인 쿼리
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 위치Kafka Topic으로 전달된 CDC 이벤트를 최종 저장하는 데이터 베이스
CDC 이벤트를 저장하고 전달하는 메시지 브로커 역할 수행

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: {}
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
Kafa와 외부 시스템을 연결하기 위한 데이터 통합 플랫폼
Connector를 실행하고 관리하는 Runtime 역할을 수행
역할
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
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 Connect 위에서 동작하는 플러그인 형태의 Component
Source 시스템의 데이터를 Kafka로 가져오거나, Kafka 데이터를 외부 시스템으로 전달하는 역할을 수행
Debezium PostgreSQL Connector를 사용
Source PostgreSQL의 Snapshot 및 WAL 변경사항을 읽어 Kafka Topic으로 전송
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
Debezium JDBC Sink Connector를 사용
Kafka Topic의 CDC 이벤트를 소비하여 Target PostgreSQL에 반영
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
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

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