하이브리드 실시간 + 오프라인

하이브리드 실시간 + 오프라인 (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)
                                └──────────────────┘

시간 경계가 동작하는 방식:

  1. Pinot이 시간 경계를 유지해요 — 오프라인과 실시간 데이터를 구분하는 타임스탬프예요.
  2. 경계 이전 시간 범위를 다루는 쿼리는 오프라인 세그먼트로 라우팅돼요.
  3. 경계 이후 시간 범위는 실시간 세그먼트로 라우팅돼요.
  4. 최근 시간 범위를 다루는 새 오프라인 세그먼트가 푸시될 때마다 경계가 전진하고, 겹치는 실시간 세그먼트가 자동으로 제거돼요.

스키마 (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)

  1. 컨트롤러가 모든 오프라인 세그먼트의 최대 종료 시간으로 시간 경계를 계산해요.
  2. 브로커가 쿼리 시점에 이 경계를 사용해요. 경계 이전 시간 범위의 세그먼트는 오프라인에서, 이후 세그먼트는 실시간에서 가져와요.
  3. 새 오프라인 세그먼트가 이전에 실시간이 서빙하던 시간을 다루게 되면, 그 시간 범위의 실시간 세그먼트는 중복이 되어 보존 매니저가 제거해요.

현재 시간 경계를 확인할 수 있어요:

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이 일치하고, 배치 작업이 같은 컬럼을 읽는지 확인하세요

더 알아보기 (Learn more)