CDC / Upsert 파이프라인

CDC / Upsert 파이프라인 (CDC / Upsert Pipeline)

CDC와 upsert를 이용해 Apache Pinot을 트랜잭션 데이터베이스와 동기화 상태로 유지하는 전체 가이드예요. 이 플레이북은 트랜잭션 데이터베이스(PostgreSQL, MySQL, MongoDB 등)에서 Change Data Capture(CDC)로 행 단위 변경을 포착하고, 이를 Kafka로 스트리밍한 뒤 upsert를 활성화한 Pinot에 적재해서 Pinot이 항상 각 행의 최신 상태를 반영하도록 만드는 패턴을 다뤄요.

출처: 문서

본문

이 패턴을 언제 쓸까 (When to use this pattern)

다음과 같은 경우 이 플레이북을 사용해요:

  • Pinot이 OLTP 데이터베이스(주문, 사용자 프로필, 재고, 티켓)의 행 현재 상태를 미러링해야 해요.
  • 행이 생성 후 자주 갱신(updated) 되거나 소프트 삭제(soft-deleted) 되며, 쿼리는 최신 버전만 봐야 해요.
  • 기본 데이터베이스에 읽기 부하를 추가하지 않고 트랜잭션 데이터 위에 실시간 분석을 원해요.
  • CDC 도구(Debezium, Maxwell, AWS DMS)가 이미 Kafka로 이벤트를 생성하고 있어요.

데이터가 추가 전용(append-only)이고 갱신되지 않는다면, 더 단순한 실시간 제품 분석 패턴이 더 잘 맞아요.

아키텍처 스케치 (Architecture sketch)

OLTP DB ──▶ Debezium ──▶ Kafka topic ──▶ Pinot REALTIME table (upsert)
  (WAL)      (CDC)        (per-table)          │
                                      ┌────────┴────────┐
                                      │ Servers with     │
                                      │ primary key map  │
                                      └─────────────────┘

핵심 구성 요소는 다음과 같아요:

  • Debezium(또는 동급 도구)이 데이터베이스 WAL을 읽고 모든 INSERT, UPDATE, DELETE에 대해 Kafka 이벤트를 생성해요.
  • Kafka 토픽이 기본 키로 파티셔닝되어 같은 행의 모든 변경이 같은 파티션에 도착해요.
  • upsert를 켠 Pinot 실시간 테이블이 기본 키에서 최신 세그먼트/doc-id로 가는 인메모리 맵을 유지해서 쿼리가 오래된 버전을 건너뛰도록 해요.

경고: Upsert는 같은 기본 키의 모든 이벤트가 같은 Kafka 파티션에 도착해야 해요. Debezium 커넥터나 Kafka 프로듀서를 기본 키 컬럼 기준으로 파티셔닝하도록 구성하세요.

스키마 (Schema)

원본 테이블의 기본 키를 Pinot의 기본 키로 사용해요. 비교 컬럼(주로 이벤트 타임스탬프나 데이터베이스 트랜잭션 시퀀스 번호)도 포함해서 Pinot이 어떤 버전이 더 새로운지 판단하게 해요.

{
  "schemaName": "orders",
  "primaryKeyColumns": ["orderId"],
  "dimensionFieldSpecs": [
    { "name": "orderId",    "dataType": "STRING" },
    { "name": "customerId", "dataType": "STRING" },
    { "name": "status",     "dataType": "STRING" },
    { "name": "region",     "dataType": "STRING" },
    { "name": "product",    "dataType": "STRING" }
  ],
  "metricFieldSpecs": [
    { "name": "amount",     "dataType": "DOUBLE" },
    { "name": "quantity",   "dataType": "INT" }
  ],
  "dateTimeFieldSpecs": [
    {
      "name": "updatedAt",
      "dataType": "TIMESTAMP",
      "format": "1:MILLISECONDS:EPOCH",
      "granularity": "1:MILLISECONDS"
    }
  ]
}

참고: primaryKeyColumns는 upsert에 필수예요. Pinot이 이 값을 이용해 기본 키 맵에서 기존 행을 조회해요. Schema and Table Shape를 참고하세요.

테이블 구성 (Table configuration)

{
  "tableName": "orders",
  "tableType": "REALTIME",
  "segmentsConfig": {
    "timeColumnName": "updatedAt",
    "retentionTimeUnit": "DAYS",
    "retentionTimeValue": "90",
    "replication": "1"
  },
  "tableIndexConfig": {
    "loadMode": "MMAP",
    "invertedIndexColumns": ["status", "region", "customerId"],
    "rangeIndexColumns": ["updatedAt"],
    "sortedColumn": [],
    "noDictionaryColumns": ["orderId", "customerId"],
    "streamConfigs": {
      "streamType": "kafka",
      "stream.kafka.topic.name": "dbserver1.public.orders",
      "stream.kafka.broker.list": "kafka:9092",
      "stream.kafka.consumer.factory.class.name": "org.apache.pinot.plugin.stream.kafka30.KafkaConsumerFactory",
      "stream.kafka.decoder.class.name": "org.apache.pinot.plugin.stream.kafka.KafkaJSONMessageDecoder",
      "realtime.segment.flush.threshold.rows": "100000",
      "realtime.segment.flush.threshold.time": "4h"
    }
  },
  "upsertConfig": {
    "mode": "FULL",
    "comparisonColumns": ["updatedAt"],
    "hashFunction": "MURMUR3",
    "enableSnapshot": true,
    "enablePreload": true,
    "metadataTTL": 0,
    "deletedKeysTTL": 0
  },
  "routing": {
    "instanceSelectorType": "strictReplicaGroup"
  },
  "tenants": {
    "broker": "DefaultTenant",
    "server": "DefaultTenant"
  },
  "metadata": {}
}

구성 요점 (Configuration highlights)

설정 이유
upsertConfig.mode: FULL 모든 갱신이 전체 행을 교체해요. CDC 이벤트가 변경된 컬럼만 담고 있다면 PARTIAL을 사용하세요. Upsert Modes 참고
comparisonColumns: ["updatedAt"] Pinot이 이 컬럼으로 순서가 어긋난 이벤트를 해결해요. updatedAt이 더 높은 행이 승리해요
enableSnapshot: true 기본 키 맵을 디스크에 영속화해서 서버 재시작 시 전체 Kafka 토픽을 재생하지 않아도 되게 해요
enablePreload: true 시작 시 스냅샷을 메모리에 로드해서 더 빠르게 복구해요
hashFunction: MURMUR3 기본 키 맵에 대해 기본 해시보다 메모리 효율적이에요
replication: 1 Upsert 테이블은 현재 복제 팩터 1이 필요해요. 일관성을 보장하려면 strictReplicaGroup 라우팅을 사용하세요
routing.instanceSelectorType 파티션의 모든 쿼리가 같은 서버로 라우팅되게 해서 올바른 upsert 의미론을 보장해요

삭제 처리 (Handling deletes)

CDC 스트림에 삭제된 행의 툼스톤(tombstone) 이벤트가 포함된다면 deleteRecordColumn을 구성하세요:

"upsertConfig": {
  "mode": "FULL",
  "comparisonColumns": ["updatedAt"],
  "deleteRecordColumn": "isDeleted",
  "deletedKeysTTL": 86400
}

스키마에 isDeleted를 불리언 차원으로 추가하세요. Debezium SMT(Single Message Transform)에서 삭제 이벤트에 대해 true로 설정하세요. 그러면 Pinot이 행을 삭제된 것으로 표시하고 쿼리 결과에서 더 이상 반환하지 않아요.

부분 upsert (Partial upsert)

CDC 이벤트가 변경된 컬럼만 담고 있다면(예: 주문에서 status만 바뀐 경우) 부분 upsert를 사용해 갱신을 기존 행에 병합해요:

"upsertConfig": {
  "mode": "PARTIAL",
  "partialUpsertStrategies": {
    "status":     "OVERWRITE",
    "amount":     "OVERWRITE",
    "quantity":   "OVERWRITE",
    "region":     "OVERWRITE"
  },
  "comparisonColumns": ["updatedAt"],
  "defaultPartialUpsertStrategy": "IGNORE"
}

기본값을 IGNORE로 두면, 들어오는 이벤트에 없는 컬럼은 이전 값을 유지해요. OVERWRITE로 나열된 컬럼은 존재할 때 교체돼요.

Debezium 구성 팁 (Debezium configuration tips)

PostgreSQL용 최소 Debezium 커넥터 구성:

{
  "name": "orders-connector",
  "config": {
    "connector.class": "io.debezium.connector.postgresql.PostgresConnector",
    "database.hostname": "postgres",
    "database.port": "5432",
    "database.user": "debezium",
    "database.password": "${DEBEZIUM_PASSWORD}",
    "database.dbname": "app",
    "topic.prefix": "dbserver1",
    "table.include.list": "public.orders",
    "key.converter": "org.apache.kafka.connect.json.JsonConverter",
    "value.converter": "org.apache.kafka.connect.json.JsonConverter",
    "transforms": "unwrap",
    "transforms.unwrap.type": "io.debezium.transforms.ExtractNewRecordState",
    "transforms.unwrap.drop.tombstones": "false",
    "transforms.unwrap.delete.handling.mode": "rewrite"
  }
}

ExtractNewRecordState 변환은 Debezium 봉투(envelope)를 평탄화해서 Pinot이 행의 컬럼만 담은 간단한 JSON 객체를 받게 해요. delete.handling.mode: rewrite는 __deleted 필드를 추가하며, 이를 ingestion transformation으로 isDeleted 컬럼에 매핑할 수 있어요.

쿼리 패턴 (Query patterns)

현재 상태 집계

SELECT status, COUNT(*) AS order_count, SUM(amount) AS total_amount
FROM orders
WHERE region = 'US'
GROUP BY status

upsert가 활성화되어 있으므로 이 쿼리는 자동으로 각 주문의 최신 버전만 봐요.

포인트 조회 (Point lookup)

SELECT *
FROM orders
WHERE orderId = 'ORD-78901'
LIMIT 1

변경 가능한 데이터에 대한 시간 범위 분석

SELECT
  DATETRUNC('hour', updatedAt, 'MILLISECONDS') AS hour_bucket,
  COUNT(*) AS updates
FROM orders
WHERE updatedAt > ago('PT24H')
GROUP BY hour_bucket
ORDER BY hour_bucket

세그먼트 압축 (Segment compaction)

시간이 지나면 upsert 테이블은 완료된 세그먼트 안에 오래된 행 버전을 쌓아요. 이런 무효 레코드는 스토리지를 낭비하고 전체 테이블 스캔을 느리게 해요. Minion을 통해 Upsert Compaction Task를 예약해 세그먼트를 다시 쓰고 오래된 행을 물리적으로 제거하세요:

{
  "task": {
    "taskTypeConfigsMap": {
      "UpsertCompactionTask": {
        "schedule": "0 0 2 * * ?",
        "invalidRecordsThresholdPercent": "30",
        "invalidRecordsThresholdCount": "100000"
      }
    }
  }
}

튜닝 지침은 Segment Compaction on Upserts를 참고하세요.

운영 체크리스트 (Operational checklist)

서비스 시작 전

  • Kafka 토픽이 기본 키로 파티셔닝되어 있는지 확인하세요. 파티션 간에 중복 키로 테스트를 돌려보세요. 잘못 구성되면 쿼리가 잘못된 결과를 반환해요.
  • 서버 재시작 시 전체 토픽 재생을 피하려면 enableSnapshot과 enablePreload를 켜세요.
  • replication: 1과 strictReplicaGroup 라우팅을 설정하세요. 현재 릴리스에서는 다중 복제 upsert가 지원되지 않아요.
  • CDC를 시작하기 전에 전체 테이블 스냅샷 초기 로드를 실행해서 과거 데이터가 빠지지 않게 하세요.
  • 세그먼트 flush 임계값을 기본 키 맵이 서버 메모리에 들어갈 만큼 작게 설정하세요. 힙 사용량을 모니터링하세요.

모니터링

  • 기본 키 맵 크기: pinot.server.upsertPrimaryKeysCount를 모니터링하세요. 사용 가능한 메모리를 초과하면 TTL 기반 축출이나 더 큰 서버를 고려하세요.
  • 무효 레코드 비율: 세그먼트당 무효화(오래된) 레코드 비율을 추적하세요. 30%를 넘으면 압축을 예약하세요.
  • Kafka 컨슈머 래그: 실시간 테이블과 마찬가지로, 래그가 커지면 Pinot 쿼리가 오래된 데이터를 보여줘요.
  • Debezium 커넥터 상태: Kafka Connect REST API로 모니터링하세요. 커넥터가 멈추면 Pinot으로 갱신이 흐르지 않아요.

흔한 함정 (Common pitfalls)

함정 해결책
순서가 어긋난 이벤트로 낡은 값이 새 값을 덮어씀 단조 증가 컬럼(데이터베이스 시퀀스, 이벤트 타임스탬프)과 함께 comparisonColumns를 사용하세요
큰 기본 키 맵으로 인한 메모리 압박 metadataTTL을 켜 시간 창 안에 갱신되지 않은 행의 키를 축출하거나, hashFunction: MURMUR3으로 키당 오버헤드를 줄이세요
쿼리가 삭제된 행을 반환함 deleteRecordColumn이 구성됐는지, CDC 변환이 삭제 이벤트에 대해 이를 올바르게 설정하는지 확인하세요
서버 재시작에 수 분이 걸림 enableSnapshot과 enablePreload를 켜세요. 없으면 Pinot이 마지막 커밋 오프셋부터 Kafka 토픽을 재생해요

더 알아보기 (Learn more)