grep

Backend

PostgreSQL to ES: (1) Kafka Connect CDC 파이프라인 구성

finn.h카카오

2025년 10월 10일

원문에서 보기 ↗

안녕하세요. 클린플랫폼에서 컨텐츠 모니터링 시스템을 개발하는 핀입니다.

저희 시스템은 약 10년 동안 운영되면서 점차 레거시가 쌓였습니다. 그중에 하나가 PostgreSQL의 데이터를 ES로 동기화하는 파이프라인이었습니다. 이번 글에서는 복잡했던 동기화 파이프라인을 대체한 Kafka Connect CDC 파이프라인을 소개합니다.

1. 문제 상황과 대안 선정

1.1 기존 파이프라인의 복잡성

기존 파이프라인에서는 아래 순서로 동기화가 진행되었습니다.

  1. PostgreSQL에 큐로 사용하는 테이블을 두어 ES에 인덱싱할 데이터를 저장

  2. 배치 앱에서 큐 테이블을 조회하여 유입 파이프라인으로 전송

  3. 유입 파이프라인에서 ES로 갈 데이터를 필터링하여 RabbitMQ로 전송

  4. 배치 앱에서 RabbitMQ 메시지를 컨슘하여 ES에 인덱싱

위 과정에서 언급된 유입 파이프라인은 저희 시스템에 연동된 서비스가 보낸 데이터를 처리하는 곳입니다. 유입된 데이터를 파싱하여 DB에 저장하고 있습니다. 굳이 ES에 인덱싱할 데이터가 거쳐 갈 필요는 없는 곳으로, 오히려 장애 지점만 늘리고 있습니다. 과거에는 나름의 이유를 갖고 이와 같이 설계되었지만 이제는 불필요하게 복잡한 설계로 남게 되었습니다. 그래서 문제가 발생했을 때 원인 파악에 어려움이 있었습니다. 기존 설계를 유지할 필요는 없었기 때문에 새로운 방식으로 파이프라인을 구성하기로 했습니다.

1.2 대안 선정

동기화 파이프라인을 새로 구성하기 위해 고민했던 툴은 Logstash, NiFi, PGSync, Kafka Connect가 있었습니다. 주로 PostgreSQL의 CDC(Change Data Capture) 이벤트를 캡처하는 툴을 찾았습니다. 그중에서 Kafka Connect를 선택했습니다. Kafka가 익숙하여 유지보수에 용이할 거라 판단했고, 가용성과 확장성 등의 이점을 그대로 가져올 수 있기 때문입니다.

2. 배경지식

이번 섹션에는 미리 알고 있으면 글을 읽을 때 도움이 되는 개념들을 작성했습니다. Kafka Connect의 구성 요소와 PostgreSQL의 Logical Decoding을 알고 계신다면 다음 섹션으로 건너뛰셔도 됩니다.

2.1 Kafka Connect 구성 요소

Kafka Connect

Kafka Connect는 자주 반복되는 카프카 기반의 데이터 작업을 실행하기 위한 JVM 프로세스입니다. Connector와 Task를 실행하고 관리하는 역할을 합니다. Connect는 단일 모드 또는 분산 모드로 실행될 수 있습니다. 단일 모드에서는 하나의 프로세스만 실행되기 때문에 장애가 발생하면 복구가 어려울 수 있습니다. 분산 모드는 여러 Connect 프로세스를 하나의 클러스터로 묶어서 관리합니다. 그래서 특정 프로세스에 장애가 발생해도 다른 프로세스에서 복구될 수 있습니다.

Kafka Connector

Kafa Connector는 수행할 작업을 구체적으로 정의하는 일종의 템플릿입니다. 어떤 작업을 할지에 대해 선언적으로 설정하여 배포할 수 있습니다. Connector 종류에는 Source Connector와 Sink Connector가 있습니다. Source Connector는 특정 소스에 있는 데이터를 Kafka로 전송합니다. 이와 반대로 Sink Connector는 Kafka 메시지를 컨슘하여 목적지로 전송합니다.

Task

Task는 Connector의 설정에 따라 실제로 데이터를 처리하는 스레드입니다. Connect 프로세스 내에서 실행됩니다.

Transforms

Transforms는 단일 메시지에 대해 변환할 수 있는 Connector의 설정입니다. 이미 구현된 플러그인을 사용할 수도 있고, 직접 커스텀하게 구현하여 사용할 수도 있습니다. Source Connector에 적용하면 Kafka 토픽에 넣기 전에 실행되고, Sink Connector에서는 목적지에 데이터를 넣기 전에 실행됩니다.

2.2 PostgreSQL의 Logical Decoding

PostgreSQL의 논리적 디코딩(Logical Decoding)은 데이터베이스의 변경 사항(WAL)을 외부에서 쉽게 이해할 수 있는 형식으로 변환하여 스트리밍하는 기술입니다. 이 스트리밍 과정은 구독자가 Replication Slot을 통해 특정 Publication을 구독하면 시작됩니다. 이후 Publication에 정의된 테이블에서 변경 사항이 발생하면, PostgreSQL은 WAL에 기록된 내용을 Output Plugin으로 디코딩 후 구독자에게 전달합니다.

Publication

데이터베이스 변경 사항을 발행할 테이블 목록입니다.

Replication Slot

데이터를 어디까지 읽었는지 기록하기 위한 책갈피입니다. 구독자가 데이터 변경 사항을 읽다가 연결이 끊겨도 Slot을 통해 끊긴 지점부터 변경 사항을 읽어올 수 있습니다.

Output Plugin

데이터베이스 변경 사항(WAL)을 특정 형식(json, sql 등)으로 변환하는 모듈입니다.

3. CDC 파이프라인 미리보기

3.1 파이프라인 구조

기본적인 파이프라인의 구조는 위와 같습니다. PostgreSQL의 데이터 변경 사항을 준실시간으로 Elasticsearch에 동기화합니다. 중간에 있는 Python 앱에서는 조직 내 비즈니스 로직에 맞게 데이터를 변환하고 있습니다.

Debezium은 Source Connector의 하나로 PostgreSQL의 CDC 이벤트를 Kafka로 전송합니다. ES Sink Connector는 CDC 이벤트 기반의 Kafka 메시지를 컨슘하여 ES에 인덱싱합니다. Debezium과 ES Sink Connector에서 필요로 하는 값들을 요구 사항에 맞게 설정하여 파이프라인을 구축하는 것이 주된 작업입니다.

3.2 데이터 흐름

위 그림은 파이프라인 구성 후의 데이터 흐름입니다. 실제로는 더 복잡하지만 구조를 쉽게 이해할 수 있도록 간단히 나타냈습니다.

4. CDC 파이프라인 구성

이제 본격적으로 CDC 파이프라인을 구성해 보겠습니다. Kafka Connect와 Kafka Connector의 설정값들을 구체적으로 설명합니다. 배포 방식은 뒤 섹션에서 설명하기 때문에, 이번 섹션에서는 어떤 설정값이 필요한지만 살펴보겠습니다.

4.1 Kafka Connect 환경 구성

Connector를 배포하기 전에 우선 Connect의 환경을 구성해야 합니다.

bootstrap.servers:  my-kafka-server:9092  
group.id:  my-connect  
  
config.storage.topic:  kakao-connect-config  
offset.storage.topic:  kakao-connect-offset  
status.storage.topic:  kakao-connect-status  
  
plugin.path:  /opt/kafka/plugins

4.2 Kafka Connector 설정

설정이 필요한 Connector는 Debezium과 ES Sink Connector입니다. 각 Connector를 설정하고 Connect 프로세스에 배포하면 Task 스레드가 실행됩니다.

4.2.1 Debezium(PostgreSQL Source Connector) 설정

PostgreSQL 관련 설정

어떤 DB에서 어떤 테이블의 변경 사항을 스트리밍 해올지 설정해야 합니다.

database.hostname:  my-db-host
database.port:  5555
database.dbname:  mydb
database.user:  ''
database.password:  ''

plugin.name:  pgoutput
publication.name:  my_publication
publication.autocreate.mode:  disabled
slot.name:  my_slot

table.include.list:  kakao.article,kakao.comment

Kafka 관련 설정

어떤 Kafka 토픽으로 CDC 이벤트를 전송할지 설정해야 합니다. 사용하는 Kafka 토픽은 기본적으로 .. 포맷입니다. 그래서 Kafka 토픽이 테이블 이름에 의존적입니다. 테이블 구분 없이 하나의 토픽만 사용하기 위해 transforms 설정을 추가했습니다.

topic.prefix:  kakao

# 단일 메시지 변환 설정
transforms:  Reroute
transforms.Reroute.type:  io.debezium.transforms.ByLogicalTableRouter
transforms.Reroute.topic.regex:  (.*)\.(.*)\.(article|comment)
transforms.Reroute.topic.replacement:  my-cdc-event

4.2.2 ES Sink Connector 설정

Kafka 관련 설정

어떤 Kafka 토픽에서 메시지를 꺼내고 어떻게 해석할지 설정해야 합니다.

topics:  my-es-index

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

schema.ignore: true

ES 관련 설정

connection.url: my-es-url:9200
behavior.on.null.values:  delete

5. 배포

Kafka Connect를 배포하는 가장 기본적인 방법은 Confluent에서 제공하는 도커 이미지를 사용하는 것입니다. 도커 컨테이너를 띄운 후 Kafka Connect에서 지원하는 REST API로 Connector를 생성하고 삭제하고 상태를 조회할 수 있습니다.

저희 조직에서는 대부분의 앱을 쿠버네티스에 올리고 GitOps로 관리하고 있습니다. Kafka Connect도 같은 방식으로 관리할 수 없을지 고민하던 중 Strimzi라는 툴을 발견했습니다. Strimzi는 쿠버네티스 환경에서 Kafka Connect를 배포할 수 있게 해줍니다. 이번 섹션에서는 Strimzi로 Connect와 Connector를 배포하는 방식을 설명합니다.

5.1 Strimzi 설치

$  helm  install  strimzi-cluster-operator  
oci://quay.io/strimzi-helm/strimzi-kafka-operator

$  kubectl  get  deployment
NAME  READY  UP-TO-DATE  AVAILABLE  AGE
strimzi-cluster-operator  1/1  1  1  3m57s

Strimzi는 helm 커맨드로 간단하게 설치 가능합니다.

5.2 Kafka Connect 배포

Kafka Connect 이미지 준비

FROM  quay.io/strimzi/kafka:0.29.0-kafka-3.0.0
USER  root:root
COPY  ./my-plugins/  /opt/kafka/plugins/

쿠버네티스에 배포할 Connect 이미지를 빌드하기 위해 위와 같이 Dockerfile을 작성했습니다. 베이스 이미지 태그는 사용할 Kafka 버전에 맞춰 설정해야 합니다. my-plugins에는 사용할 Connector 또는 Transforms의 JAR 파일이 들어있어야 합니다. Debezium과 ES Sink Connector 플러그인도 Confluent Hub에서 다운받아서 이미지에 넣어줬습니다.

커스텀 리소스 작성

apiVersion:  kafka.strimzi.io/v1beta2
kind:  KafkaConnect
metadata:
    name:  my-connect
    annotations:
	    strimzi.io/use-connector-resources:  true
spec:
    version:  3.0.0
    replicas:  2
    bootstrapServers:  my-kafka-server:9092
    image:  my-kafka-connect:latest
    config:  
	    group.id:  my-connect
	    config.storage.topic:  kakao-connect-config  
	    offset.storage.topic:  kakao-connect-offset
	    status.storage.topic:  kakao-connect-status  
	    plugin.path:  /opt/kafka/plugins

이제 Strimzi의 커스텀 리소스를 통해 Connect를 배포할 수 있습니다. 일부 필드만 살펴보겠습니다.

배포

$  kubectl  apply  -f  connect.yaml
kafkaconnect.kafka.strimzi.io/my-connect  created

$  kubectl  get  KafkaConnect
NAME  DESIRED  REPLICAS  READY
my-connect  2  True

$  kubectl  get  pods
NAME  READY  STATUS  RESTARTS  AGE
my-connect-connect-0  1/1  Running  0  33s
my-connect-connect-1  1/1  Running  0  33s

배포하면 파드 두개가 뜬 것을 볼 수 있습니다. 이제 Connector를 배포하기 위한 준비가 되었습니다.

5.3 Debezium 배포

커스텀 리소스 작성

apiVersion:  kafka.strimzi.io/v1beta2
kind:  KafkaConnector
metadata:
	name:  debezium
	labels:
		strimzi.io/cluster:  my-connect
spec:
	class:  io.debezium.connector.postgresql.PostgresConnector
	tasksMax:  1
	config:
		database.hostname:  my-db-host
		database.port:  5555
		database.dbname:  mydb
		database.user:  ''
		database.password:  ''

	    plugin.name:  pgoutput
	    publication.name:  my_publication
	    publication.autocreate.mode:  disabled
	    slot.name:  my_slot

	    table.include.list:  kakao.article,kakao.comment

	    topic.prefix:  kakao

	    transforms:  Reroute
	    transforms.Reroute.type:  			
	    io.debezium.transforms.ByLogicalTableRouter
	    transforms.Reroute.topic.regex:  (.*)\.(.*)\.(article|comment)
	    transforms.Reroute.topic.replacement:  my-cdc-event

Kafka Connect의 매니페스트를 작성한 것과 크게 다르지 않습니다. 일부만 살펴보겠습니다.

배포

$  kubectl  apply  -f  debezium.yaml
kafkaconnector.kafka.strimzi.io/debezium  created

$  kubectl  get  KafkaConnector
NAME  CLUSTER  CONNECTOR  CLASS  MAX  TASKS  READY
debezium  my-connect PostgresConnector  1  True

배포하면 KafkaConnector 리소스로 조회할 수 있습니다. Connect 파드의 로그를 보면 Task가 실행되는 것을 확인할 수 있습니다.

5.4 ES Sink Connector 배포

커스텀 리소스 작성

apiVersion:  kafka.strimzi.io/v1beta2
kind:  KafkaConnector
metadata:
    name:  es-sink-connector
    labels:
	    strimzi.io/cluster:  my-connect
spec:
    class:  io.confluent.connect.elasticsearch.ElasticsearchSinkConnector
    tasksMax:  4
    config:
	    topics:  my-es-index

	    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
	    schema.ignore: true

	    connection.url: my-es-url:9200
	    behavior.on.null.values:  delete

Debezium과 설정값만 다를 뿐 구조는 동일합니다.

배포

$  kubectl  apply  -f  es_sink_connector.yaml
kafkaconnector.kafka.strimzi.io/es-sink-connector  created

$  kubectl  get  KafkaConnector
NAME  CLUSTER  CONNECTOR  CLASS  MAX  TASKS  READY
es-sink-connector  my-connect  ElasticsearchSinkConnector  5  True

배포 후 KafkaConnector 리소스로 확인할 수 있습니다.

6. 정리

지금까지 Kafka Connect를 활용하여 CDC 파이프라인을 구성해 봤습니다. 데이터를 어떻게 전달하고 읽어올지에 따라 Connector의 설정값과 파이프라인 구성 방식이 달라질 수 있습니다. 예를 들면 위에 구성한 Debezium에서 불필요한 데이터를 제거한 뒤 Kafka로 보낼 수도 있습니다. 조직과 환경의 요구 사항에 맞춰 적절하게 설정하시면 될 것 같습니다.

성과 측면에서는, 우선 간헐적으로 발생하던 정합성 문제가 모두 해소되었습니다. 또한 불필요한 의존성을 제거하고 파이프라인을 단순화했기 때문에 문제가 발생해도 원인을 쉽게 찾을 수 있게 됐습니다. ES 인덱싱 속도는 기존 대비 최대 47% 향상되기도 했습니다.

다음 글에서는 CDC 파이프라인을 구성하면서 겪었던 트러블슈팅을 중점적으로 소개합니다. 이번 글에서 구성한 파이프라인을 적용했을 때 어떤 문제가 발생했고 어떻게 해결했는지 공유할 예정입니다.

참고 링크