하이브리드 실시간 + 오프라인
하이브리드 실시간 + 오프라인 (Hybrid Real-Time + Offline)
하나의 Pinot 테이블에서 저지연 스트리밍과 고품질 배치 백필(backfill)을 결합하는 전체 가이드예요. 이 플레이북은 하이브리드 테이블 패턴을 다뤄요. 즉, 하나의 논리적 테이블이 저지연 스트리밍 데이터용 실시간(real-time) 테이블과, 일괄 수정·중복 제거·강화된 과거 데이터용 오프라인(offline) 테이블 양쪽으로 뒷받침돼요. 쿼리는 두 테이블을 매끄럽게 아우르고, Pinot의 시간 경계(time boundary) 메커니즘이 각 시간 범위가 최고 품질의 소스로 서빙되도록 보장해요.
출처: 문서
본문
이 패턴을 언제 쓸까 (When to use this pattern)
다음과 같은 경우 이 플레이북을 사용해요:
- 같은 쿼리에서 실시간 신선함(초 단위 지연) 그리고 배치 품질의 과거 데이터(중복 제거, 강화, 재처리)가 모두 필요해요.
- 야간 또는 시간별 배치 작업이 스트림으로 도착하는 데이터보다 더 높은 품질의 데이터를 만들 수 있어요(예: 차원 테이블 조인, 지각 데이터 수정, ML 강화 후).
- 많은 작은 실시간 세그먼트를 더 적고 최적화된 오프네일 세그먼트로 교체해 스토리지 비용을 줄이고 싶어요.
- 워크로드가 Lambda 아키텍처와 비슷하지만, 두 시스템을 각각 조회하는 대신 단일 쿼리 인터페이스를 원해요.
실시간 데이터만 필요하고 백필이 전혀 없다면 실시간 제품 분석 플레이북을 사용하세요. 데이터가 제자리에서 변형된다면 CDC / Upsert 파이프라인을 보세요.
아키텍처 스케치 (Architecture sketch)
┌──────────────────┐
Kafka topic ──────────────────▶ │ REALTIME table │ (fresh, last few hours)
└────────┬─────────┘
│
Pinot query │ time boundary
spans both │
│
┌────────┴─────────┐
Spark / Flink batch job ──────▶ │ OFFLINE table │ (historical, optimized)
└──────────────────┘
시간 경계가 동작하는 방식:
- Pinot이 시간 경계를 유지해요 — 오프라인과 실시간 데이터를 구분하는 타임스탬프예요.
- 경계 이전 시간 범위를 다루는 쿼리는 오프라인 세그먼트로 라우팅돼요.
- 경계 이후 시간 범위는 실시간 세그먼트로 라우팅돼요.
- 최근 시간 범위를 다루는 새 오프라인 세그먼트가 푸시될 때마다 경계가 전진하고, 겹치는 실시간 세그먼트가 자동으로 제거돼요.
스키마 (Schema)
스키마는 오프라인과 실시간 테이블이 공유해요. 둘 다 정확히 같은 schemaName을 사용해야 해요.
{
"schemaName": "web_analytics",
"dimensionFieldSpecs": [
{ "name": "pageUrl", "dataType": "STRING" },
{ "name": "userId", "dataType": "STRING" },
{ "name": "country", "dataType": "STRING" },
{ "name": "browser", "dataType": "STRING" },
{ "name": "campaign", "dataType": "STRING" },
{ "name": "enrichedSegment", "dataType": "STRING" }
],
"metricFieldSpecs": [
{ "name": "durationMs", "dataType": "LONG" },
{ "name": "revenue", "dataType": "DOUBLE" }
],
"dateTimeFieldSpecs": [
{
"name": "eventTimestamp",
"dataType": "TIMESTAMP",
"format": "1:MILLISECONDS:EPOCH",
"granularity": "1:MILLISECONDS"
}
]
}
참고:
enrichedSegment컬럼은 배치 작업이 채우지만 실시간 데이터에서는 null일 수 있어요. Pinot은 null 값을 자연스럽게 처리하므로 문제없어요. Null Value Support 참고.
실시간 테이블 구성 (Real-time table configuration)
{
"tableName": "web_analytics",
"tableType": "REALTIME",
"segmentsConfig": {
"timeColumnName": "eventTimestamp",
"retentionTimeUnit": "DAYS",
"retentionTimeValue": "3",
"replication": "2",
"segmentPushType": "APPEND"
},
"tableIndexConfig": {
"loadMode": "MMAP",
"invertedIndexColumns": ["country", "browser", "campaign"],
"rangeIndexColumns": ["eventTimestamp"],
"noDictionaryColumns": ["userId", "pageUrl"],
"streamConfigs": {
"streamType": "kafka",
"stream.kafka.topic.name": "web-analytics-events",
"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": "500000",
"realtime.segment.flush.threshold.time": "6h"
}
},
"tenants": {
"broker": "DefaultTenant",
"server": "DefaultTenant"
},
"metadata": {}
}
retentionTimeValue는 짧은 창(예: 3일)으로 설정하세요. 이것은 단지 폴백(fallback)일 뿐이에요. 정상 운영에서는 시간 경계 메커니즘이 오프라인 세그먼트로 교체되면서 실시간 세그먼트를 제거해요.
오프라인 테이블 구성 (Offline table configuration)
{
"tableName": "web_analytics",
"tableType": "OFFLINE",
"segmentsConfig": {
"timeColumnName": "eventTimestamp",
"retentionTimeUnit": "DAYS",
"retentionTimeValue": "365",
"replication": "2",
"segmentPushType": "APPEND"
},
"tableIndexConfig": {
"loadMode": "MMAP",
"invertedIndexColumns": ["country", "browser", "campaign", "enrichedSegment"],
"rangeIndexColumns": ["eventTimestamp"],
"sortedColumn": ["country"],
"noDictionaryColumns": ["userId", "pageUrl"],
"starTreeIndexConfigs": [
{
"dimensionsSplitOrder": ["country", "browser", "campaign"],
"functionColumnPairs": ["COUNT__*", "SUM__revenue", "SUM__durationMs"],
"maxLeafRecords": 10000
}
]
},
"tenants": {
"broker": "DefaultTenant",
"server": "DefaultTenant"
},
"metadata": {}
}
오프라인 테이블이 더 공격적으로 인덱싱될 수 있는 이유
- 스타트리 인덱스: 배치 세그먼트는 한 번만 쓰이므로 스타트리 빌드 비용이 한 번만 발생해요. 실시간 consuming 세그먼트는 flush마다 스타트리를 다시 빌드해서 더 비싸요.
- 정렬 컬럼: 오프라인 세그먼트는 배치 작업 중에 전역적으로 정렬되어 완벽한 정렬 순서를 얻을 수 있어요. 실시간 세그먼트는 각 flush 내에서만 정렬돼요.
- 더 긴 보존: 과거 데이터는 오프라인 세그먼트에만 있으므로 보존을 수 개월 또는 수 년으로 설정하세요.
배치 적재 작업 (Batch ingestion job)
Spark 또는 독립 실행형 적재 작업을 사용해 매일 오프라인 세그먼트를 푸시하세요. Spark 작업 사양 예시:
executionFrameworkSpec:
name: spark
segmentGenerationJobRunnerClassName: org.apache.pinot.plugin.ingestion.batch.spark3.SparkSegmentGenerationJobRunner
segmentTarPushJobRunnerClassName: org.apache.pinot.plugin.ingestion.batch.spark3.SparkSegmentTarPushJobRunner
jobType: SegmentCreationAndTarPush
inputDirURI: s3://data-lake/web_analytics/dt=2026-03-31/
outputDirURI: s3://pinot-segments/web_analytics/
overwriteOutput: true
pinotFSSpecs:
- scheme: s3
className: org.apache.pinot.plugin.filesystem.S3PinotFS
recordReaderSpec:
dataFormat: parquet
className: org.apache.pinot.plugin.inputformat.parquet.ParquetRecordReader
tableSpec:
tableName: web_analytics
schemaURI: http://pinot-controller:9000/schemas/web_analytics
tableConfigURI: http://pinot-controller:9000/tables/web_analytics
segmentNameGeneratorSpec:
type: normalizedDate
configs:
segment.name.prefix: web_analytics
exclude.sequence.id: false
pushJobSpec:
pushAttempts: 3
pushParallelism: 2
데이터 파이프라인이 끝난 뒤 매일(또는 시간별로) 이 작업을 예약하세요. 각 푸시는 시간 경계를 자동으로 전진시켜요.
자세한 내용은 Batch Ingestion Guide와 Spark Connector를 참고하세요.
시간 경계가 전진하는 방식 (How the time boundary advances)
- 컨트롤러가 모든 오프라인 세그먼트의 최대 종료 시간으로 시간 경계를 계산해요.
- 브로커가 쿼리 시점에 이 경계를 사용해요. 경계 이전 시간 범위의 세그먼트는 오프라인에서, 이후 세그먼트는 실시간에서 가져와요.
- 새 오프라인 세그먼트가 이전에 실시간이 서빙하던 시간을 다루게 되면, 그 시간 범위의 실시간 세그먼트는 중복이 되어 보존 매니저가 제거해요.
현재 시간 경계를 확인할 수 있어요:
GET /tables/web_analytics/timeBoundary
전체 알고리즘은 Time Boundary를 참고하세요.
쿼리 패턴 (Query patterns)
쿼리는 단일 테이블 쿼리와 동일해요. 브로커가 라우팅을 투명하게 처리해요:
SELECT
DATETRUNC('day', eventTimestamp, 'MILLISECONDS') AS day,
country,
COUNT(*) AS page_views,
SUM(revenue) AS total_revenue
FROM web_analytics
WHERE eventTimestamp > ago('P30D')
GROUP BY day, country
ORDER BY day
LIMIT 10000
최근 데이터(지난 몇 시간)는 쿼리가 실시간 세그먼트를 조회하고, 오래된 데이터는 오프라인 세그먼트를 조회해요. 결과는 매끄럽게 병합돼요.
운영 체크리스트 (Operational checklist)
서비스 시작 전
- 오프라인과 실시간 테이블이 같은
schemaName과tableName을 쓰는지 확인하세요. - 실시간 테이블 보존을 배치 작업 간격의 최소 2배로 설정하세요(예: 일일 작업이면 3일) — 배치 지연을 견디기 위해서예요.
- 첫 배치 푸시 후
GET /tables/{tableName}/timeBoundary로 시간 경계가 전진하는지 검증하세요. - 오프라인 세그먼트에 시간 공백이 없는지 확인하세요. 날짜가 빠지면 시간 경계가 전진을 멈추고, 해당 공백의 실시간 데이터는 보존이 만료될 때 유실돼요.
모니터링
- 시간 경계 신선도: 경계가 배치 간격 하나보다 더 뒤처져 있으면 배치 파이프라인이 실패하고 있을 수 있어요.
- 세그먼트 중복: 오프라인 세그먼트가 다루는 시간 범위의 실시간 세그먼트가 보존에 의해 정리되는지 확인하세요.
- 오프라인 세그먼트 푸시 성공: 푸시 실패를 컨트롤러 API로 모니터링하세요. 푸시 실패 시 실시간 테이블의 오래된 데이터가 계속 서빙되지만, 이는 안전하되 품질이 낮은 상태예요.
흔한 함정 (Common pitfalls)
| 함정 | 해결책 |
|---|---|
| 시간 경계가 전진하지 않음 | 오프라인 세그먼트의 시간 범위가 연속적이고 예상 경계까지 채우는지 확인하세요 |
| 오프라인 푸시 후 실시간 세그먼트가 제거되지 않음 | 실시간 보존이 구성됐는지, 보존 매니저가 실행 중인지 확인하세요 |
| 경계 시간 범위에 대한 이중 집계 | 오프라인과 실시간 세그먼트가 같은 밀리초를 겹치면 발생할 수 있어요. Pinot의 시간 경계 로직이 처리하지만, 경계에서 COUNT 쿼리로 검증하세요 |
| 배치 작업이 잘못된 시간 컬럼의 세그먼트를 생성함 | 두 테이블 구성에서 timeColumnName이 일치하고, 배치 작업이 같은 컬럼을 읽는지 확인하세요 |