스트림 수집 + Upsert

스트림 수집 + Upsert

Apache Pinot의 업서트 지원을 설명하는 페이지예요.

출처: Stream Ingestion with Upsert

본문

Pinot은 수집 중 네이티브 업서트를 지원해요. 레코드에 수정이 필요한 시나리오(라이드 요금 수정, 배달 상태 업데이트 등)가 있을 수 있어요.

부분 업서트(partial upsert)는 값이 바뀌는 컬럼만 지정하고 나머지는 무시하면 되므로 편리해요.

테이블 타입 지원

업서트는 REALTIME, OFFLINE, HYBRID 테이블 타입 모두에서 지원돼요. 사용 가능한 모드는 테이블 타입에 따라 달라져요:

테이블 타입 (Table type) FULL upsert PARTIAL upsert 참고 (Notes)
REALTIME 예 예 전체 업서트 기능을 갖춘 스트림 기반 수집
OFFLINE 예 아니오 배치 수집; 전체 행만 교체
HYBRID 예 아니오 오프라인과 실시간 시간 범위가 겹치지 않도록 해야 함

OFFLINE 테이블 업서트 설정에 대한 자세한 내용은 오프라인 테이블 업서트 (Offline Table Upsert) 참고.

Pinot의 업서트 개요

Pinot에서 업서트가 어떻게 동작하는지 개요를 보세요.

{% embed url="https://youtu.be/byzF91PQ6hE" %} Apache Pinot 1.0 Upserts overview {% endembed %}

Pinot에서 업서트 활성화

Pinot 테이블에서 업서트를 활성화하려면 다음을 수행해요:

  1. 스키마에서 기본 키 정의
  2. 테이블 설정에서 업서트 활성화

스키마에서 기본 키 정의

레코드를 업데이트하려면 레코드를 고유하게 식별하는 기본 키가 필요해요. 기본 키를 정의하려면 스키마 정의에 primaryKeyColumns 필드를 추가해요. 예를 들어 퀵스타트 예시의 UpsertMeetupRSVP 스키마 정의는 다음과 같아요.

{
    "primaryKeyColumns": ["event_id"]
}

이 필드는 기본 키가 복합(composite)일 수 있으므로 컬럼 목록을 기대해요.

같은 기본 키의 두 레코드가 수집되면 비교 값이 더 큰 레코드(기본적으로 timeColumn)가 사용돼요. 레코드들이 같은 기본 키와 이벤트 시간을 가지면 순서는 결정되지 않아요. 대부분의 경우 나중에 수집된 레코드가 사용되지만, 테이블에 정렬할 컬럼이 있는 경우에는 그렇지 않을 수 있어요.

{% hint style="warning" %} 입력 스트림을 기본 키로 파티셔닝하기

Pinot 업서트 테이블의 중요한 요구 사항은 입력 스트림을 기본 키로 파티셔닝하는 것입니다. Kafka 메시지의 경우 프로듀서가 send API에서 키를 설정해야 합니다. 원래 스트림이 파티셔닝되지 않았다면, Pinot 수집을 위해 입력 스트림을 셔플·재파티셔닝하는 스트리밍 처리 잡(예: Flink)이 필요합니다.

추가로 segmentPartitionConfig를 사용해 Broker 세그먼트 프루닝을 활용한다면, 사용하는 파티션 함수가 Kafka 프로듀서 측과 Pinot 양쪽에서 일치하도록 하는 것이 중요합니다. Java 클라이언트의 Kafka 기본값은 32비트 murmur2 해시이고, Python 같은 다른 언어에서는 CRC32 (Cyclic Redundancy Check 32-bit)입니다. {% endhint %}

테이블 설정에서 업서트 활성화

업서트를 활성화하려면 테이블 설정에서 다음 구성들을 해요.

업서트 모드

전체 업서트 (Full upsert)

업서트 모드는 기본적으로 FULL이에요. FULL upsert는 같은 기본 키를 가지면 새 레코드가 이전 레코드를 완전히 교체한다는 뜻이에요. 예시 설정:

{
  "upsertConfig": {
    "mode": "FULL"
  }
}

부분 업서트 (Partial upserts)

부분 업서트는 특정 컬럼만 업데이트하고 나머지는 무시하도록 선택할 수 있어요.

부분 업서트를 활성화하려면 mode를 PARTIAL로 설정하고 부분 업서트 컬럼에 partialUpsertStrategies를 지정해요. release-0.10.0부터 전략이 지정되지 않은 컬럼에는 OVERWRITE가 기본 전략으로 사용돼요. 모든 컬럼의 기본 전략을 바꾸기 위한 defaultPartialUpsertStrategy도 도입됐어요.

{% hint style="info" %} 부분 업서트가 동작하려면 null handling이 활성화되어 있어야 합니다. {% endhint %}

예를 들어:

release-0.8.0:

{
  "upsertConfig": {
    "mode": "PARTIAL",
    "partialUpsertStrategies":{
      "rsvp_count": "INCREMENT",
      "group_name": "IGNORE",
      "venue_name": "OVERWRITE"
    }
  },
  "tableIndexConfig": {
    "nullHandlingEnabled": true
  }
}

release-0.10.0:

{
  "upsertConfig": {
    "mode": "PARTIAL",
    "defaultPartialUpsertStrategy": "OVERWRITE",
    "partialUpsertStrategies":{
      "rsvp_count": "INCREMENT",
      "group_name": "IGNORE"
    }
  },
  "tableIndexConfig": {
    "nullHandlingEnabled": true
  }
}

부분 업서트의 커스텀 행 병합기

컬럼 수준 partialUpsertStrategies가 충분히 표현적이지 않다면 upsertConfig.partialUpsertMergerClass로 커스텀 행 병합기 클래스를 제공할 수 있어요.

{
  "upsertConfig": {
    "mode": "PARTIAL",
    "partialUpsertMergerClass": "org.apache.pinot.segment.local.upsert.merger.PartialUpsertMyCustomMerger"
  },
  "tableIndexConfig": {
    "nullHandlingEnabled": true
  }
}

partialUpsertMergerClass가 설정되면 Pinot은 내장 컬럼형 부분 업서트 병합기 대신 그 PartialUpsertMerger 구현을 인스턴스화해요. 커스텀 병합기 클래스는 서버 클래스패스에 있어야 하고 (List<String> primaryKeyColumns, List<String> comparisonColumns, UpsertConfig upsertConfig) 생성자를 노출해야 해요.

partialUpsertMergerClass는 partialUpsertStrategies와 상호 배타적이에요. Pinot은 둘을 동시에 설정하려는 테이블 설정을 거부해요.

Pinot은 다음 부분 업서트 전략을 지원해요:

전략 (Strategy) 설명 (Description)
OVERWRITE 마지막 레코드의 컬럼을 덮어씀
INCREMENT 새 값을 기존 값에 더함
APPEND 새 항목을 Pinot 순서 없는 집합에 추가
UNION 기존에 없으면 새 항목을 Pinot 순서 없는 집합에 추가
IGNORE 새 값을 무시하고 기존 값 유지 (v0.10.0+)
MAX 기존 값과 새 값 중 최댓값 유지 (v0.12.0+)
MIN 기존 값과 새 값 중 최솟값 유지 (v0.12.0+)
FORCE_OVERWRITE null을 포함해 기존 값을 항상 들어오는 값으로 교체 (v1.4.0+)

{% hint style="info" %} FORCE_OVERWRITE가 아닌 부분 업서트 전략의 경우, 기존 레코드나 들어오는 레코드 중 한쪽의 값이 null이면 Pinot은 업서트 전략을 무시하고 null이 아닌 값을 유지합니다:

(null, newValue) -> newValue

(oldValue, null) -> oldValue

(null, null) -> null {% endhint %}

들어오는 null이 이전에 저장된 값을 지워야 할 때 FORCE_OVERWRITE를 사용해요:

(oldValue, null) -> null

Post-Partial-Upsert 변환 (파생 컬럼)

부분 업서트를 사용할 때, 행이 들어오는 레코드와 기존 레코드에서 병합된 후 다시 계산해야 하는 파생 컬럼이 있을 수 있어요. postPartialUpsertTransformConfigs 기능을 사용하면 완전히 병합된 행에서 파생 컬럼을 계산하는 변환 함수를 적용할 수 있어요.

사용 사례

주문을 추적하는 전자상거래 테이블을 생각해 보세요:

  • order_id: 기본 키
  • score: 주문에서 얻은 포인트
  • bonus: 부여된 보너스 포인트
  • total: score + bonus여야 하는 파생 컬럼

부분 업서트에서 들어오는 레코드는 score 또는 bonus의 업데이트된 값만 포함할 수 있어요. 수집 시점 변환은 들어오는 레코드만 보므로 부분 병합된 행에서 total을 올바르게 계산할 수 없어요. postPartialUpsertTransformConfigs 기능을 사용하면 부분 업서트 병합 후 완전한 병합 행에서 total을 다시 계산할 수 있어요.

설정

post-partial-upsert 변환을 활성화하려면 테이블의 upsertConfig에 postPartialUpsertTransformConfigs 설정을 추가해요:

{
  "upsertConfig": {
    "mode": "PARTIAL",
    "defaultPartialUpsertStrategy": "OVERWRITE",
    "partialUpsertStrategies": {
      "score": "OVERWRITE",
      "bonus": "OVERWRITE"
    },
    "postPartialUpsertTransformConfigs": [
      {
        "columnName": "total",
        "transformFunction": "plus(score,bonus)"
      }
    ]
  },
  "tableIndexConfig": {
    "nullHandlingEnabled": true
  }
}

postPartialUpsertTransformConfigs는 수집 시점 변환과 같은 TransformConfig 모양을 사용해요: 각 항목은 대상 columnName과 transformFunction을 제공해요.

Pinot은 테이블을 받아들이기 전에 이 설정들을 검증해요:

  • PARTIAL 업서트 테이블에서만 지원됨.
  • 대상 컬럼은 스키마에 존재해야 함.
  • 대상 컬럼은 기본 키, 비교 컬럼, deleteRecordColumn, 또는 outOfOrderRecordColumn일 수 없음.
  • 각 대상 컬럼은 최대 한 번 나타날 수 있음.
  • 변환 함수는 자기 자신의 대상 컬럼을 참조할 수 없음.

평가 시맨틱

  • Post-partial-upsert 변환은 부분 업서트 병합이 완료된 후에 평가됨
  • 들어오는 레코드만이 아니라 완전한 병합 행에서 동작
  • 변환 표현식에는 들어오는 값과 기존 값 모두 사용 가능
  • 변환은 수집 시점 변환과 같은 함수 구문을 사용
  • 변환 결과는 최종 레코드의 일부로 파생 컬럼에 저장됨

수집 변환과의 상호작용

수집 시점 변환과 post-partial-upsert 변환은 서로 다른 목적을 제공해요:

측면 (Aspect) 수집 변환 (Ingestion Transforms) Post-Partial-Upsert 변환
실행 시점 (Execution timing) Pinot 수집 전 부분 업서트 병합 후, 수집 중
입력 레코드 (Input record) 들어오는 소스 레코드 병합된 행 (incoming + existing)
사용 사례 (Use case) 원시 입력 데이터 정규화/정리 병합된 상태에서 파생 컬럼 재계산
적용 대상 (Applies to) 모든 테이블 타입 (업서트/비-업서트) 부분 업서트 테이블만
예시 (Example) 타임스탬프 형식 변환 total = plus(score, bonus) (score·bonus는 병합 행에서)

둘을 함께 사용할 수 있어요:

  1. 수집 변환으로 들어오는 레코드 정규화
  2. 정규화된 들어오는 레코드가 부분 업서트 병합에 참여
  3. Post-partial-upsert 변환이 완전한 병합 행에서 파생 컬럼 재계산

예시 워크플로

이 설정을 가진 부분 업서트 테이블:

{
  "upsertConfig": {
    "mode": "PARTIAL",
    "partialUpsertStrategies": {
      "score": "OVERWRITE",
      "bonus": "OVERWRITE"
    },
    "postPartialUpsertTransformConfigs": [
      {
        "columnName": "total",
        "transformFunction": "plus(score,bonus)"
      }
    ]
  },
  "tableIndexConfig": {
    "nullHandlingEnabled": true
  }
}

이 레코드들을 처리하는 과정:

  1. 최초 레코드 (order_id=123):
    • 들어오는 값: {order_id: 123, score: 100, bonus: 10}
    • 병합: (첫 레코드, 기존 행 없음)
    • Post-transform: total = plus(100, 10) = 110
    • 최종: {order_id: 123, score: 100, bonus: 10, total: 110}
  2. 업데이트 레코드 (order_id=123):
    • 들어오는 값: {order_id: 123, score: 150} (score만 업데이트)
    • 병합: {order_id: 123, score: 150, bonus: 10} (bonus는 기존 행에서 보존)
    • Post-transform: total = plus(150, 10) = 160
    • 최종: {order_id: 123, score: 150, bonus: 10, total: 160}
  3. 또 다른 업데이트 (order_id=123):
    • 들어오는 값: {order_id: 123, bonus: 25} (bonus만 업데이트)
    • 병합: {order_id: 123, score: 150, bonus: 25} (score는 기존 행에서 보존)
    • Post-transform: total = plus(150, 25) = 175
    • 최종: {order_id: 123, score: 150, bonus: 25, total: 175}

{% hint style="info" %} post-partial-upsert 변환이 계산한 파생 컬럼은 다른 컬럼처럼 쿼리할 수 있습니다. 이런 파생 컬럼을 추가 업서트 전략이나 변환에 사용해야 한다면 스키마에 정의되어 있는지 확인하세요. {% endhint %}

None 업서트

mode를 NONE으로 설정하면 업서트가 비활성화돼요.

비교 컬럼 (Comparison column)

기본적으로 Pinot은 시간 컬럼(tableConfig의 timeColumn)의 값으로 최신 레코드를 결정해요. 즉 같은 기본 키의 두 레코드에 대해 시간 컬럼 값이 더 큰 레코드가 최신 업데이트로 선택돼요. 그러나 순서를 결정하는 데 다른 컬럼을 사용해야 하는 경우가 있어요. 그런 경우 comparisonColumn 옵션으로 비교에 사용할 컬럼을 재정의할 수 있어요. 예를 들어,

{
  "upsertConfig": {
    "mode": "FULL",
    "comparisonColumn": "anotherTimeColumn"
  }
}

부분 업서트 테이블의 경우 순서가 어긋난 이벤트는 소비·인덱싱되지 않아요. 예를 들어 같은 기본 키의 두 레코드 중 비교 컬럼 값이 더 작은 레코드가 다른 레코드보다 나중에 도착했다면 건너뛰어져요.

{% hint style="info" %} 참고: 단일 비교 컬럼에는 comparisonColumn 대신 comparisonColumns를 사용하세요. comparisonColumn은 현재 deprecated입니다. 이전 설정을 사용하면 unrecognizedProperties가 보일 수 있지만, 테이블 추가 시 자동으로 comparisonColumns로 변환됩니다. {% endhint %}

복수 비교 컬럼

특히 부분 업서트가 사용될 수 있는 경우, 데이터의 여러 프로듀서가 각자 상호 배타적인 컬럼 집합에 쓰고 기본 키만 공유하는 경우가 있어요. 이 경우 프로듀서 그룹마다 비교 컬럼을 하나씩 사용하면 각 그룹이 다른 프로듀서 그룹과 버전 관리를 조정할 필요 없이 자체 버전 시맨틱을 관리할 수 있어요.

{
  "upsertConfig": {
    "mode": "PARTIAL",
    "defaultPartialUpsertStrategy": "OVERWRITE",
    "partialUpsertStrategies":{},
    "comparisonColumns": ["secondsSinceEpoch", "otherComparisonColumn"]
  }
}

Pinot에 쓰여지는 문서는 comparisonColumns 집합 중 정확히 1개의 null이 아닌 값을 가질 것으로 기대돼요; 컬럼 중 1개 이상에 값이 있으면 문서가 거부돼요. 새 문서가 쓰일 때 null이 아닌 비교 컬럼은 이전 문서들 중 같은 기본 키로 본 그 비교 컬럼에만 대조돼요. 문서들이 배열에 지정된 순서대로 도착한다고 가정한 다음 예시를 생각해 보세요.

[
  {
    "event_id": "aa",
    "orderReceived": 1,
    "description" : "first",
    "secondsSinceEpoch": 1567205394
  },
  {
    "event_id": "aa",
    "orderReceived": 2,
    "description" : "update",
    "secondsSinceEpoch": 1567205397
  },
  {
    "event_id": "aa",
    "orderReceived": 3,
    "description" : "update",
    "secondsSinceEpoch": 1567205396
  },
  {
    "event_id": "aa",
    "orderReceived": 4,
    "description" : "first arrival, other column",
    "otherComparisonColumn": 1567205395
  },
  {
    "event_id": "aa",
    "orderReceived": 5,
    "description" : "late arrival, other column",
    "otherComparisonColumn": 1567205392
  },
  {
    "event_id": "aa",
    "orderReceived": 6,
    "description" : "update, other column",
    "otherComparisonColumn": 1567205398
  }
]

다음이 발생해요:

  1. orderReceived: 1
  • 결과: 유지됨
  • 이유: 기본 키 "aa"에 대한 첫 문서
  1. orderReceived: 2
  • 결과: 유지됨 (orderReceived: 1 대체)
  • 이유: 비교 컬럼 (secondsSinceEpoch)이 이전에 본 것보다 큼
  1. orderReceived: 3
  • 결과: 거부됨
  • 이유: 비교 컬럼 (secondsSinceEpoch)이 이전에 본 것보다 작음
  1. orderReceived: 4
  • 결과: 유지됨 (orderReceived: 2 대체)
  • 이유: 비교 컬럼 (otherComparisonColumn)이 이전에 본 것보다 큼 (이전에 본 적 없음), 값이 secondsSinceEpoch에서 본 것보다 작음에도 불구하고
  1. orderReceived: 5
  • 결과: 거부됨
  • 이유: 비교 컬럼 (otherComparisonColumn)이 이전에 본 것보다 작음
  1. orderReceived: 6
  • 결과: 유지됨 (orderReceived: 4 대체)
  • 이유: 비교 컬럼 (otherComparisonColumn)이 이전에 본 것보다 큼

메타데이터 TTL

Pinot에서 메타데이터 맵은 힙 메모리에 저장돼요. 인메모리 데이터를 줄이고 성능을 개선하려면 기본 키 항목이 메타데이터 맵에 저장되는 시간(메타데이터 TTL)을 최소화해요. TTL 제한은 카디널리티가 높고 업데이트가 빈번한 기본 키에 특히 유용해요.

메타데이터 TTL은 첫 번째 비교 컬럼에 적용되므로, 업서트 TTL의 시간 단위는 첫 번째 비교 컬럼과 같아요.

기본 키가 메타데이터에 저장되는 기간 구성

기본 키가 메타데이터에 저장되는 기간을 구성하려면 metadataTTL에 시간 길이를 지정해요. 예를 들어:

{
  "upsertConfig": {
    "mode": "FULL",
    "snapshot": "ENABLE",
    "preload": "ENABLE",
    "metadataTTL": 86400
  }
}

이 예시에서 Pinot은 기본 키를 메타데이터에 1일 동안 보유해요.

인메모리 validDocsIDs 복구를 위한 메타데이터 TTL에는 업서트 스냅샷 활성화가 필요하다는 점에 유의하세요.

삭제 컬럼 (Delete column)

Upsert Pinot 테이블은 기본 키의 소프트 삭제(soft-delete)를 지원해요. 이를 위해 들어오는 레코드가 기본 키의 삭제 마커 역할을 하는 전용 부울 단일-필드 컬럼을 포함해야 해요. 실시간 엔진이 delete 컬럼이 true로 설정된 레코드를 만나면 그 기본 키는 더 이상 쿼리 가능한 document 집합의 일부가 아니게 돼요. 즉 기본 키는 쿼리 옵션 skipUpsert=true로 명시적으로 요청하지 않는 한 쿼리에 보이지 않아요.

{ 
    "upsertConfig": {  
        ... 
        "deleteRecordColumn": <column_name>
    } 
}

delete 컬럼은 단일 값 부울 컬럼이어야 해요.

// In the Schema
{
    ...
    {
      "name": "<delete_column_name>",
      "dataType": "BOOLEAN"
    },
    ...
}

{% hint style="warning" %} 기존 업서트 테이블에서는 deleteRecordColumn을 불변으로 취급하세요. 서버가 업서트 상태를 초기화할 때 delete-column 선택을 캐시하기 때문에, Pinot은 강제 업데이트를 하지 않는 한 컨트롤러 업데이트 API로 이를 추가·제거·변경하는 것을 거부합니다. {% endhint %}

삭제된 기본 키는 같은 기본 키를 갖되 비교 컬럼 값이 더 높은 레코드를 수집하면 되살아날 수 있어요.

부분 업서트 테이블에서 기본 키를 되살릴 때, 되살아난 레코드는 모든 컬럼의 진실 원본으로 취급된다는 점에 유의하세요. 즉 이전의 모든 컬럼 업데이트는 무시되고 새 레코드의 값으로 덮어써져요.

삭제 키 TTL (Deleted Keys time-to-live)

위의 deleteRecordColumn 설정은 기본 키를 소프트 삭제할 뿐이에요. 인메모리 데이터를 줄이고 성능을 개선하려면 삭제된-기본-키 항목이 메타데이터 맵에 저장되는 시간(삭제 키 TTL)을 최소화해요. TTL 제한은 미래 업데이트가 예상되지 않는 삭제된-기본-키에 특히 유용해요.

삭제된-기본-키가 메타데이터에 저장되는 기간 구성

기본 키가 메타데이터에 저장되는 기간을 구성하려면 deletedKeysTTL에 시간 길이를 지정해요. 예를 들어:

  "upsertConfig": {
    "mode": "FULL",
    "deleteRecordColumn": <column_name>,
    "deletedKeysTTL": 86400
  }

이 예시에서 Pinot은 삭제된-기본-키를 메타데이터에 1일 동안 보유해요.

{% hint style="info" %} 이 deletedKeysTTL 필드의 값은 비교 컬럼의 단위와 같아야 합니다. 비교 컬럼이 초에 해당하는 값을 갖는다면, 이 설정도 초 단위여야 합니다 (위 예시 참고). metadataTTL과 deletedKeysTTL은 복수 비교 컬럼에서는 동작하지 않으며, 비교/시간 컬럼은 NUMERIC 타입이어야 합니다. {% endhint %}

삭제와 컴팩션을 함께 사용할 때의 데이터 일관성

deletedKeysTTL을 UpsertCompactionTask와 함께 사용하면, 삭제 레코드(deleteRecordColumn = true가 기본 키에 설정됨)를 포함하는 세그먼트가 먼저 컴팩션되고 이전의 오래된 레코드는 아직 컴팩션되지 않는 시나리오가 있을 수 있어요. 서버 재시작 중 이제 이전 오래된 레코드가 메타데이터 매니저 맵에 추가되고 삭제되지 않은 것으로 취급돼요. 이 시나리오의 데이터 불일치를 막기 위해 enableDeletedKeysCompactionConsistency라는 새 설정을 추가했는데, true로 설정하면 삭제된 기본 키에 대한 다른 모든 세그먼트의 이전 레코드가 컴팩션될 때까지 삭제 레코드가 컴팩션되지 않도록 보장해요.

{
  "upsertConfig": {
    "mode": "FULL",
    "deleteRecordColumn": <column_name>,
    "deletedKeysTTL": 86400,
    "enableDeletedKeysCompactionConsistency": true
  }
}

쿼리와 업서트가 동시에 발생할 때의 데이터 일관성

Pinot의 업서트는 실시간 업데이트를 가능하게 하고 쿼리가 항상 레코드의 최신 버전을 검색하도록 보장하므로, 변경 가능한(mutable) 데이터를 효율적으로 관리하는 강력한 기능이에요. 그러나 매우 높은 QPS와 높은 수집율의 애플리케이션에서는 쿼리와 업서트가 동시에 발생하면서 때로 쿼리 결과에 불일치가 생길 수 있어요.

예를 들어 1백만 기본 키를 가진 테이블을 생각해 보세요. 새 레코드가 수집되고 이전 레코드가 무효화되는 방식과 무관하게 distinct count 쿼리는 항상 1백만을 반환해야 해요. 그러나 높은 수집·쿼리율에서는 쿼리가 때로 백만보다 약간 위/아래로 반환할 수 있어요. 이것은 쿼리가 여러 세그먼트에서 validDocIds 비트맵(현재 유효한 document를 나타내는)을 획득해 유효 레코드를 결정하기 때문이에요. 이 비트맵 획득이 진행 중인 업서트와 원자적이지 않으므로, 쿼리가 데이터의 일관성 없는 뷰를 포착해 유효 레코드를 과다/과소 집계할 수 있어요.

이것은 읽기와 쓰기가 동시에 발생해 일시적인 불일치를 초래하는 전형적인 동시성 문제예요. 보통 이러한 문제는 락이나 스냅샷으로 쿼리 실행 중 안정적인 데이터 뷰를 유지해 해결해요. 이를 해결하기 위해 업서트 활성화 테이블에 SYNC와 SNAPSHOT이라는 두 가지 새 일관성 모드가 도입되어, 쿼리와 업서트가 동시에 그리고 매우 높은 처리량으로 발생해도 일관된 쿼리 결과를 보장해요.

기본적으로 일관성 모드는 NONE이며, 이전처럼 시스템이 동작해요. SYNC 모드는 쿼리 실행 중 업서트를 차단함으로써 일관성을 보장해, 쿼리가 항상 안정적인 업서트 데이터 뷰를 보도록 해요. 그러나 쓰기 지연을 유발할 수 있어요. 대안으로 SNAPSHOT 모드는 쿼리가 사용할 validDocIds 비트맵의 일관된 스냅샷을 생성해요. 이렇게 하면 쿼리를 차단하지 않고 업서트가 계속되므로, 쿼리와 쓰기율이 모두 높은 워크로드에 더 적합해요.

이 새 일관성 모드들은 애플리케이션이 특정 요구 사항에 따라 일관성 보장과 성능 트레이드오프 사이에서 균형을 맞출 수 있는 유연성을 제공해요.

{
  "upsertConfig": {
    "consistencyMode": "SYNC", // or "SNAPSHOT", "NONE"
...
  }
}

SNAPSHOT 모드의 경우 upsertViewRefreshIntervalMs라는 upsertConfig(기본 3000ms)로 업서트 뷰를 얼마나 자주 리프레시해야 하는지 구성할 수 있어요. 쓰기와 쿼리 스레드 모두 이 설정에 따라 뷰가 낡으면 업서트 뷰를 리프레시할 수 있어요. 이 설정을 바꾸면 서버 재시작이 필요해요.

Pinot은 또한 newSegmentTrackingTimeMs(기본 10000) 동안 서버에서 새롭게 추가된 세그먼트를 추적해요. 그 윈도우 동안 브로커 라우팅이 따라잡는 동안 Pinot이 그 새로 추가된 세그먼트를 선택적 세그먼트로 포함할 수 있어, 세그먼트 추가 직후 쿼리가 더 완전한 업서트 뷰를 보도록 도와요. newSegmentTrackingTimeMs를 0으로 설정하면 이 추적이 비활성화돼요. consistencyMode가 SYNC 또는 SNAPSHOT이면 newSegmentTrackingTimeMs는 양수여야 해요.

쿼리 시간 중 서버 재시작 없이 뷰의 신선도를 조정하려면 upsertViewFreshnessMs라는 쿼리 옵션을 사용할 수 있어요. 기본적으로 이 쿼리 옵션은 upsertConfig upsertViewRefreshIntervalMs와 일치하지만, 쿼리가 더 작은 값으로 설정하면 그 쿼리에 대해 업서트 뷰가 더 일찍 리프레시될 수 있어요; 0으로 설정하면 쿼리가 매번 업서트 뷰를 강제로 리프레시해요.

디버깅 목적에는 skipUpsertView라는 쿼리 옵션이 있어요. true로 설정하면 SYNC 또는 SNAPSHOT 모드가 유지하는 일관된 업서트 뷰를 우회해요. 이로써 NONE 모드인 것처럼 쿼리가 실행돼요.

라우팅에 strictReplicaGroup 사용

Upsert Pinot 테이블은 입력 스트림에 대해 저수준 소비자만 사용할 수 있어요. 결과적으로 세그먼트에 파티셔닝된 replica-group 배정을 암묵적으로 사용해요. 게다가 업서트는 같은 파티션의 모든 세그먼트가 같은 서버에서 서빙되어야 세그먼트 간 데이터 일관성이 보장된다는 추가 요구를 제기해요. 따라서 라우팅 전략으로 strictReplicaGroup을 사용해야 해요. 그러려면 Routing에서 instanceSelectorType을 다음과 같이 구성해요:

{
  "routing": {
    "instanceSelectorType": "strictReplicaGroup"
  }
}

{% hint style="warning" %} 저수준 소비자의 암묵적 파티셔닝된 replica-group 배정을 사용하면 인스턴스 배정(파티션에서 서버로의 매핑)이 ZooKeeper에 지속되지 않고, 새로 추가된 서버가 명시적 재배정(보통 rebalance) 없이 자동으로 포함됩니다. 이로 인해 같은 파티션의 새 세그먼트가 다른 서버에 배정되어 업서트 요구 사항이 깨질 수 있습니다.

이를 방지하려면 명시적 파티셔닝된 replica-group 인스턴스 배정을 사용해 인스턴스 배정이 지속되도록 권장합니다. replicaGroupPartitionConfig에서 numInstancesPerPartition은 항상 1이어야 합니다. {% endhint %}

{% hint style="warning" %} 업서트 테이블에서 Adaptive Server Selection을 활성화하지 마세요. 적응형 선택기는 strictReplicaGroup 경계를 존중하지 않으며 여러 replica-group에 걸쳐 서버로 쿼리를 라우팅해, 같은 파티션의 모든 세그먼트가 같은 서버에서 서빙되어야 한다는 업서트 일관성 보장을 깨뜨릴 수 있습니다. 자세한 내용은 #12507 참고. {% endhint %}

업서트 메타데이터 복구를 위한 validDocIds 스냅샷 활성화

업서트 스냅샷 지원은 release-0.12.0에도 추가됐어요. 스냅샷을 활성화하려면 snapshot을 ENABLE로 설정해요. 예를 들어:

{
  "upsertConfig": {
    "mode": "FULL",
    "snapshot": "ENABLE"
  }
}

업서트는 특정 세그먼트에서 어떤 docId가 유효한지(ValidDocIndexes)를 포함하는 메타데이터를 메모리에 유지해요. 이 메타데이터는 서버 재시작 시 손실되며 다시 재생성되어야 해요.

ValidDocIndexes는 TTL이 지나 기본 키가 제거된 후에는 쉽게 복구할 수 없어요. 스냅샷을 활성화하면 불변 세그먼트에 대한 validDocIds 스냅샷을 저장·복구하는 함수를 추가해 이 문제를 해결해요.

스냅샷은 갑작스러운 종료 시 지속된 데이터와 일관성을 보장하기 위해 모든 세그먼트 커밋 시에 찍혀요.

재시작 시 서버 부팅 시간을 빠르게 하려면 이 기능을 활성화하는 것을 권장해요.

{% hint style="info" %} metadataTTL 또는 deletedKeysTTL을 사용하는 업서트 테이블의 경우, 세그먼트 리로드는 불변 세그먼트의 모든 행을 다시 스캔하는 대신 지속된 validDocIds 스냅샷에서 업서트 메타데이터를 재구축합니다. 이는 TTL 만료나 삭제 처리가 이미 업서트 메타데이터에서 제거한 키를 리로드가 되살리지 않도록 방지합니다.

리로드가 세그먼트의 다른 복사본을 다운로드해야 한다면, 로컬 docId 기반 스냅샷이 다운로드된 행과 더 이상 일치하지 않을 수 있으므로 세그먼트 CRC가 바뀌면 Pinot은 리로드를 실패시킵니다. 이는 forceDownload=true 리로드와 CRC 기반 재다운로드를 트리거하는 일반 리로드에 적용됩니다. {% endhint %}

{% hint style="info" %} validDocIds 스냅샷의 수명주기는 다음과 같습니다,

  1. 스냅샷이 활성화되면, 다음 소비 중인 세그먼트가 시작될 때 기존 세그먼트의 스냅샷이 찍히거나 리프레시됩니다.
  2. 스냅샷 파일은 세그먼트가 제거될 때(예: 데이터 보존 또는 수동 삭제)까지 디스크에 유지됩니다.
  3. 스냅샷이 비활성화되면, 세그먼트의 기존 스냅샷은 세그먼트가 서버에 로드될 때(예: 서버 재시작) 정리됩니다. {% endhint %}

더 빠른 서버 재시작을 위한 프리로드 활성화

업서트 프리로드는 서버가 재시작될 때 업서트 상태를 더 빠르게 복원할 수 있게 해요. REALTIME 테이블과, PR #19271을 포함하는 빌드의 OFFLINE 테이블에 적용돼요. 프리로드 기능을 활성화하려면 preload를 ENABLE로 설정해요. 스냅샷도 활성화되어야 해요. 예를 들어:

{
  "upsertConfig": {
    "mode": "FULL",
    "snapshot": "ENABLE",
    "preload": "ENABLE"
  }
}

내부적으로 validDocIds 스냅샷을 사용해 유효한 document를 식별하고, 전체 업서트 비교 흐름을 수행하는 대신 업서트 메타데이터를 빠르게 복원해요. 이 흐름은 서버가 ready로 표시되기 전에 트리거되며, 그 후 서버는 스냅샷 없는 나머지 세그먼트 로드를 시작해요(그래서 프리로드라는 이름).

이 기능은 또한 서버 설정에서 pinot.server.instance.max.segment.preload.threads: N을 지정해야 해요(N은 프리로드에 사용할 스레드 수). 기본값은 0으로 프리로드 기능을 비활성화해요.

{% hint style="warning" %} v1.2.0에서 스냅샷과 프리로드 복구가 활성화되었지만 max.segment.preload.threads가 0으로 남아 있을 때, 프리로딩 메커니즘은 여전히 활성화되지만 프리로딩할 스레드가 없어 세그먼트 로드가 실패하는 버그가 도입되었습니다. 이는 최신 버전에서 수정됐지만, v1.2.0에서는 max.segment.preload.threads도 양수로 설정하는 것을 기억하세요. 설정 변경 적용에는 서버 재시작이 필요합니다. {% endhint %}

스토리지 최적화를 위한 커밋 시간 컴팩션 활성화

{% hint style="warning" %} 기존 테이블에 커밋 시간 컴팩션을 활성화한다면, 먼저 그 테이블의 수집을 일시정지하고, 테이블 설정을 업데이트해 이 기능을 활성화한 다음, 수집을 재개하는 것을 권장합니다. {% endhint %}

많은 Upsert 사용 사례는 세그먼트 커밋 윈도우 안에 많은 Update 이벤트가 있어요. 예를 들어 Uber Eats 주문의 주문 상태에 대한 Upsert 테이블이 있다면, 1시간 윈도우 안에 같은 주문에 대한 많은 업데이트 이벤트가 있을 것으로 예상할 수 있어요. 이런 사용 사례에서 커밋된 세그먼트는 많은 죽은 튜플로 끝나고, 그것을 정리하려면 몇 시간이 걸릴 수 있는 Segment Compaction 태스크를 기다려야 해요.

커밋 시간 컴팩션(commit time compaction)은 세그먼트 커밋 프로세스 자체 중에 무효·오래된 레코드를 제거하는 업서트 테이블용 성능 최적화 기능이에요. 이는 테이블의 스토리지 팽창을 즉시 줄여줄 뿐 아니라 세그먼트 커밋 시간도 낮출 수 있어요.

커밋 시간 컴팩션을 활성화하려면 업서트 설정에서 enableCommitTimeCompaction을 true로 설정해요. 예를 들어:

{
  "upsertConfig": {
    "mode": "FULL",
    "enableCommitTimeCompaction": true
  }
}

동작 방식

세그먼트 커밋 중 커밋 시간 컴팩션은:

  • 무효 document ID를 걸러냄. 유효 레코드와 소프트 삭제 레코드를 유지.
  • 컴팩션된 세그먼트에 대한 정확한 컬럼 통계 생성
  • 오래된 데이터 제거 시 올바른 document 순서 유지
  • minion 태스크 없이 세그먼트 크기를 즉시 줄임

설정 요구 사항

  • 기능은 업서트 설정에서 enableCommitTimeCompaction=true로 설정해 테이블별 활성화
  • 변경은 한 세그먼트 커밋 주기 후에 적용됨 (현재 소비 중인 세그먼트는 컴팩션 없이 커밋됨)
  • 모든 종류의 업서트 테이블과 호환

순서 어긋난 이벤트 처리

순서 어긋난 이벤트 처리와 관련된 설정 2개가 추가됐어요.

dropOutOfOrderRecord

순서 어긋난 레코드 버리기를 활성화하려면 dropOutOfOrderRecord를 true로 설정해요. 예를 들어:

{
  "upsertConfig": {
    ...,
    "dropOutOfOrderRecord": true
  }
}

이 기능은 어떤 순서 어긋난 이벤트도 소비 중인 세그먼트에 지속하지 않아요. 지정하지 않으면 기본값은 false예요.

  • false일 때 순서 어긋난 레코드는 소비 중인 세그먼트에 지속되지만, MetadataManager 매핑은 업데이트되지 않으므로 이 레코드는 쿼리나 향후 업데이트에서 참조되지 않아요. skipUpsert 쿼리 옵션을 사용하면 이 레코드를 여전히 볼 수 있어요.
  • true일 때 순서 어긋난 레코드는 전혀 지속되지 않고 MetadataManager 매핑도 업데이트되지 않으므로 쿼리나 향후 업데이트에서 참조되지 않아요. skipUpsert 쿼리 옵션을 사용해도 이 레코드를 볼 수 없어요.

outOfOrderRecordColumn

이것은 순서 어긋난 이벤트를 프로그래밍 방식으로 식별하기 위한 것이에요. 이 설정을 활성화하려면 테이블 스키마에 isOutOfOrder 같은 부울 필드를 추가하고 이 설정으로 활성화해요. 예를 들어:

{
  "upsertConfig": {
    ...,
    "outOfOrderRecordColumn": "isOutOfOrder"
  }
}

이 기능은 이벤트의 순서성에 따라 isOutOfOrder 필드에 true / false 값을 지속해요. skipUpsert를 사용하는 동안 순서 어긋난 이벤트를 걸러 혼란을 피할 수 있어요. 예를 들어:

select key, val from tbl1 where isOutOfOrder = false option(skipUpsert=false)

{% hint style="info" %} dropOutOfOrderRecord와 outOfOrderRecordColumn은 consistencyMode가 설정되지 않았을 때만 지원됩니다 (즉 consistencyMode = NONE). consistencyMode가 활성화되면 유효 document가 업데이트되기 전에 행이 추가되기 때문입니다. 결과적으로 순서 어긋난 레코드는 업서트 테이블에서 버려지거나 표시될 수 없어, 이 옵션들의 목적이 무의미해집니다. {% endhint %}

커스텀 메타데이터 매니저 사용

Pinot은 레코드와 세그먼트 업데이트를 처리하는 커스텀 PartitionUpsertMetadataManager를 지원해요.

{
  "upsertConfig": {
    "metadataManagerClass": org.apache.pinot.segment.local.upsert.CustomPartitionUpsertMetadataManager
  }
}

커스텀 업서트 매니저 추가

다음과 같이 커스텀 PartitionUpsertMetadataManager를 추가할 수 있어요:

  • 새 Java 프로젝트 생성. 패키지 이름을 org.apache.pinot.segment.local.upsert.xxx로 유지.
  • Java 프로젝트에 의존성 포함

Maven:

<dependency>
  <groupId>org.apache.pinot</groupId>
  <artifactId>pinot-segment-local</artifactId>
  <version>1.0.0</version>
 </dependency>

Gradle:

include 'org.apache.pinot:pinot-common:1.0.0'
  • PartitionUpsertMetadataManager 인터페이스를 구현하는 커스텀 파티션 매니저 추가
//Example custom partition manager

class CustomPartitionUpsertMetadataManager implements PartitionUpsertMetadataManager {}
  • BaseTableUpsertMetadataManager 인터페이스를 구현하는 커스텀 TableUpsertMetadataManager 추가
//Example custom table upsert metadata manager

public class CustomTableUpsertMetadataManager extends BaseTableUpsertMetadataManager {}
  • 컴파일된 JAR을 pinot의 /plugins 디렉터리에 배치. 이미 실행 중이면 모든 Pinot 인스턴스를 재시작해야 함.
  • 이제 테이블 설정에서 커스텀 업서트 매니저를 다음과 같이 사용할 수 있음:
{
  "upsertConfig": {
    "metadataManagerClass": org.apache.pinot.segment.local.upsert.CustomPartitionUpsertMetadataManager
  }
}

:warning: 업서트 매니저 클래스 이름은 대소문자를 구분하지 않아요.

불변 업서트 설정 필드

{% hint style="danger" %} 일부 업서트·스키마 설정 필드는 테이블 생성 후 수정할 수 없습니다.

기존 업서트 테이블에서 이 필드를 변경하면, 특히 서버가 재시작하고 세그먼트를 커밋할 때 데이터 불일치 또는 데이터 손실이 발생할 수 있습니다. Pinot은 이 설정들을 기반으로 document를 검증·무효화하므로, 데이터가 수집된 후 이를 변경하면 기존 validDocId 스냅샷이 새 설정과 불일치하게 됩니다.

테이블 생성 후 불변인 필드:

스키마 필드:

  • primaryKeyColumns

upsertConfig 필드:

  • mode (FULL, PARTIAL, NONE)
  • hashFunction
  • comparisonColumns
  • timeColumnName (기본 비교 컬럼으로 사용될 때)
  • deleteRecordColumn
  • dropOutOfOrderRecord
  • outOfOrderRecordColumn

이 필드를 업데이트하려 하면 오류가 반환됩니다:

Failed to update table '<tableName>': Cannot modify [<field>] as it may lead to data inconsistencies. Please create a new table instead.

권장 우회법: 원하는 설정으로 새 테이블을 만들고 모든 데이터를 다시 수집하세요.

대안 (주의해서 사용): 테이블을 재생성하지 않고 이 필드를 수정하려면 테이블 설정 업데이트 API에서 force=true 쿼리 파라미터를 사용할 수 있습니다. 그 전에 upsertConfig에서 SNAPSHOT 모드를 비활성화하고, 소비를 일시정지하고, 모든 서버를 재시작하세요. 이 방법은 새로 수집되는 키에 대해서만 일관성을 보장하며, 기존 데이터는 불일치할 수 있습니다. {% endhint %}

{% hint style="warning" %} PARTIAL 업서트 테이블의 경우, Pinot은 이제 기존 테이블에서 컨트롤러 업데이트 API를 통해 partialUpsertStrategies와 defaultPartialUpsertStrategy를 업데이트하는 것을 허용합니다.

이 업데이트는 소급 적용되지 않습니다:

  • 기존 병합 값은 이미 저장된 그대로 유지됩니다.
  • 새 전략은 각 소비 서버가 재시작되고 partial-upsert 핸들러를 재구축한 후에만 적용됩니다.
  • 롤링 재시작 중에는 복제본이 일시적으로 다른 전략 버전으로 소비해 새로 병합된 행에서 분기할 수 있습니다.

롤아웃 후 행 값 드리프트가 관찰되면, Pinot이 공통 세그먼트 데이터에서 재구축할 수 있도록 영향을 받는 소비 중인 세그먼트를 리셋하세요. 세그먼트 수명주기와 복구 (Segment Lifecycle and Repair) 참고. {% endhint %}

업서트 테이블 제한 사항

업서트 Pinot 테이블에는 몇 가지 제한이 있어요.

  • 부분 업서트는 REALTIME 테이블에서만 지원돼요. OFFLINE 테이블은 FULL 업서트만 지원해요. 자세한 내용은 오프라인 테이블 업서트 (Offline Table Upsert) 참고.
  • star-tree 인덱스는 인덱싱에 사용할 수 없어요. star-tree 인덱스가 수집 중 사전 집계를 수행하기 때문이에요.
  • append-only 테이블과 달리, 순서 어긋난 이벤트(들어오는 레코드의 비교 값이 최신 사용 가능 값보다 작은 경우)는 Pinot 부분 업서트 테이블에서 소비·인덱싱되지 않으며, 이런 늦은 이벤트는 건너뛰어져요.
  • 업서트/중복제거 테이블이 생성된 후에는 소스 토픽의 파티션 수를 변경할 수 없어요 (모범 사례에서 언급한 것처럼 상대적으로 높은 파티션 수로 시작하세요).

불일치 처리

소비 중인 세그먼트가 커밋되면 서버는 변경 가능한 세그먼트를 새 불변 세그먼트로 교체해요. 이 전환 중에 인메모리 업서트 메타데이터(기본 키 → 최신 레코드 위치)가 복제본 간에 분기할 수 있어요.

이 분기는 FULL Upsert 테이블에서는 복제본이 결국 수렴하므로 일반적으로 안전하지만, 다음에는 안전하지 않아요:

  • 부분 업서트 테이블: 병합 정확성은 정확한 "최신" 레코드 위치에 의존함; 잘못된 포인터는 새 항목에 잘못된 값을 도입할 수 있음.
  • dropOutOfOrderRecord=true 또는 outOfOrderRecordColumn을 사용하는 전체 업서트 테이블: 순서 어긋남 감지가 현재 위치에 의존함; 잘못된 메타데이터는 잘못된 승인·거부를 일으킬 수 있음.

이를 완화하기 위해 Helix 클러스터 설정(프로세스별 pinot-server.conf 키가 아님)인 pinot.server.consuming.segment.consistency.mode를 추가했어요. 세 가지 모드:

RESTRICTED (기본)

Partial Upsert와 DropOutOfOrder 테이블에 대해 force-commit과 reload를 차단해요. 소비 중인 세그먼트는 자연스럽게만 커밋할 수 있어요. 일관성을 보장해요.

PROTECTED

교체 후 조정(reconciliation)을 위해 키의 이전 불변 세그먼트 위치를 추적하는 임시 맵을 사용해 force-commit/reload를 허용해요.

조정:

  • 여전히 교체된 세그먼트를 가리키는 키 → 이전 불변 위치로 복귀.
  • 이전 위치가 없는 키 → 제거.
  • 조정 불가능한 키 → 기록되고, 사용자가 조치하도록 메트릭이 방출. ParallelSegmentConsumptionPolicy가 항상 ∈ {DISALLOW_ALWAYS, ALLOW_DURING_BUILD_ONLY}인지 확인.

UNSAFE

조정 없이 force-commit/reload를 허용해요. 커밋 중 불일치 가능성이 있어요. 이 모드는 안전하지 않으며 프로덕션 설정에서 권장되지 않아요.

모니터링

  • pinot.server.tableName.realtimeUpsertInconsistentRows: 전체 업서트 테이블에 대해 세그먼트 교체가 복제본 간 불일치 메타데이터를 감지한 후 교체되지 않은 기본 키 수. dropOutOfOrderRecord=true 또는 outOfOrderRecordColumn을 사용하는 테이블 포함.
  • pinot.server.tableName.partialUpsertKeysNotReplaced: 부분 업서트 테이블에 대해 세그먼트 교체가 복제본 간 불일치 메타데이터를 감지한 후 교체되지 않은 기본 키 수.

모범 사례

다른 실시간 테이블과 달리 Upsert 테이블은 레코드 위치를 메모리에 기록해야 하므로 더 많은 메모리 리소스를 차지해요. 따라서 사전에 용량을 계획하고 리소스 사용을 모니터링하는 것이 중요해요. Upsert 테이블 사용 권장 사항은 다음과 같아요.

더 많은 파티션으로 토픽/스트림 생성.

입력 스트림의 파티션 수가 Pinot 테이블의 파티션 수를 결정해요. 입력 토픽/스트림에 파티션이 많을수록 Pinot 테이블을 분산시킬 Pinot 서버가 많아져 더 수평으로 확장할 수 있어요. 주의 업서트 활성화 테이블은 향후 파티션을 늘릴 수 없으므로 충분히 좋은 파티션 수로 시작해야 해요 (pinot 서버 수의 최소 2-3배).

메모리 사용

Upsert 테이블은 기본 키에서 레코드 위치로의 인메모리 맵을 유지해요. 따라서 단순한 기본 키 타입을 사용하고 메모리 비용을 줄이기 위해 복합 기본 키를 피하는 것을 권장해요. JSON 컬럼을 기본 키로 사용할 때 주의하세요. 같은 키-값이라도 순서가 다르면 다른 기본 키로 간주됩니다. 추가로 Upsert 설정의 hashFunction 설정을 고려해 보세요. UUID, MD5, MURMUR3가 될 수 있어요.

기본 키 컬럼이 유효한 UUID이고 많은 기본 키 때문에 메모리가 부족하다면, UUID 해시 함수가 해시 충돌 위험 없이 메모리 요구를 최대 35% 낮출 수 있어요.

기본 키가 유효한 UUID가 아니면, 이 해시 함수는 기본 키를 있는 그대로 저장하고 UUID 기반 압축을 건너뛰어요.

MD5와 MURMUR3도 메모리 요구를 낮추는 데 도움이 돼요. 모든 종류의 기본 키 값에 대해 작동하지만 작은 해시 충돌 위험이 있어요. MD5와 MURMUR3의 생성 해시는 128비트이므로, 기본 키 값이 128비트보다 클 때 유용해요.

모니터링

테이블 파티션의 기본 키 수를 보기 위해 pinot.server.upsertPrimaryKeysCount.tableName 메트릭 위에 대시보드를 설정해요. 메모리 사용 증가에 비례하는 성장을 추적하는 데 유용해요. ****** 업서트의 총 메모리 사용은 대략 (primaryKeysCount * (sizeOfKeyInBytes + 24))예요.

용량 계획

나중에 리소스 제약에 부딪히지 않도록 사전에 용량을 계획하는 것이 유용해요. 간단한 방법은 입력 스트림의 파티션별 기본 키 비율을 측정하고 그 데이터를 특정 기간(테이블 보존 기반)으로 외삽해 메모리 사용을 근사하는 것이에요. 힙 덤프도 upsert 테이블 인스턴스의 지금까지 메모리 사용을 확인하는 데 유용해요.

예시

이를 종합하면 퀵스타트 예시의 테이블 설정은 다음과 같아요:

{
  "tableName": "upsertMeetupRsvp",
  "tableType": "REALTIME",
  "tenants": {},
  "segmentsConfig": {
    "timeColumnName": "mtime",
    "retentionTimeUnit": "DAYS",
    "retentionTimeValue": "1",
    "replication": "1"
  },
  "tableIndexConfig": {
    "segmentPartitionConfig": {
      "columnPartitionMap": {
        "event_id": {
          "functionName": "Hashcode",
          "numPartitions": 2
        }
      }
    }
  },
  "instanceAssignmentConfigMap": {
    "CONSUMING": {
      "tagPoolConfig": {
        "tag": "DefaultTenant_REALTIME"
      },
      "replicaGroupPartitionConfig": {
        "replicaGroupBased": true,
        "numReplicaGroups": 1,
        "partitionColumn": "event_id",
        "numPartitions": 2,
        "numInstancesPerPartition": 1
      }
    }
  },
  "routing": {
    "segmentPrunerTypes": [
      "partition"
    ],
    "instanceSelectorType": "strictReplicaGroup"
  },
  "ingestionConfig": {
    "streamIngestionConfig": {
      "streamConfigMaps": [
        {
          "streamType": "kafka",
          "stream.kafka.topic.name": "upsertMeetupRSVPEvents",
          "stream.kafka.decoder.class.name": "org.apache.pinot.plugin.inputformat.json.JSONMessageDecoder",
          "stream.kafka.consumer.factory.class.name": "org.apache.pinot.plugin.stream.kafka30.KafkaConsumerFactory",
          "stream.kafka.broker.list": "localhost:19092"
        }
      ]
    }
  },
  "upsertConfig": {
    "mode": "FULL",
    "snapshot": "ENABLE",
    "preload": "ENABLE"
  },
  "fieldConfigList": [
    {
      "name": "location",
      "encodingType": "RAW",
      "indexType": "H3",
      "properties": {
        "resolutions": "5"
      }
    }
  ],
  "metadata": {
    "customConfigs": {}
  }
}
{
  "tableName": "upsertPartialMeetupRsvp",
  "tableType": "REALTIME",
  "tenants": {},
  "segmentsConfig": {
    "timeColumnName": "mtime",
    "retentionTimeUnit": "DAYS",
    "retentionTimeValue": "1",
    "replication": "1"
  },
  "tableIndexConfig": {
    "segmentPartitionConfig": {
      "columnPartitionMap": {
        "event_id": {
          "functionName": "Hashcode",
          "numPartitions": 2
        }
      }
    },
    "nullHandlingEnabled": true
  },
  "instanceAssignmentConfigMap": {
    "CONSUMING": {
      "tagPoolConfig": {
        "tag": "DefaultTenant_REALTIME"
      },
      "replicaGroupPartitionConfig": {
        "replicaGroupBased": true,
        "numReplicaGroups": 1,
        "partitionColumn": "event_id",
        "numPartitions": 2,
        "numInstancesPerPartition": 1
      }
    }
  },
  "routing": {
    "segmentPrunerTypes": [
      "partition"
    ],
    "instanceSelectorType": "strictReplicaGroup"
  },
  "ingestionConfig": {
    "streamIngestionConfig": {
      "streamConfigMaps": [
        {
          "streamType": "kafka",
          "stream.kafka.topic.name": "upsertPartialMeetupRSVPEvents",
          "stream.kafka.decoder.class.name": "org.apache.pinot.plugin.inputformat.json.JSONMessageDecoder",
          "stream.kafka.consumer.factory.class.name": "org.apache.pinot.plugin.stream.kafka30.KafkaConsumerFactory",
          "stream.kafka.broker.list": "localhost:19092"
        }
      ]
    }
  },
  "upsertConfig": {
    "mode": "PARTIAL",
    "partialUpsertStrategies": {
      "rsvp_count": "INCREMENT",
      "group_name": "UNION",
      "venue_name": "APPEND"
    }
  },
  "fieldConfigList": [
    {
      "name": "location",
      "encodingType": "RAW",
      "indexType": "H3",
      "properties": {
        "resolutions": "5"
      }
    }
  ],
  "metadata": {
    "customConfigs": {}
  }
}

{% hint style="info" %} Pinot 서버는 업서트 활성화 테이블에서 서빙되는 모든 세그먼트에 걸쳐 기본 키에서 레코드 위치로의 맵을 유지합니다. 결과적으로 기존 업서트 테이블의 설정을 업데이트할 때(예: 기본 키의 컬럼 변경, 비교 컬럼 변경) 변경을 적용하고 맵을 재구축하려면 서버를 재시작해야 합니다. {% endhint %}

실시간 수집 OOM 보호

Pinot은 JVM 힙 사용이 높을 때 서버 측 실시간 수집 역압(backpressure)을 적용할 수 있어요. 기본 서버 모드 pinot.server.instance.ingestion.oom.protection.mode=DISABLE은 기능을 꺼두며, UPSERT_DEDUP_ONLY로 설정하면 기본적으로 실시간 업서트·중복제거 테이블을 보호하거나, ingestionConfig.streamIngestionConfig.oomProtection으로 개별 테이블을 옵트인/아웃할 수 있어요.

전체 설정 키, 기본값, 런타임 동작, 메트릭은 자동 쿼리 킬링을 사용한 OOM 보호 참고.

커밋, 다운로드, 교체 중 병렬 소비

부분 업서트 테이블의 경우 Pinot은 이전 세그먼트가 아직 확정되는 동안 다음 소비 중인 세그먼트를 일시정지해, 복제본이 커밋 처리 중 다른 병합 행 상태로 진행되지 않도록 할 수 있어요.

  • 기본적으로 부분 업서트 테이블은 커밋 처리 중 병렬로 소비를 계속하지 않아요.
  • 일시정지 없는(pauseless) 소비가 활성화되면 Pinot은 parallelSegmentConsumptionPolicy에 따라 다운로드와 교체 중에는 중지하면서 빌드 단계 중에는 계속할 수 있어요.

하위 호환성을 위해 부분 업서트 테이블은 여전히 deprecated 테이블 수준 플래그 upsertConfig.allowPartialUpsertConsumptionDuringCommit을 받아들여요. true로 설정하면 이전 동작으로 복원되어 세그먼트 다운로드·교체를 포함해 커밋 처리 내내 복제본이 소비를 계속할 수 있게 해요:

{
  "upsertConfig": {
    "mode": "PARTIAL",
    "allowPartialUpsertConsumptionDuringCommit": true
  }
}

테이블 수준 플래그가 기본 false로 남으면 서버 수준 폴백 pinot.server.instance.upsert.default.allow.partial.upsert.consumption.during.commit이 해당 서버에서 부분 업서트 테이블에 대해 같은 레거시 동작을 활성화할 수 있어요. 새 배포는 병렬 소비를 명시적으로 제어해야 할 때 streamIngestionConfig의 parallelSegmentConsumptionPolicy를 선호해야 해요.

고급 서버 설정

소비 중인 세그먼트 일관성 모드

부분 업서트 테이블 또는 dropOutOfOrderRecord=true나 outOfOrderRecordColumn이 구성된 테이블의 경우, Helix 클러스터 설정 키 pinot.server.consuming.segment.consistency.mode(컨트롤러 cluster-config API/UI로 설정, 예: POST /cluster/configs)로 컨트롤러와 서버가 세그먼트 리로드·force commit을 어떻게 처리할지 구성해요. 이 키를 pinot-server.conf에만 두는 것은 효과가 없어요 — 컨트롤러와 서버 모두 ConsumingSegmentConsistencyModeListener를 통해 클러스터 설정에서 그것을 로드해요.

모드 (Mode) 설명 (Description)
RESTRICTED (기본) dropOutOfOrderRecord나 outOfOrderRecordColumn이 구성된 부분 업서트 테이블과 업서트 테이블의 소비 중인 세그먼트 리로드를 건너뛰고 명시적 force commit을 거부.
PROTECTED 세그먼트 교체 중 업서트 메타데이터 복귀와 함께 리로드/force commit을 활성화. ParallelSegmentConsumptionPolicy가 DISALLOW_ALWAYS 또는 ALLOW_DURING_BUILD_ONLY로 설정 필요.
UNSAFE 메타데이터 복귀 없이 리로드를 허용. 불일치가 허용되거나 외부에서 처리되는 경우에만 사용.

참고: 이 클러스터 설정은 단일 서버에서 쿼리 vs 업서트 동시성을 제어하는 테이블 수준 upsertConfig.consistencyMode 설정(SYNC / SNAPSHOT / NONE)과 구별됩니다.

Deprecated 설정 필드에서 마이그레이션

Pinot 1.4.0부터 다음 업서트 설정 필드가 이름이 바뀌었어요:

Deprecated 필드 새 필드 (New field) 값 (Values)
enableSnapshot snapshot ENABLE, DISABLE, 또는 DEFAULT
enablePreload preload ENABLE, DISABLE, 또는 DEFAULT

새 필드는 부울 값 대신 Enablement 열거형(ENABLE, DISABLE, DEFAULT)을 사용해요. DEFAULT는 서버 수준 설정에 위임하며, 기능이 인스턴스 수준에서 활성화될 때 테이블 수준 재정의를 허용해요.

Deprecated된 부울 필드는 여전히 동작하지만 향후 릴리스에서 제거될 예정이에요. 테이블 설정을 새 필드 이름으로 업데이트하세요.

퀵스타트 (Quick Start)

전체 업서트가 어떻게 동작하는지 설명하기 위해 Pinot 바이너리는 퀵스타트 예시와 함께 제공돼요. 다음 명령으로 실시간 업서트 테이블 meetupRSVP를 생성해요.

# stop previous quick start cluster, if any
bin/quick-start-upsert-streaming.sh

부분 업서트 데모는 다음 명령으로도 실행할 수 있어요.

# stop previous quick start cluster, if any
bin/quick-start-partial-upsert-streaming.sh

스트림으로 데이터가 흐르는 즉시 Pinot 테이블이 그것을 소비하고 쿼리 준비가 돼요. Query Console로 가서 실시간 데이터를 확인해 보세요.

부분 업서트의 경우 지정된 부분 업서트 전략에 따라 설정된 컬럼의 값만 바뀐 것을 볼 수 있어요.

부분 업서트 예시가 아래에 나와 있어요. 수집 중 각 event_id는 고유하게 유지되는 동시에 rsvp_count 값은 증가해요.

업서트가 아닌 테이블과의 차이를 보려면 쿼리 옵션 skipUpsert를 사용해 쿼리 결과에서 업서트 효과를 건너뛸 수 있어요.

FAQ

기존 업서트 테이블을 재생성하지 않고 스키마 컬럼을 추가할 수 있나요?

추가적(additive) 컬럼이면 예. 스키마를 업데이트한 다음 리로드 및/또는 forceCommit하여 새 소비 중인 세그먼트가 컬럼을 받도록 해요. 부분 업서트 테이블과 순서 어긋남 처리 구성 업서트 테이블은 Helix 클러스터 설정 pinot.server.consuming.segment.consistency.mode가 허용할 때만 force commit을 제한합니다 (pinot-server.conf가 아니라 소비 중인 세그먼트 일관성 모드 참고). 업서트 설정 변경은 다릅니다: 허용된 부분 업서트 전략 변경은 통제된 재시작이 필요하고 소급 적용되지 않으며, 기본/비교 컬럼과 모드 같은 불변 설정은 새 테이블과 재수집이 필요해요. 스키마 진화 결정표 참고.

기존 업서트 테이블에서 기본 키 컬럼, 비교 컬럼 같은 설정을 변경할 수 있나요?

권장하지 않아요. 기존 세그먼트는 이전 설정으로 계산된 validDocId 스냅샷을 포함해요. 설정을 변경하면 기존 스냅샷이 정리되지 않아 데이터 불일치가 생길 수 있는데, 특히 복제 서버는 validDocId 스냅샷 없이 재시작할 때 그렇지 않지만 한 서버가 validDocId 스냅샷과 함께 재시작되는 경우가 그렇지 않아요.

변경을 피할 것: 기본 키 컬럼, 비교 컬럼, 업서트 모드, hashFunction, deleteRecordColumn.

Pinot은 이제 컨트롤러 업데이트 API에서 이 가드를 강제해요. 기본적으로 PUT /tables/{tableName}과 PUT /tableConfigs/{tableName}은 하위 호환되지 않는 업서트 또는 dedup 설정 변경을 400 Bad Request로 거부해요. 업서트 테이블의 경우 비교 컬럼, 해시 함수, 모드, deleteRecordColumn, 순서 어긋남 설정, 그리고 Pinot이 기본 비교 컬럼으로 사용할 때의 테이블 시간 컬럼이 포함돼요. dedup 테이블의 경우 dedup 해시 함수, dedup 시간 컬럼, 그리고 Pinot이 기본 dedup 시간 컬럼으로 사용할 때의 테이블 시간 컬럼이 포함돼요.

partialUpsertStrategies와 defaultPartialUpsertStrategy는 PARTIAL 업서트 테이블에서 예외예요. Pinot은 force=true 없이 그 업데이트를 받아들여요. 하지만 새 전략은 각 서버가 재시작되고 테이블 설정을 다시 로드한 후에 발생하는 병합에만 영향을 미쳐요. 기존 병합 값은 자동으로 다시 쓰여지지 않고, 롤링 재시작 중 복제본이 일시적으로 분기할 수 있어요. 그런 경우 롤아웃 후 영향을 받는 소비 중인 세그먼트를 리셋하세요.

PUT /tables/{tableName}에서 force=true 또는 PUT /tableConfigs/{tableName}에서 forceTableSchemaUpdate=true로 가드를 우회할 수 있지만, Pinot은 통제된 복구나 마이그레이션 워크플로에서만 사용을 권장해요.

변경이 불가피하다면:

최선의 옵션: 새 테이블을 만들고 모든 데이터를 다시 수집하세요.

대안: SNAPSHOT을 비활성화하고, 소비를 일시정지하고, 모든 서버를 재시작하세요. 이는 새로 들어오는 키에서만 동작하며, 기존 데이터의 일관성은 보장되지 않아요.

더 알아보기 (Learn more)