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 토픽을 재생해요 |