본문 바로가기
J-H-T

Kafka를 통한 MongoDB 데이터 스트리밍 싱크 구성 가이드 !!

by METAVERSE STORY 2026. 8. 2.
반응형

 

 

1. 결론부터 말씀드리면

Kafka에 들어온 데이터를 MongoDB로 실시간에 가깝게 저장하는 구성은 가능합니다.

일반적으로 다음과 같이 구성합니다.

```mermaid
flowchart LR
    A["업무 시스템"] --> B["Kafka Topic"]
    B --> C["Kafka Connect"]
    C --> D["MongoDB Sink Connector"]
    D --> E["MongoDB Primary"]
    E --> F["MongoDB Secondary 1"]
    E --> G["MongoDB Secondary 2"]
```

여기서 매우 중요한 부분은 다음입니다.

Kafka가 MongoDB Pod 3개에 각각 데이터를 넣는 것이 아닙니다.
Kafka Sink Connector는 MongoDB Replica Set에 접속해 Primary에 데이터를 저장하고, MongoDB가 Secondary로 복제합니다.

MongoDB 공식 Sink Connector는 Kafka 토픽의 레코드를 읽어 MongoDB Collection에 기록하는 용도로 제공됩니다. MongoDB Kafka Sink Connector


2. 먼저 “싱크”와 “복제”를 구분해야 합니다

이번 구성에는 서로 다른 두 가지 데이터 이동이 존재합니다.

구분데이터 흐름담당 시스템목적
Kafka 스트리밍 싱크 Kafka → MongoDB Kafka Connect Sink Connector Kafka 데이터를 MongoDB에 저장
MongoDB 복제 Primary → Secondary MongoDB Replica Set MongoDB 장애 대응 및 고가용성
스토리지 저장 MongoDB Pod → PVC Kubernetes·Storage 각 MongoDB 인스턴스의 데이터 보관

즉, 다음 세 가지는 동일한 기능이 아닙니다.

Kafka Sink
≠ MongoDB Replica Set
≠ Volume 공유

전체 흐름은 다음과 같습니다.

업무 시스템
   ↓ 데이터 발생
Kafka Topic
   ↓ 메시지 소비
MongoDB Sink Connector
   ↓ Primary에 Insert 또는 Upsert
MongoDB Primary
   ↓ MongoDB 내부 Replication
MongoDB Secondary 1
MongoDB Secondary 2

3. 권장 구성안

3.1 전체 구성

```mermaid
flowchart TB
    A["Producer 애플리케이션"] --> B["Kafka Cluster"]
    B --> C["Kafka Connect Cluster"]
    C --> D["MongoDB Sink Connector"]
    D --> E["MongoDB Replica Set Service"]

    E --> P["mongodb-0 / Primary / PVC-0"]
    E --> S1["mongodb-1 / Secondary / PVC-1"]
    E --> S2["mongodb-2 / Secondary / PVC-2"]
```

구성요소별 역할

구성요소역할
Producer 업무 데이터를 Kafka로 전송
Kafka Broker 메시지를 Topic에 저장하고 일정 기간 보관
Kafka Topic MongoDB로 보낼 데이터가 쌓이는 통로
Kafka Connect Connector를 실행하고 작업 상태와 Offset 관리
MongoDB Sink Connector Kafka 데이터를 MongoDB 문서로 변환해 저장
MongoDB Primary 실제 쓰기 요청 처리
MongoDB Secondary Primary의 데이터를 복제
PVC MongoDB Pod별 데이터 파일 저장

3.2 Kubernetes 구성 예시

MongoDB가 StatefulSet으로 구성되었다고 가정하겠습니다.

mongodb-0
├─ 역할: Primary 또는 Secondary
└─ PVC: mongodb-data-mongodb-0

mongodb-1
├─ 역할: Primary 또는 Secondary
└─ PVC: mongodb-data-mongodb-1

mongodb-2
├─ 역할: Primary 또는 Secondary
└─ PVC: mongodb-data-mongodb-2

Replica Set에서 Primary는 고정된 Pod 번호가 아닙니다.

예를 들어 처음에는 다음과 같을 수 있습니다.

mongodb-0 = Primary
mongodb-1 = Secondary
mongodb-2 = Secondary

mongodb-0에 장애가 발생하면 선거를 통해 다음처럼 바뀔 수 있습니다.

mongodb-0 = 장애
mongodb-1 = 새로운 Primary
mongodb-2 = Secondary

따라서 Sink Connector가 mongodb-0 Pod 하나만 바라보도록 구성하면 안 됩니다.

다음처럼 Replica Set 정보가 포함된 MongoDB 연결 주소를 사용해야 합니다.

mongodb://mongodb-0.mongodb:27017,
          mongodb-1.mongodb:27017,
          mongodb-2.mongodb:27017/
          targetdb?replicaSet=rs0

실제로는 한 줄로 작성합니다.

connection.uri=mongodb://mongodb-0.mongodb:27017,mongodb-1.mongodb:27017,mongodb-2.mongodb:27017/targetdb?replicaSet=rs0

Connector는 이 주소를 통해 현재 Primary를 찾아 데이터를 저장합니다. Primary가 변경되면 MongoDB 드라이버가 새로운 Primary를 다시 찾습니다. MongoDB Connector는 connection.uri를 사용해 MongoDB Cluster에 연결합니다. 운영 환경에서는 접속정보가 URI에 직접 노출되지 않도록 ConfigProvider나 Secret 사용을 권장합니다. MongoDB 연결 설정


4. 데이터가 저장되는 실제 과정

예를 들어 업무 애플리케이션이 다음 메시지를 Kafka로 보낸다고 가정하겠습니다.

Kafka Key

{
  "customerId": "CUST-1001"
}

Kafka Value

{
  "customerId": "CUST-1001",
  "name": "홍길동",
  "status": "ACTIVE",
  "updatedAt": "2026-07-31T10:00:00+09:00"
}

처리 과정은 다음과 같습니다.

  1. Producer가 메시지를 Kafka Topic에 저장합니다.
  2. Kafka가 메시지를 복제하고 보관합니다.
  3. Kafka Connect가 해당 메시지를 읽습니다.
  4. MongoDB Sink Connector가 메시지를 MongoDB 문서로 변환합니다.
  5. Connector가 현재 MongoDB Primary에 문서를 저장합니다.
  6. Primary가 MongoDB 내부 복제를 통해 Secondary에 전달합니다.
  7. 처리가 완료되면 Kafka Connect가 처리 위치인 Offset을 저장합니다.

MongoDB 결과는 다음과 같은 형태가 될 수 있습니다.

{
  "_id": "CUST-1001",
  "customerId": "CUST-1001",
  "name": "홍길동",
  "status": "ACTIVE",
  "updatedAt": "2026-07-31T10:00:00+09:00"
}

5. 가장 중요한 설계 결정: Insert인가, Upsert인가?

Kafka 메시지를 MongoDB에 어떻게 저장할지를 먼저 결정해야 합니다.

5.1 Insert 방식

메시지를 받을 때마다 새로운 문서를 추가합니다.

같은 customerId 메시지 3건
→ MongoDB 문서도 3건

예시:

{"customerId":"CUST-1001","status":"CREATED"}
{"customerId":"CUST-1001","status":"APPROVED"}
{"customerId":"CUST-1001","status":"COMPLETED"}

적합한 경우

  • 이벤트 이력 저장
  • 로그성 데이터
  • 센서 데이터
  • 거래 이력
  • 모든 변경 이력을 남겨야 하는 경우

주의점

Kafka가 메시지를 재처리하면 동일한 문서가 중복 저장될 수 있습니다.


5.2 Upsert 방식

같은 업무 Key를 가진 데이터가 이미 존재하면 수정하고, 없으면 새로 추가합니다.

CUST-1001 최초 수신
→ Insert

CUST-1001 변경 데이터 수신
→ Update 또는 Replace

MongoDB Connector의 기본적인 ReplaceOneDefaultStrategy는 _id가 같은 문서가 있으면 교체하고, 없으면 새 문서를 추가하는 Upsert 방식으로 사용할 수 있습니다. 업무 Key를 기준으로 교체하는 전략도 지원합니다. MongoDB Write Model 전략

적합한 경우

  • 고객의 현재 상태
  • 장비의 최신 상태
  • 주문의 현재 진행 상태
  • 계정 정보
  • 최종 결과 데이터

권장 사항

현재 상태를 저장하는 시스템이라면 대부분 다음 구조를 권장합니다.

Kafka Message Key = 업무 고유번호
MongoDB _id       = 업무 고유번호
저장 방식          = Upsert

예:

orderId    → MongoDB _id
customerId → MongoDB _id
deviceId   → MongoDB _id

이렇게 해야 같은 메시지가 다시 처리되더라도 문서가 무한히 중복 생성되는 문제를 줄일 수 있습니다.


6. 중복 데이터에 특히 주의해야 합니다

Kafka 기반 처리에서는 장애 복구나 재시도 과정에서 같은 메시지를 다시 처리할 가능성이 있습니다.

예를 들어:

1. Connector가 MongoDB에 저장
2. Connector가 Kafka Offset을 저장하기 전에 장애 발생
3. Connector 재기동
4. 같은 메시지를 다시 읽음
5. MongoDB에 다시 저장

Insert 방식이라면:

문서 1건 → 동일 문서 2건

Upsert 방식이라면:

동일한 _id 문서를 다시 갱신
→ 중복 문서 발생 가능성 감소

Kafka Connect의 정확히 한 번 처리 가능 여부는 Connector와 대상 시스템의 구현에 따라 달라집니다. 설정만으로 전체 구간의 완전한 Exactly-once가 자동 보장된다고 가정하면 안 됩니다. Kafka Connect 전달 보장

따라서 설계서에는 다음과 같이 표현하는 것이 안전합니다.

Kafka 메시지는 장애 및 재시도 과정에서 중복 처리될 수 있으므로, 업무 고유 Key를 MongoDB _id 또는 Unique Index로 사용하고 Upsert 기반의 멱등성 처리를 적용한다.

여기서 멱등성은 “같은 데이터를 여러 번 처리해도 최종 결과가 동일하게 유지되는 성질”입니다.


7. Kafka Key와 Partition 설계가 중요합니다

Kafka는 같은 Partition 안에서의 순서를 보장합니다. Topic 전체의 메시지 순서를 무조건 보장하는 것은 아닙니다.

따라서 동일한 업무 데이터는 같은 Kafka Key를 사용해야 합니다.

예:

Kafka Key = customerId
CUST-1001 데이터 → Partition 0
CUST-1001 데이터 → Partition 0
CUST-1001 데이터 → Partition 0

이렇게 하면 CUST-1001의 변경 메시지가 동일한 Partition에 들어가므로 순서를 관리하기 쉬워집니다.

반대로 Kafka Key를 매번 무작위로 생성하면 다음과 같은 문제가 발생할 수 있습니다.

10:00 고객 상태 = ACTIVE
10:01 고객 상태 = SUSPENDED

병렬 처리 때문에 MongoDB에 도착하는 순서가 바뀌면 최종 상태가 다시 ACTIVE가 될 수 있습니다.

권장 설계

데이터 종류권장 Kafka Key
고객 데이터 customerId
주문 데이터 orderId
장비 데이터 deviceId
GPU 상태 gpuId 또는 nodeId
사용자 세션 sessionId
업무 요청 requestId

8. 권장 MongoDB Sink Connector 설정 예시

아래는 구조를 이해하기 위한 예시입니다. 실제 클래스명과 세부 옵션은 사용하는 Connector 버전에 맞춰 검증해야 합니다.

{
  "name": "mongodb-sink-connector",
  "config": {
    "connector.class": "com.mongodb.kafka.connect.MongoSinkConnector",
    "tasks.max": "3",

    "topics": "business-data",

    "connection.uri": "${file:/opt/connect-secrets/mongodb.properties:connection.uri}",
    "database": "business",
    "collection": "customer",

    "key.converter": "org.apache.kafka.connect.json.JsonConverter",
    "key.converter.schemas.enable": "false",

    "value.converter": "org.apache.kafka.connect.json.JsonConverter",
    "value.converter.schemas.enable": "false",

    "document.id.strategy": "com.mongodb.kafka.connect.sink.processor.id.strategy.ProvidedInKeyStrategy",
    "writemodel.strategy": "com.mongodb.kafka.connect.sink.writemodel.strategy.ReplaceOneDefaultStrategy",

    "errors.tolerance": "all",
    "errors.deadletterqueue.topic.name": "business-data-dlq",
    "errors.deadletterqueue.context.headers.enable": "true",

    "max.num.retries": "10",
    "retries.defer.timeout": "5000"
  }
}

주요 항목 설명

설정의미
topics 읽을 Kafka Topic
connection.uri MongoDB Replica Set 접속 주소
database 저장할 MongoDB Database
collection 저장할 Collection
tasks.max 동시에 실행할 최대 작업 수
document.id.strategy MongoDB _id 생성 기준
writemodel.strategy Insert·Replace·Update 방식
errors.tolerance 오류 발생 시 처리 정책
deadletterqueue 실패 메시지를 보관할 Kafka Topic

tasks.max를 크게 설정하더라도 Kafka Partition 수보다 무조건 많이 병렬 처리되는 것은 아닙니다.

예:

Topic Partition = 3개
tasks.max        = 10개
실질적인 병렬도  = 최대 3개 수준

따라서 Partition 수, Connector Task 수, MongoDB 처리 성능을 함께 산정해야 합니다.


9. 실패 데이터는 DLQ로 분리해야 합니다

다음과 같은 메시지가 들어올 수 있습니다.

  • JSON 형식 오류
  • 필수 필드 누락
  • 데이터 타입 오류
  • 너무 큰 메시지
  • MongoDB Unique Index 충돌
  • 변환 규칙 오류
  • 날짜 형식 오류

이때 잘못된 메시지 한 건 때문에 전체 Connector가 중지되지 않도록 **DLQ(Dead Letter Queue)**를 구성할 수 있습니다.

```mermaid
flowchart LR
    A["Kafka Topic"] --> B["MongoDB Sink Connector"]
    B -->|정상| C["MongoDB"]
    B -->|실패| D["DLQ Topic"]
    D --> E["분석·수정·재처리"]
```

MongoDB Connector는 오류 메시지를 별도의 Kafka DLQ Topic으로 보낼 수 있으며, 오류 원인을 Header에 포함하도록 설정할 수 있습니다. MongoDB Connector 오류 처리

다만 다음 설정만 해놓고 끝내면 안 됩니다.

errors.tolerance=all

이 설정만 사용하면 오류를 지나치면서 실제 데이터 누락을 발견하지 못할 수 있습니다.

따라서 다음 구성이 필요합니다.

오류 허용
+ DLQ 저장
+ DLQ 건수 모니터링
+ 오류 알림
+ 재처리 절차

10. Schema 관리가 필요합니다

Producer마다 필드 이름과 데이터 타입을 제각각 보내면 MongoDB 데이터 구조가 깨질 수 있습니다.

예를 들어 같은 age 필드가 다음처럼 들어올 수 있습니다.

{"age": 30}
{"age": "30"}
{"age": "삼십"}

MongoDB는 유연한 문서 구조를 지원하지만, 아무 데이터나 허용해도 된다는 뜻은 아닙니다.

권장 방법

초기 또는 소규모 구성

JSON
+ 필수 필드 정의
+ Producer 검증
+ MongoDB Schema Validation

규모가 큰 운영 환경

Avro 또는 Protobuf
+ Schema Registry
+ Schema 버전 관리
+ 호환성 검사

설계 단계에서 최소한 다음 항목은 정해야 합니다.

  • Kafka Key 구조
  • Kafka Value 구조
  • 필수 필드
  • 필드 타입
  • 날짜와 시간 형식
  • Schema 버전
  • 필드 추가·삭제 정책
  • 잘못된 Schema 처리 방법

11. 초기 데이터와 실시간 데이터의 경계가 필요합니다

Kafka 스트리밍을 시작하기 전에 MongoDB에 기존 데이터를 먼저 적재해야 할 수 있습니다.

예를 들어 기존 시스템에 고객 데이터 1,000만 건이 있고, 오늘부터 변경 데이터만 Kafka로 전달한다고 가정하겠습니다.

다음 순서가 필요합니다.

1. 기준 시점 T0 결정
2. 기존 데이터 전체 추출
3. MongoDB에 초기 적재
4. T0 이후 Kafka 메시지 적용
5. 원본과 MongoDB 데이터 비교
6. 실시간 전환

주의할 점은 초기 적재 중 발생한 변경 데이터가 빠지거나 중복되지 않아야 한다는 것입니다.

설계서에는 다음 항목을 포함하는 것이 좋습니다.

  • 초기 적재 기준 시각
  • 초기 적재 대상 범위
  • Kafka Offset 시작 위치
  • 초기 적재와 실시간 메시지 중복 처리 방법
  • 원본·Kafka·MongoDB 건수 검증
  • 전환 실패 시 Rollback 방법

12. 원본 DB와 Kafka에 동시에 저장하면 Dual Write 문제가 발생합니다

애플리케이션이 다음 작업을 각각 수행한다고 가정하겠습니다.

1. 원본 DB 저장
2. Kafka 메시지 전송

다음 장애가 발생할 수 있습니다.

원본 DB 저장 성공
Kafka 전송 실패

그러면 원본 DB에는 데이터가 있지만 MongoDB에는 전달되지 않습니다.

반대 상황도 가능할 수 있습니다.

Kafka 전송 성공
원본 DB 트랜잭션 실패

이를 Dual Write 불일치 문제라고 합니다.

권장 대응

방법 1: Transactional Outbox

업무 데이터 저장
+ Outbox 테이블 저장
→ 같은 DB 트랜잭션으로 처리

Outbox 데이터
→ Kafka로 전송
→ MongoDB Sink

방법 2: CDC

원본 DB 변경 로그
→ CDC Connector
→ Kafka
→ MongoDB Sink Connector
→ MongoDB

원본 DB와 MongoDB 간 데이터 정합성이 중요하다면, 애플리케이션에서 단순히 “DB 저장 후 Kafka 전송”하는 방식보다 Outbox 또는 CDC 방식을 우선 검토하는 것이 좋습니다.


13. MongoDB 쓰기 안정성 설정

MongoDB에 데이터를 쓸 때는 Write Concern을 함께 검토해야 합니다.

예를 들어 w=majority를 사용하면 Replica Set 과반수에 쓰기가 반영된 것을 기준으로 성공을 판단할 수 있습니다.

개념적으로는 다음과 같습니다.

Connector → Primary 저장
          → Secondary 복제 확인
          → 성공 응답

장점:

  • Primary 장애 시 데이터 유실 가능성 감소
  • 더 높은 데이터 안정성

단점:

  • Secondary 복제 확인 시간이 추가됨
  • 쓰기 지연시간이 늘어날 수 있음

중요 업무 데이터라면 보통 안정성을 우선해 다음과 같은 방향을 검토합니다.

replicaSet=rs0
retryWrites=true
w=majority

다만 실제 적용은 처리량, 지연시간, MongoDB 버전 및 드라이버 설정을 기준으로 부하 테스트 후 결정해야 합니다.


14. 성능 설계 시 확인할 항목

14.1 Kafka 측

  • 초당 메시지 발생량
  • 평균 메시지 크기
  • 최대 메시지 크기
  • Topic Partition 수
  • Kafka 복제 계수
  • Retention 기간
  • Producer acks 정책
  • 압축 방식
  • Consumer Lag

14.2 Kafka Connect 측

  • Connect Worker 수
  • Distributed Mode 구성
  • Connector Task 수
  • CPU와 Memory
  • Batch 크기
  • 재시도 횟수
  • DLQ 처리량
  • Connector 재기동 시간

운영 환경에서는 Kafka Connect를 한 개로만 실행하기보다 여러 Worker를 사용하는 Distributed Mode 구성이 적합합니다.

예:

Kafka Connect Worker 1
Kafka Connect Worker 2
Kafka Connect Worker 3

한 Worker에 장애가 발생하면 다른 Worker가 Connector Task를 재배치해 이어받을 수 있습니다.

14.3 MongoDB 측

  • Primary의 CPU와 Memory
  • Disk IOPS와 지연시간
  • Collection별 쓰기 처리량
  • Index 개수
  • Replica 지연
  • WiredTiger Cache
  • Connection 수
  • Document 평균 크기
  • Unique Index 충돌
  • PVC 성능과 용량

MongoDB Connector는 여러 쓰기 요청을 묶는 Bulk Write를 사용합니다. Ordered 방식은 순서를 보장하지만 중간 오류 이후의 작업이 처리되지 않을 수 있고, Unordered 방식은 성능을 높일 수 있지만 순서가 중요한 데이터에는 주의해야 합니다. MongoDB Bulk Write 전략


15. 반드시 모니터링해야 할 항목

구간주요 지표의미
Kafka Consumer Lag MongoDB 적재가 얼마나 밀렸는지
Kafka Under Replicated Partitions Kafka 복제 이상
Kafka Connect Connector Status RUNNING·FAILED 상태
Kafka Connect Task Status 개별 Task 장애
Kafka Connect 처리 성공·실패 건수 적재 결과
DLQ 메시지 수 적재 실패 데이터
MongoDB Replication Lag Secondary 복제 지연
MongoDB Primary 상태 쓰기 가능 여부
MongoDB Insert·Update Latency 적재 지연
MongoDB Connection 수 Connector 연결 상태
MongoDB Disk Usage PVC 용량
MongoDB Disk Latency 스토리지 병목

MongoDB Kafka Connector는 Kafka Connect 메트릭 체계를 이용해 JMX 지표를 제공하므로 Prometheus JMX Exporter 등과 연계할 수 있습니다. MongoDB Connector 모니터링

권장 알림 조건 예시

Connector FAILED 발생
Consumer Lag 지속 증가
DLQ 메시지 발생
MongoDB Primary 없음
MongoDB Replication Lag 임계치 초과
PVC 사용량 80% 이상
MongoDB 쓰기 지연 증가
Kafka Topic Retention 소진 임박

16. 보안 구성 시 주의점

네트워크

  • Kafka Broker 외부 공개 금지
  • Kafka Connect에서 MongoDB로 필요한 포트만 허용
  • Kubernetes NetworkPolicy 적용
  • 관리용 REST API 접근 제한

인증 및 암호화

  • Kafka SASL·TLS 적용
  • MongoDB 사용자 인증 적용
  • MongoDB TLS 적용
  • 접속 비밀번호를 Connector JSON에 평문 저장하지 않기
  • Kubernetes Secret 또는 외부 Secret Manager 사용

권한

MongoDB Sink 전용 계정을 생성하고 필요한 Database와 Collection에만 최소 권한을 부여합니다.

관리자 계정 사용 금지
전체 DB 권한 부여 금지
Sink 전용 계정 사용
대상 Collection 쓰기 권한만 부여

개인정보

DLQ와 로그에도 원본 메시지가 들어갈 수 있습니다.

따라서 다음 항목을 확인해야 합니다.

  • 개인정보 마스킹
  • DLQ 접근 권한
  • DLQ Retention 기간
  • 로그 본문 출력 여부
  • 암호화
  • 감사 로그

17. Kafka가 백업을 대신하는 것은 아닙니다

Kafka Topic에 데이터가 일정 기간 남아 있어도 이것을 MongoDB 백업으로 보면 안 됩니다.

Kafka Retention
= 메시지 재처리를 위한 보관

MongoDB Backup
= 특정 시점의 DB 복원

별도로 다음이 필요합니다.

  • MongoDB 정기 백업
  • 백업 보관 기간
  • Point-in-Time Recovery 검토
  • 복구 테스트
  • 다른 저장소 또는 DR센터 보관
  • RPO·RTO 정의

Kafka 메시지 재처리만으로 MongoDB 전체를 복원하려면 모든 Topic 데이터가 충분한 기간 동안 보존되어야 하고, Schema·삭제 이벤트·처리 순서까지 완전하게 유지되어야 하므로 운영상 위험합니다.


18. 설계 시 권장하는 최종안

권장 구조

① 업무 애플리케이션 또는 CDC
          ↓
② Kafka Topic
   - 업무 고유번호를 Message Key로 사용
   - 최소 3개 이상의 Partition 검토
   - 복제 계수 3 검토
          ↓
③ Kafka Connect Cluster
   - Distributed Mode
   - Worker 2~3개 이상
          ↓
④ MongoDB Sink Connector
   - Upsert 적용
   - 고정된 _id 또는 Business Key 사용
   - 재시도 설정
   - DLQ 설정
          ↓
⑤ MongoDB Replica Set
   - Primary 1
   - Secondary 2
   - 각 Pod별 독립 PVC
          ↓
⑥ 모니터링 및 백업
   - Consumer Lag
   - DLQ
   - Replica Lag
   - PVC 용량
   - MongoDB 백업

19. 설계·분석 문서에 넣을 수 있는 문구

구성방안

Kafka에 수집된 데이터를 MongoDB로 실시간 적재하기 위해 Kafka Connect 기반의 MongoDB Sink Connector를 구성한다. Sink Connector는 지정된 Kafka Topic을 구독하여 메시지를 MongoDB 문서로 변환하고, MongoDB Replica Set의 현재 Primary 노드에 저장한다. Primary에 저장된 데이터는 MongoDB 자체 Replication 기능을 통해 Secondary 노드로 복제된다.

스토리지 구성

MongoDB Replica Set을 구성하는 각 Pod는 동일한 Volume을 공유하지 않고, Pod별 독립 PVC를 사용한다. 데이터의 일관성과 복제는 공유 스토리지가 아니라 MongoDB Replica Set의 내부 복제 기능으로 보장한다.

중복 방지

Kafka Connect의 장애 복구 및 재처리 과정에서 동일 메시지가 중복 처리될 가능성을 고려하여, Kafka Message Key와 MongoDB _id를 업무 고유 Key로 통일하고 Upsert 기반의 멱등성 저장 방식을 적용한다.

순서 보장

동일 업무 객체의 이벤트 순서를 최대한 유지하기 위해 customerId, orderId, deviceId 등 업무 고유 식별자를 Kafka Message Key로 사용하여 동일 Partition으로 라우팅한다. Kafka의 메시지 순서 보장은 Topic 전체가 아닌 Partition 단위임을 고려한다.

오류 처리

변환 실패, Schema 오류 및 MongoDB 쓰기 실패 메시지는 DLQ Topic으로 분리하여 보관한다. DLQ 발생 건수, Connector Task 상태 및 Consumer Lag을 모니터링하고, 실패 메시지에 대한 원인 분석과 재처리 절차를 마련한다.

고가용성

Kafka Connect는 Distributed Mode로 다중 Worker를 구성하고, MongoDB는 최소 3개 인스턴스의 Replica Set으로 구성한다. MongoDB 연결 문자열에는 Replica Set 정보를 포함하여 Primary 변경 시 Connector가 새로운 Primary를 자동 탐색할 수 있도록 구성한다.


20. 최종 핵심 정리

Kafka가 mongodb-0, 1, 2에 각각 쓰는 것이 아님
                    ↓
Kafka Sink Connector가 MongoDB Replica Set에 접속
                    ↓
현재 Primary에 데이터 저장
                    ↓
MongoDB가 Secondary로 복제
                    ↓
각 MongoDB Pod는 자기 PVC에 데이터 보관

가장 중요한 설계 원칙은 다음 7가지입니다.

  1. MongoDB Pod별로 독립 PVC를 사용합니다.
  2. Sink Connector는 Replica Set 주소로 연결합니다.
  3. 업무 고유번호를 Kafka Key와 MongoDB _id로 사용합니다.
  4. 현재 상태 데이터는 Insert보다 Upsert 방식을 권장합니다.
  5. 오류 데이터는 DLQ로 분리하고 반드시 모니터링합니다.
  6. 초기 적재와 실시간 적재의 기준 시점을 정의합니다.
  7. Kafka Retention과 MongoDB 백업을 별도로 운영합니다.
 
 
반응형

댓글