Backend
PostgreSQL to ES: (2) Kafka Connect 트러블슈팅
finn.h카카오
2025년 10월 10일
원문에서 보기 ↗안녕하세요. 클린플랫폼에서 컨텐츠 모니터링 시스템을 개발하는 핀입니다.
지난 글에서 저희 조직 내의 동기화 파이프라인을 대체한 Kafka Connect 기반의 CDC 파이프라인을 소개했습니다. 이번 글에서는 제가 직면했던 구체적인 문제들과 이를 해결한 방법을 말씀드리겠습니다.
1. CDC 파이프라인 구조

이전 글에서 구성한 CDC 파이프라인은 위와 같습니다. 대부분의 문제는 Debezium과 ES Sink Connector의 설정을 추가하거나 처리하는 데이터 구조를 변경하여 해결할 수 있었습니다. 그래서 이번 글에서도 파이프라인 구조 자체는 변경되지 않을 예정입니다.
2. Debezium 트러블슈팅
2.1 PostgreSQL의 Disk Full
2.1.1 Debezium 기본 동작

Debezium이 테이블의 변경 사항을 캡처하는 과정은 위 그림과 같습니다. 우선 구독하는 테이블의 변경 사항을 스트리밍 해온 후, CDC 이벤트를 Kafka로 전송합니다. 그리고 WAL(Write-Ahead Log)의 LSN(Log Sequence Number)를 Kafka에 전송합니다. WAL의 어느 지점까지 처리했는지 기록하기 위한 작업입니다. 마지막으로 PostgreSQL에게 LSN 피드백을 보냅니다. 그럼 PostgreSQL은 Debezium이 사용하는 Slot에 LSN을 업데이트합니다.
2.1.2 문제 상황

문제는 캡처하려는 테이블의 변경 사항이 거의 없는 경우 발생합니다.
Debezium 입장에서는 캡처하는 테이블의 변경 사항이 있을 때만 이후의 프로세스를 진행합니다. 만약 변경 사항이 없다면 PostgreSQL에게 LSN 피드백을 보내지 않게 되고 Slot이 갖는 LSN이 갱신되지 않습니다.
PostgreSQL은 클러스터 내의 모든 변경 사항을 WAL에 기록하면서 오래된 WAL 파일부터 삭제합니다. 그런데 Slot이 필요로 하는 WAL 파일이 있다면 해당 파일 이후부터는 삭제하지 않고 유지합니다. WAL은 여러 Slot이 공유하는 스트림 파일입니다. 그래서 특정 Slot의 LSN이 갱신되지 않는다면 WAL 파일이 쌓이면서 Disk Full 문제가 발생할 수 있습니다.
2.1.3 해결 방안 - heartbeat 설정

Debezium에서는 Disk Full 문제를 해결하기 위한 heartbeat 설정을 지원하고 있습니다. heartbeat 설정을 추가하면 위 그림처럼 테이블의 변경 사항이 없더라도 일정 주기로 타이머가 작동하여 Kafka로 heartbeat 메시지를 전송합니다. 이후 동작은 변경 사항이 발생했을 때와 동일합니다. LSN을 Kafka에 기록하고 PostgreSQL에게도 피드백을 보냅니다. Slot의 LSN이 업데이트되면 WAL 파일이 쌓이던 문제를 해결할 수 있습니다.
Debezium에 추가하는 heatbeat 설정은 아래와 같습니다.
heartbeat.interval.ms: 60000
topic.heartbeat.prefix: my-heartbeat
topic.prefix: kakao
topic.delimiter: -
-
heartbeat.interval.ms
- heartbeat 메시지를 전송할 주기입니다.
-
topic.heartbeat.prefix
- heartbeat 메시지를 전송할 Kafka 토픽의 prefix입니다.
- 토픽 이름은 . 포맷으로 설정됩니다. 위와 같이 topic.delimiter까지 설정하면 heartbeat에 사용되는 토픽은 my-heartbeat-kakao입니다.
2.2 Debezium의 LSN 유실
2.2.1 Debezium 재시작 시 정상 동작

Debezium을 재시작하면 Kafka에서 마지막으로 기록한 WAL의 LSN을 가져옵니다. 그 후 PostgreSQL에 LSN 이후 데이터를 요청하고 변경 데이터에 대한 스트리밍이 시작됩니다.
2.2.2 문제 상황

문제는 Debezium이 PostgreSQL에 LSN 이후 데이터를 요청할 때 발생했습니다. 간헐적으로 PostgreSQL에서 LSN이 유효하지 않다고 판단하고 있었습니다. 거의 발생하지 않는 문제였으나 데이터 유실이 발생할 수 있어서 치명적이었습니다. PostgreSQL 담당 크루와 원인 파악에 나섰지만, 문제가 재현되지 않아 결국 원인을 파악하지 못했습니다. 결국 문제 자체를 해결하는 대신, 문제가 생겼을 때 데이터 유실이 발생하지 않는 방법을 찾기로 했습니다.
2.2.3 해결 아이디어
PostgreSQL의 WAL 파일만으로는 특정 테이블의 전체 데이터를 알 수가 없습니다. 그래서 보통 Debezium이 최초로 실행될 때, 캡처하려는 테이블의 모든 데이터를 읽고 Kafka로 전송하는 초기화 과정을 거칩니다. 이 작업을 스냅샷(Snapshot)이라고 합니다. 이때 Kafka로 보내는 메시지는 실시간 변경 데이터를 처리하는 스트리밍 단계의 메시지 구조와 동일합니다. 그래서 메시지를 처리하는 입장에서는 스냅샷으로 들어온 메시지와 실시간 변경으로 들어온 메시지를 동일하게 처리할 수 있습니다. 스냅샷이 끝난 후에는 스트리밍 모드로 전환합니다.
스냅샷을 찍으면 데이터 유실을 막을 수 있습니다. 하지만 기본적으로 제공되는 스냅샷 쿼리는 테이블의 모든 데이터를 조회합니다. 테이블의 로우 개수가 1억 개이면 이를 모두 조회한다는 뜻입니다. 데이터가 많으면 현실적으로 어려운 방식이고, 굳이 모든 데이터를 조회할 필요도 없습니다. 그런데 만약 최신 데이터만 조회하도록 스냅샷 쿼리를 변경할 수 있다면 빠르게 스냅샷을 실행할 수 있습니다. 다행히 Debezium에서는 스냅샷 쿼리를 직접 설정할 수 있도록 지원하고 있습니다.
2.2.4 해결 방안 1 - snapshot.select.statement.overrides 설정
snapshot.mode: when_needed
snapshot.select.statement.overrides: kakao.article,kakao.comment
snapshot.select.statement.overrides.kakao.article: SELECT * FROM
kakao.article WHERE updated_at >= NOW() - INTERVAL '5 minutes'
snapshot.select.statement.overrides.kakao.comment: SELECT * FROM
kakao.comment WHERE updated_at >= NOW() - INTERVAL '5 minutes'
가장 간단하게 스냅샷 쿼리를 오버라이딩하는 방법입니다. 위 설정을 Debezium에 추가하면 됩니다.
-
snapshot.mode
- 언제 스냅샷을 실행할지 설정합니다.
- when_needed로 설정하면 Debezium이 Kafka에 기록한 LSN을 가져오지 못하거나, 가져온 LSN이 PostgreSQL에서 유효하지 않을 때 스냅샷을 실행합니다. 현재 문제는 후자에 해당합니다.
-
snapshot.select.statement.overrides
- 스냅샷 쿼리를 적용할 테이블을 설정합니다.
-
snapshot.select.statement.overrides. < schema>.< table>
- 테이블에 적용할 스냅샷 쿼리를 설정합니다.
- 현재 설정된 쿼리는 updated_at 칼럼을 기준으로 5분 전 데이터만 가져오는 쿼리입니다.
이 방식은 간단하다는 장점이 있지만, Debezium이 여러 테이블을 캡처하고 있다면 테이블마다 일일이 설정해줘야 한다는 번거로움이 있습니다. 특히 파티션 테이블을 사용 중이라면 부모 테이블로 적용할 수 있는지 확인이 필요합니다.
2.2.5 해결 방안 2 - SnapshotQuery SPI 구현
스냅샷 쿼리를 직접 작성하기 위해 SnapshotQuery SPI를 구현할 수 있습니다. 기존 코드를 참고하면 작성이 어렵지 않습니다. 아래 코드는 실제로 사용 중인 구현체입니다.
import io.debezium.annotation.ConnectorSpecific;
import io.debezium.connector.postgresql.PostgresConnector;
import io.debezium.snapshot.spi.SnapshotQuery;
@ConnectorSpecific(connector = PostgresConnector.class)
public class CustomIntervalSnapshotQuery implements SnapshotQuery {
private String date_column;
private String interval;
@Override
public String name() {
return "custom-interval-snapshot-query";
}
@Override
public void configure(Map properties) {
this.interval = (String) properties.get("snapshot.query.interval.minutes");
this.date_column = (String) properties.get("snapshot.query.date.column");
}
@Override
public Optional snapshotQuery(String tableId, List
snapshotSelectColumns) {
String whereClause = "WHERE " + this.date_column + " >= NOW() -
INTERVAL '" + this.interval + " minutes'";
return Optional.of(snapshotSelectColumns.stream()
.collect(Collectors.joining(", ", "SELECT ", " FROM " +
tableId + " " + whereClause)));
}
}
구현해야 하는 메서드는 세 가지입니다.
-
name()
- 스냅샵 쿼리의 이름을 반환합니다. 추후 Debezium 설정에 사용됩니다.
-
configure()
- Debezium 설정으로부터 주입받아 사용할 필드를 설정합니다.
- interval과 date_column을 주입받아서 어떤 칼럼을 기준으로 몇 분 전 데이터를 가져올지 쿼리하도록 했습니다.
-
snapshotQuery()
- 실제 쿼리문을 반환합니다.
- 주입받은 필드 값으로 실제 쿼리를 완성합니다.
위 SPI 구현체를 작성한 뒤에 Kafka Connect가 사용할 수 있도록 해야 합니다. Connector 플러그인을 Connect 플러그인 경로에 넣어줬던 것과 동일하게, SPI 구현체도 JAR로 빌드한 후 플러그인 경로에 넣어주면 됩니다. 이때 ServiceLoader가 구현체를 등록할 수 있도록 META-INF/services/에 설정 파일을 추가한 뒤 빌드해야 합니다.
이제 Debezium에 아래와 같이 설정하여 직접 구현한 커스텀 쿼리를 사용할 수 있습니다.
snapshot.mode: when_needed
snapshot.query.mode: custom
snapshot.query.mode.custom.name: custom-interval-snapshot-query
snapshot.query.date.column: updated_at
snapshot.query.interval.minutes: 5
-
snapshot.query.mode
- 직접 구현한 스냅샷 쿼리를 사용하기 위해 custom으로 설정합니다.
- 기본값은 select_all입니다.
-
snapshot.query.mode.custom.name
- name() 메서드에서 반환하는 값으로 설정합니다.
-
snapshot.query.date.column
- 어떤 칼럼을 조건으로 사용할지 설정합니다.
-
snapshot.query.interval.minutes
- 몇 분 전 데이터를 가져올지 설정합니다.
SnapshotQuery SPI 구현 방식은 간단하지 않지만, 여러 Debezium에서 테이블에 구애받지 않고 사용할 수 있다는 장점이 있습니다. 현재는 이 방식으로 스냅샷 쿼리를 적용하고 있습니다. 적용 이후 데이터 유실이 발생하지 않았습니다.
3. ES Sink Connector 트러블슈팅
3.1 ES 인덱스를 날짜 기반으로 분리
3.1.1 문제 상황
ES Sink Connector는 컨슘하는 Kafka 토픽 이름을 ES의 인덱스로 사용합니다. 직접 ES 인덱스를 설정하고 싶다면 Transforms를 사용하면 됩니다. 공식 문서에서 안내하는 Transforms는 RegexRouter와 TimestampRouter입니다.
설정하고 싶은 ES 인덱스는 데이터 생성 날짜가 인덱스 접미사에 들어간 형식이었습니다(ex. article-2025_09). 아래와 같이 TimestampRouter를 설정하면 날짜 형식이 들어간 인덱스를 사용할 수 있었습니다.
transforms: router
transforms.router.type: org.apache.kafka.connect.transforms.TimestampRouter
transforms.router.topic.format: article-${timestamp}
transforms.router.timestamp.format: YYYY_MM
-
transforms.router.type
- 사용할 Transforms 플러그인을 TimestampRouter로 설정합니다.
-
transforms.router.topic.format
- 사용할 ES 인덱스의 포맷입니다.
-
transforms.router.timestamp.format
- topic.format에서 {timestamp}에 들어갈 날짜 형식입니다.
그런데 과거에 생성된 데이터가 변경될 때 문제가 생겼습니다. 8월에 생성된 데이터는 article-2025_08 인덱스에 저장되는데, 해당 데이터가 9월에 변경되면 article-2025_09 인덱스에 중복으로 저장됐습니다. 생성 날짜가 들어간 인덱스로 수정되는 것이 요구 사항이었기 때문에 TimestampRouter로는 한계가 있었습니다. 다른 Transforms도 찾아봤지만 적합한 플러그인을 찾지 못했습니다.
3.1.2 해결 방안 - Transforms 구현
직접 Transforms를 구현하기로 했습니다.

우선 메시지 구조부터 위 그림처럼 수정했습니다. ES 인덱스를 Transforms 안에서 만드는 것보다 앞단(Python 앱)에서 전달하는 게 깔끔해 보였습니다. 그래서 ES 인덱스를 앞단에 있는 Python 앱에서 만들어서 메시지 키로 전달하고, Transforms에서는 메시지 키에서 필요한 값을 꺼내도록 했습니다. 기존에 메시지 키로 들어가던 도큐먼트 id도 함께 처리했습니다.

Transforms가 하는 일을 자세히 보면 위 그림과 같습니다. 실제로 구현할 부분은 단순히 메시지 키에 있는 인덱스를 토픽으로 바꿔치기하고, 도큐먼트 id는 메시지 키에 다시 넣는 것이었습니다. 이렇게 수정하는 이유는 토픽을 인덱스로 사용하고, 메시지 키를 도큐먼트 id로 사용하는 ES Sink Connector의 특징 때문입니다.
import org.apache.kafka.connect.transforms.Transformation;
public class ESRouter> implements
Transformation {
private static final String INDEX_FIELD = "index.field";
private static final String ID_FIELD = "id.field";
public static final ConfigDef CONFIG_DEF = new ConfigDef()
.define(ID_FIELD, ConfigDef.Type.STRING,
ConfigDef.NO_DEFAULT_VALUE, ConfigDef.Importance.HIGH, "ES doc id")
.define(INDEX_FIELD, ConfigDef.Type.STRING,
ConfigDef.NO_DEFAULT_VALUE, ConfigDef.Importance.HIGH, "ES index name");
@Override
public void configure(Map props) {
final SimpleConfig config = new SimpleConfig(CONFIG_DEF, props);
idField = config.getString(ID_FIELD);
indexField = config.getString(INDEX_FIELD);
}
@Override
private R apply(R record) {
final Map key = requireMap(record.key(), PURPOSE);
String newTopic = key.get(indexField).toString();
String newKey = key.get(idField).toString();
return record.newRecord(
newTopic,
record.kafkaPartition(),
record.keySchema(),
newKey,
record.valueSchema(),
record.value(),
record.timestamp(),
record.headers()
);
}
// 생략
}
기존의 간단한 Transforms 코드를 참고해서 구현했습니다. 코드를 살펴보면 Transforms 인터페이스를 구현하기 위해 필요한 메서드는 두 가지입니다.
-
configure()
- 필드 값을 주입받을 때 사용합니다.
- 위 코드에서는 Kafka 메시지 키에서 추출할 인덱스 필드와 id 필드를 주입받고 있습니다.
-
apply()
- 단일 메시지 변환이 진행되는 핵심 메서드입니다.
- Kafka 메시지 키에서 인덱스 필드와 id 필드를 꺼낸 후 각각 토픽과 메시지 키로 설정합니다. 실제로는 스키마를 처리하는 코드가 필요할 수 있지만, 처리하는 메시지 구조가 명확하기 때문에 간단히 나타냈습니다.
위 구현체를 JAR로 빌드한 후, Kafka Connect가 사용할 수 있도록 플러그인 경로에 넣어주면 됩니다. 그럼 아래와 같이 ES Sink Connector 설정을 추가할 수 있습니다.
transforms: esRouter
transforms.esRouter.type: io.kakao.kafka.connect.transforms.ESRouter
transforms.esRouter.index.field: index
transforms.esRouter.id.field: id
-
transforms.esRouter.type
- 구현한 클래스의 FQN(Fully Qualified Name)으로 설정합니다.
-
transforms.esRouter.index.field
- 메시지 키에 있는 인덱스 필드를 설정합니다.
-
transforms.esRouter.id.field
- 메시지 키에 있는 도큐먼트 id 필드를 설정합니다.
앞단(Python 앱)에서 ES 인덱스를 메시지 키에 넣고 ES Sink Connector에서 Transforms를 통해 변환하니 요구 사항에 맞게 ES 인덱스를 사용할 수 있었습니다. 그리고 ES Sink Connector가 컨슘하는 Kafka 토픽이 ES 인덱스로 사용되지 않기 때문에 Kafka 토픽을 자유롭게 사용할 수 있게 되기도 했습니다.
3.2 ES 도큐먼트의 버전 충돌
3.2.1 도큐먼트 버전 관리 방식

ES Sink Connector가 external 하게 버전 관리를 하면 위 그림과 같이 토픽의 offset을 도큐먼트의 버전으로 설정합니다. 현재 메시지 키를 기반으로 파티션을 배치하고 있어서 데이터 순서 보장이 되는 구조입니다. 참고로 internal 하게 관리하는 경우는 데이터 스트림에 저장하거나 Connector 설정에서 key.ignore를 true로 설정했을 때입니다. 현재 파이프라인에서는 해당하지 않아서 고려하지 않았습니다.
3.2.2 문제 상황

문제는 파티션을 추가했을 때 발생했습니다. 0번 파티션으로 가던 메시지가 새로 생성된 1번 파티션으로 배치될 수 있기 때문입니다. 생성된 파티션은 기존 파티션보다 offset이 작았기 때문에 도큐먼트의 업데이트가 무시되는 문제가 있었습니다.
3.2.3 해결 방안 - 버전 값 직접 설정

ES Sink Connector는 컨슘한 메시지의 헤더에서 버전으로 사용할 값을 설정할 수 있습니다. 그래서 앞단(Python 앱)에서 버전으로 사용할 값을 메시지 헤더에 담아 전달하도록 했습니다. ES Sink Connector에서는 헤더에 있는 값을 도큐먼트 버전 값으로 사용했습니다.
그런데 여기서 또 중요한 건 어떤 값을 헤더에 넣을지입니다. 저희는 WAL의 LSN(Log Sequence Number)을 사용했습니다. Debezium으로부터 받은 CDC 이벤트에는 LSN 값도 포함돼 있어서 쉽게 사용할 수 있습니다. LSN을 도큐먼트 버전 값으로 사용하여 PostgreSQL에서부터 보장된 순서를 ES에까지 적용할 수 있었습니다.
ES Sink Connector에 추가해야 하는 설정은 아래와 같습니다.
external.version.header: lsn
- external.version.header
- 메시지 헤더에서 버전으로 사용할 값을 설정합니다.
3.2.4 원인 파악 과정
이번 문제는 해결하는 시간보다 원인 파악에 오래 걸렸습니다. 보통은 Kafka Connect의 로그를 보고 원인을 파악하는데 처음에는 로그를 발견할 수 없었습니다. Kafka Connect의 로그 레벨을 debug로 변경하고 나서야 아래의 단서를 찾을 수 있었습니다. 참고로 Connect의 로그 레벨은 REST API로 변경할 수 있습니다.
2025-09-25 10:25:00,000 DEBUG [es-sink-connector-v2|task-0] Ignoring
EXTERNAL version conflict for operation INDEX on document '1' version 12 in
index 'my-index'. (io.confluent.connect.elasticsearch.ElasticsearchClient)
4. 정리
이번 글에서는 Kafka Connect 기반의 파이프라인을 구성하면서 겪었던 시행착오들을 정리했습니다. Kafka Connect를 실제로 적용해보니 생각보다 헤매는 경우가 많았습니다. Kafka에 대한 이해도 중요하지만 PostgreSQL과 ES 등 연결된 시스템에 대한 이해 또한 중요하다고 느꼈습니다. 예를 들어 PostgreSQL의 논리적 디코딩 개념을 제대로 알고 있었다면 문제의 원인을 파악할 때 조금은 수월했을 것 같습니다.
지난 글에서도 말씀드렸지만 Connector 설정값과 파이프라인 구조는 요구 사항과 취향에 맞춰서 다양하게 구성할 수 있습니다. 같은 문제를 다른 방식으로 풀 수 있기 때문에 소개드린 시행착오들은 하나의 사례로 참고해주시면 될 것 같습니다. Kafka Connect 도입을 검토 중이거나 적용 중에 문제가 생긴 분들에게 이번 글들이 조금은 도움이 되었으면 좋겠습니다.
읽어주셔서 감사합니다.