실시간 제품 분석

실시간 제품 분석 (Real-Time Product Analytics)

Apache Pinot로 Kafka 이벤트 스트림 위에 밀리초 단위 대시보드를 구축하는 전체 가이드예요. 이 플레이북은 스트리밍 데이터 위에 사람이 직접 보거나 내부용으로 쓰는 분석 대시보드를 밀리초 단위 쿼리 지연으로 구동하는 전체 구성을 안내해요. 페이지뷰 카운터, 거래 대시보드, 클릭률(CTR) 모니터, 그리고 이벤트가 끊임없이 도착하고 대시보드가 실시간으로 갱신돼야 하는 어떤 워크로드에도 적용할 수 있어요.

출처: 문서

본문

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

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

  • 이벤트가 Kafka(또는 Kinesis/Pulsar)로 생성되고 도착 후 수 초 안에 조회할 수 있어야 해요.
  • 대시보드가 시간 창(지난 1시간, 지난 24시간)에 걸쳐 집계하며 국가, 기기, 제품 카테고리 같은 차원의 필터를 사용해요.
  • 높은 쿼리 동시성(수백 명의 동시 대시보드 사용자)에 P99 지연이 1초 미만이어야 해요.
  • 데이터가 append-only(추가 전용)이고 이전에 적재한 행을 갱신할 필요가 없어요. 갱신이 필요하다면 CDC / Upsert 파이프라인 플레이북을 대신 보세요.

아키텍처 스케치 (Architecture sketch)

Producers ──▶ Kafka topic ──▶ Pinot Real-Time Table ──▶ Broker ──▶ Dashboard
                                    │
                          ┌─────────┴─────────┐
                          │  Real-time servers │
                          │  (consuming segments)
                          └───────────────────┘

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

  • Kafka 토픽 — 로드 분산을 위해 고카디널리티 키(예: userId)로 파티셔닝해요.
  • 실시간 테이블 (Real-time table) — 서버마다 Kafka 파티션당 하나의 consuming 세그먼트를 가져요.
  • 스타트리 인덱스 (Star-tree index) — 가장 흔한 대시보드 group-by/filter 조합에 대해 사전 집계된 롤업(rollup)을 제공해요.
  • 브로커 (Brokers) — 쿼리를 서버로 팬아웃(fan out)하고 결과를 병합해요.

스키마 (Schema)

대시보드가 실행할 쿼리를 기준으로 스키마를 설계해요. 전형적인 제품 분석 이벤트 스트림은 다음과 같아요:

{
  "schemaName": "product_events",
  "dimensionFieldSpecs": [
    { "name": "eventType",  "dataType": "STRING" },
    { "name": "userId",     "dataType": "STRING" },
    { "name": "country",    "dataType": "STRING" },
    { "name": "device",     "dataType": "STRING" },
    { "name": "productId",  "dataType": "STRING" },
    { "name": "category",   "dataType": "STRING" },
    { "name": "sessionId",  "dataType": "STRING" }
  ],
  "metricFieldSpecs": [
    { "name": "revenue",    "dataType": "DOUBLE" },
    { "name": "quantity",   "dataType": "INT" }
  ],
  "dateTimeFieldSpecs": [
    {
      "name": "eventTimestamp",
      "dataType": "TIMESTAMP",
      "format": "1:MILLISECONDS:EPOCH",
      "granularity": "1:MILLISECONDS"
    }
  ]
}

참고: 스키마를 좁게 유지하세요. 추가하는 차원 컬럼마다 세그먼트 크기가 커지고 쿼리가 느려질 수 있어요. 대시보드가 실제로 필터하거나 그룹화하는 컬럼만 포함하세요. 설계 지침은 Schema and Table Shape를 참고하세요.

테이블 구성 (Table configuration)

{
  "tableName": "product_events",
  "tableType": "REALTIME",
  "segmentsConfig": {
    "timeColumnName": "eventTimestamp",
    "retentionTimeUnit": "DAYS",
    "retentionTimeValue": "30",
    "replication": "2",
    "segmentPushType": "APPEND"
  },
  "tableIndexConfig": {
    "loadMode": "MMAP",
    "invertedIndexColumns": ["eventType", "country", "device", "category"],
    "rangeIndexColumns": ["eventTimestamp"],
    "sortedColumn": ["eventType"],
    "bloomFilterColumns": ["userId"],
    "noDictionaryColumns": ["sessionId", "userId"],
    "starTreeIndexConfigs": [
      {
        "dimensionsSplitOrder": ["country", "device", "category", "eventType"],
        "skipStarNodeCreationForDimensions": [],
        "functionColumnPairs": [
          "COUNT__*",
          "SUM__revenue",
          "SUM__quantity"
        ],
        "maxLeafRecords": 10000
      }
    ],
    "streamConfigs": {
      "streamType": "kafka",
      "stream.kafka.topic.name": "product-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": {}
}

구성 요점 (Configuration highlights)

설정 이유
차원의 invertedIndexColumns 대시보드 필터 드롭다운에서 빠른 필터링을 가능하게 해요
eventTimestamp의 rangeIndexColumns WHERE eventTimestamp > ago('PT1H') 같은 시간 범위 조건(predicate)을 빠르게 해요
starTreeIndexConfigs 가장 흔한 롤업 쿼리를 사전 집계해 수 밀리초 안에 응답하게 해요
userId의 bloomFilterColumns 대시보드가 단일 사용자를 드릴다운할 때 포인트 조회를 가속화해요
고카디널리티 ID의 noDictionaryColumns 메모리를 절약해요. 수백만 개 고유 값을 가진 컬럼에 사전(dictionary) 인코딩은 낭비예요
retentionTimeValue: 30 30일보다 오래된 데이터를 자동으로 정리해요. 보존 요구에 맞게 조정하세요

인덱싱 전략 (Indexing strategy)

위에 보인 인덱스로 시작해서 반복적으로 다듬어요:

  1. WHERE 등가 필터에 쓰는 모든 컬럼에 역인덱스 (Inverted index).
  2. 시간 컬럼과 범위 필터에 쓰는 모든 숫자 컬럼에 범위 인덱스 (Range index).
  3. 가장 흔한 등가 필터가 걸리는 컬럼에 정렬 인덱스 (Sorted index) (테이블당 정렬 컬럼은 하나만).
  4. 트래픽을 지배하는 상위 3~5개 대시보드 쿼리에 스타트리 인덱스 (Star-tree index). 추가하기 전에 쿼리를 프로파일링하세요. 스타트리 구성마다 적재 오버헤드와 세그먼트 크기가 늘어나요.
  5. 매우 높은 카디널리티의 포인트 조회에 쓰는 컬럼에 블룸 필터 (Bloom filter).

결정 프레임워크는 Choosing Indexes를 참고하세요.

쿼리 패턴 (Query patterns)

대시보드 시계열 집계

SELECT
  DATETIMECONVERT(eventTimestamp, '1:MILLISECONDS:EPOCH', '1:MINUTES:EPOCH', '5:MINUTES') AS ts_bucket,
  country,
  COUNT(*) AS event_count,
  SUM(revenue) AS total_revenue
FROM product_events
WHERE eventTimestamp > ago('PT1H')
  AND eventType = 'PURCHASE'
GROUP BY ts_bucket, country
ORDER BY ts_bucket
LIMIT 1000

필터가 있는 Top-N

SELECT
  productId,
  SUM(revenue) AS total_revenue,
  COUNT(*) AS purchase_count
FROM product_events
WHERE eventTimestamp > ago('PT24H')
  AND country = 'US'
  AND eventType = 'PURCHASE'
GROUP BY productId
ORDER BY total_revenue DESC
LIMIT 20

단일 사용자 드릴다운

SELECT eventType, eventTimestamp, productId, revenue
FROM product_events
WHERE userId = 'u-123456'
  AND eventTimestamp > ago('PT7D')
ORDER BY eventTimestamp DESC
LIMIT 100

참고: 대시보드 쿼리에 OPTION(timeoutMs=5000)을 사용해서 단일 느린 쿼리가 커넥션 풀을 붙잡지 않게 하세요. Query Options를 참고하세요.

운영 체크리스트 (Operational checklist)

서비스 시작 전

  • Kafka 토픽 파티션이 원하는 병렬성을 충족하는지 확인하세요. Pinot은 파티션·서버당 consuming 세그먼트를 하나씩 만들어요.
  • realtime.segment.flush.threshold.rows와 realtime.segment.flush.threshold.time을 설정해서 완료 세그먼트 크기가 적정(300 MB–1 GB)하도록 하세요. 너무 작은 세그먼트는 메타데이터 오버헤드를 늘리고, 너무 큰 세그먼트는 쿼리를 느리게 해요.
  • 스타트리 인덱스가 가장 빈번한 쿼리를 커버하는지 EXPLAIN PLAN을 실행하고 출력에서 StarTreeIndex를 확인해 검증하세요.
  • 여러 팀이 클러스터를 공유한다면 테이블 수준 query quotas를 설정하세요.

모니터링

  • CONSUMING 세그먼트 래그: pinot.server.realtimeConsumptionCatchupRatio로 Kafka 컨슈머 래그를 모니터링하세요. 래그가 커지면 서버나 파티션을 더 추가하세요.
  • 쿼리 지연 P99: 브로커 메트릭으로 추적하세요. P99가 SLA를 넘으면 스타트리 인덱스를 추가하거나 브로커를 확장하는 걸 고려하세요.
  • 테이블당 세그먼트 수: 아주 많은 수의 작은 세그먼트는 쿼리 성능을 떨어뜨려요. 완료된 세그먼트를 압축하려면 Minion Merge Rollup Task를 사용하세요.
  • 힙 및 직접 메모리: 실시간 서버는 consuming 세그먼트를 메모리에 유지해요. JVM 힙과 off-heap 사용량을 모니터링하세요.

종합적인 메트릭 목록은 Monitoring와 Running Pinot in Production을 참고하세요.

흔한 함정 (Common pitfalls)

함정 해결책
대시보드 쿼리가 30일 전체 데이터를 스캔함 시간 범위 조건을 추가하세요. 없으면 브로커가 모든 세그먼트로 쿼리를 라우팅해요
스타트리가 활성화되지 않음 쿼리의 GROUP BY와 집계 함수가 스타트리 구성과 정확히 일치하는지 확인하세요
서버에서 잦은 GC 일시정지 고카디널리티 noDictionary 컬럼을 사전에서 빼거나, MMAP용 직접 메모리를 늘리세요
Kafka 리밸런스로 잠깐 쿼리 공백 발생 replication: 2를 설정해서 한 서버가 따라잡는 동안 다른 서버가 데이터를 제공하게 하세요

더 알아보기 (Learn more)