수집 집계

수집 집계 (Ingestion Aggregations)

실시간 데이터 수집 중에 집계를 수행해 저장 공간을 줄이고 쿼리 성능을 높이는 방법을 다루는 페이지예요.

출처: Ingestion Aggregations

본문

많은 데이터 분석 사용 사례는 집계된 데이터만 필요로 합니다. 예를 들어 차트에 사용되는 데이터는 차원 조합당 시간 버킷당 한 행으로 집계될 수 있습니다.

이렇게 하면 저장 공간이 훨씬 줄고 쿼리 성능이 좋아집니다. 테이블에 대해 이를 구성하는 것은 table config의 Aggregation Config를 통해 이루어집니다.

{% hint style="warning" %} 수집 집계는 실시간 Pinot 테이블에서만 동작한다는 점에 유의하세요. 또한 세그먼트 수준에서 수행됩니다. 교차 세그먼트 집계는 여전히 쿼리 시간 처리가 필요합니다. {% endhint %}

Aggregation Config

집계 설정은 실시간 데이터 수집 중 발생하는 집계를 제어합니다. 오프라인 집계는 별도로 처리해야 합니다.

아래는 테이블 설정의 수집 설정에 정의된 설정에 대한 설명입니다.

{
  "tableConfig": {
    "tableName": "...",
    "ingestionConfig": {
      "aggregationConfigs": [{
        "columnName": "aggregatedFieldName",
        "aggregationFunction": "<aggregationFunction>(<originalFieldName>)"
      }]
    }
  }
}

요구사항

수집 집계가 동작하려면 다음이 필요합니다:

  • 수집 집계 설정은 실시간 테이블에서만 유효합니다. (오프라인 테이블에 대한 수집 시간 집계 지원은 없습니다. Merge/Rollup Task를 사용하거나 Spark/MapReduce 같은 배치 처리 엔진으로 오프라인 데이터 흐름에서 집계를 전처리해야 합니다.)
  • 스트림 수집 타입은 lowLevel이어야 합니다.
  • 모든 메트릭에는 집계 설정이 있어야 합니다.
  • 모든 메트릭 컬럼은 단일 값이어야 하고 noDictionaryColumns로 구성되어야 합니다.
  • 집계 키로 사용되는 모든 차원 및 시간 컬럼은 단일 값이어야 합니다. 테이블 설정에서 dictionary 또는 no-dictionary 컬럼으로 구성할 수 있습니다.
  • 메트릭 집계는 upsert나 dedup과 함께 활성화할 수 없으며, 스키마에 COMPLEX 컬럼이 있으면 지원되지 않습니다.
  • aggregatedFieldName은 Pinot 스키마에 있어야 하고 originalFieldName은 Pinot 스키마에 존재하지 않아야 합니다.

{% hint style="info" %} 메트릭 집계가 활성화되면 Pinot는 소비 중인 세그먼트의 행을 차원·시간 키 컬럼의 dictionary id로 그룹핑합니다. 키 컬럼 중 하나가 테이블 설정에서 no-dictionary로 구성되어 있어도, 수집 집계가 동작하도록 Pinot는 소비 세그먼트에 그 컬럼의 일시적 dictionary를 여전히 생성합니다. 커밋된 세그먼트는 테이블 설정에서 재구축되므로, 커밋된 세그먼트는 여전히 구성된 no-dictionary 설정을 존중합니다. {% endhint %}

예제 시나리오

제품별 일일 판매 집계만 필요한 판매 데이터의 예입니다.

RealtimeQuickStart를 실행할 때도 dailySales라는 테이블에서 이를 찾을 수 있습니다.

예제 입력 데이터

{"customerID":205,"product_name": "car","price":1500.00,"timestamp":1571900400000}
{"customerID":206,"product_name": "truck","price":2200.00,"timestamp":1571900400000}
{"customerID":207,"product_name": "car","price":1300.00,"timestamp":1571900400000}
{"customerID":208,"product_name": "truck","price":700.00,"timestamp":1572418800000}
{"customerID":209,"product_name": "car","price":1100.00,"timestamp":1572505200000}
{"customerID":210,"product_name": "car","price":2100.00,"timestamp":1572505200000}
{"customerID":211,"product_name": "truck","price":800.00,"timestamp":1572678000000}
{"customerID":212,"product_name": "car","price":800.00,"timestamp":1572678000000}
{"customerID":213,"product_name": "car","price":1900.00,"timestamp":1572678000000}
{"customerID":214,"product_name": "car","price":1000.00,"timestamp":1572678000000}

스키마

스키마는 최종 테이블 구조만 반영한다는 점에 유의하세요.

{
  "schemaName": "dailySales",
  "dimensionFieldSpecs": [
    {
      "name": "product_name",
      "dataType": "STRING"
    }
  ],
  "metricFieldSpecs": [
    {
      "name": "sales_count",
      "dataType": "LONG"
    },
    {
      "name": "total_sales",
      "dataType": "DOUBLE"
    }
  ],
  "dateTimeFieldSpecs": [
    {
      "name": "daysSinceEpoch",
      "dataType": "LONG",
      "format": "1:MILLISECONDS:EPOCH",
      "granularity": "1:MILLISECONDS"
    }
  ]
}

테이블 설정

아래 집계 설정 예제에서 price는 입력 데이터에 있고 total_sales는 Pinot 스키마에 있다는 점에 유의하세요.

{
  "tableName": "daily_sales",
  "ingestionConfig": {
    "transformConfigs": [
      {
        "columnName": "daysSinceEpoch",
        "transformFunction": "toEpochDays(\"timestamp\")"
      }
    ],
    "aggregationConfigs": [
      {
        "columnName": "total_sales",
        "aggregationFunction": "SUM(price)"
      },
      {
        "columnName": "sales_count", 
        "aggregationFunction": "COUNT(*)"
      }
    ]
  }
  "tableIndexConfig": {
    "noDictionaryColumns": [
      "sales_count",
      "total_sales"
    ]
  }
}

최종 테이블 예제

product_name sales_count total_sales daysSinceEpoch
car 2 2800.00 18193
truck 1 2200.00 18193
truck 1 700.00 18199
car 2 3200.00 18200
truck 1 800.00 18202
car 3 3700.00 18202

허용되는 집계 함수

함수 이름 참고
MAX
MIN
SUM
COUNT COUNT(*)로 지정
DISTINCTCOUNTHLL DISTINCTCOUNTHLL(field, log2m)로 지정, 기본은 12. log2m 정의는 함수 레퍼런스 참고. 나중에 변경할 수 없으며 새 필드를 사용해야 합니다. 출력 필드의 스키마는 BYTES 타입이어야 합니다.
DISTINCTCOUNTHLLPLUS DISTINCTCOUNTHLLPLUS(field, s, p)로 지정. s와 p 정의는 함수 레퍼런스 참고, 나중에 변경할 수 없습니다. 출력 필드의 스키마는 BYTES 타입이어야 합니다.
SUMPRECISION SUMPRECISION(field, precision)로 지정, precision은 반드시 정의해야 합니다. 필드의 가능한 최대 크기를 계산하는 데 사용됩니다. 나중에 변경할 수 없으며 새 필드를 사용해야 합니다. 출력 필드의 스키마는 BIG_DECIMAL 타입이어야 합니다.

자주 묻는 질문

왜 Startree를 사용하지 않나요?

Startree는 실시간 세그먼트가 sealed된 후에만 추가할 수 있으며, startree를 만드는 것은 CPU 집약적입니다. 수집 집계는 소비 중인 세그먼트에 대해 동작하며 추가 CPU를 사용하지 않습니다.

Startree는 저장에 추가 메모리를 차지 하지만, 수집 집계는 원래 데이터셋보다 적은 데이터를 저장합니다.

수집 집계를 언제 사용하지 않나요?

비집계 형태의 원래 행이 필요하다면 수집 집계를 사용할 수 없습니다.

이미 aggregateMetrics 설정을 쓰고 있는데요?

aggregateMetrics는 수집 집계와 동일하게 동작하지만 SUM 함수만 허용합니다.

현재 변경은 하위 호환이 되므로, 다른 집계 함수가 필요하지 않다면 테이블 설정을 바꿀 필요가 없습니다.

이 설정은 오프라인 데이터에서는 동작하나요?

수집 집계는 실시간 수집에서만 동작합니다. 오프라인 데이터의 경우 오프라인 프로세스가 집계를 별도로 생성해야 합니다.

왜 모든 메트릭을 집계해야 하나요?

메트릭을 집계하지 않으면 고유 차원 집합마다 한 행을 초과하는 결과가 생깁니다.

집계 키 컬럼은 no-dictionary일 수 있나요?

네. 단일 값 차원·시간 컬럼은 no-dictionary 컬럼으로 구성해도 수집 집계에 참여할 수 있습니다. Pinot는 집계 키를 만들기 위해 소비 세그먼트에만 그 키 컬럼에 대한 일시적 dictionary를 생성합니다. 메트릭 컬럼은 다릅니다: 집계 값이 제자리에서 업데이트되므로 단일 값 no-dictionary 컬럼으로 유지되어야 합니다.

AggregationConfigs를 활성화했는데 왜 데이터가 안 보이나요?

  1. AggregationConfigs 없이 수집이 정상인지 확인하세요. 이는 문제를 격리하기 위함입니다.
  2. Pinot Server 로그에서 경고·오류 로그를 확인하세요. 특히 MutableSegmentImpl 클래스와 aggregateMetrics 메서드 관련 로그를 봅니다.
  3. JSON 데이터의 경우 숫자를 따옴표로 감싸지 마세요. 내부적으로 문자열로 파싱되어 sum 같은 값 기반 집계를 할 수 없습니다. 위 예제에서 {"customerID":205,"product_name": "car","price":"1500.00","timestamp":1571900400000} 행으로는 데이터 수집이 동작하지 않습니다. 핵심 문제는 price 숫자가 따옴표로 감싸져 보이지 않는다는 것입니다. 샘플 스택트레이스:
2024/11/04 00:24:27.760 ERROR [RealtimeSegmentDataManager_dailySales__0__0__20241104T0824Z] [dailySales__0__0__20241104T0824Z] Caught exception while indexing the record at offset: 9 , row: {
  "fieldToValueMap" : {
    "price" : "1000.00",
    "daysSinceEpoch" : 18202,
    "sales_count" : 0,
    "total_sales" : 0.0,
    "product_name" : "car",
    "timestamp" : 1572678000000
  },
  "nullValueFields" : [ "sales_count", "total_sales" ]
}
java.lang.ClassCastException: class java.lang.String cannot be cast to class java.lang.Number (java.lang.String and java.lang.Number are in module java.base of loader 'bootstrap')
	at org.apache.pinot.segment.local.aggregator.SumValueAggregator.applyRawValue(SumValueAggregator.java:25) ~[classes/:?]
	at org.apache.pinot.segment.local.indexsegment.mutable.MutableSegmentImpl.aggregateMetrics(MutableSegmentImpl.java:855) ~[classes/:?]
	at org.apache.pinot.segment.local.indexsegment.mutable.MutableSegmentImpl.index(MutableSegmentImpl.java:577) ~[classes/:?]
	at org.apache.pinot.core.data.manager.realtime.RealtimeSegmentDataManager.processStreamEvents(RealtimeSegmentDataManager.java:641) ~[classes/:?]
	at org.apache.pinot.core.data.manager.realtime.RealtimeSegmentDataManager.consumeLoop(RealtimeSegmentDataManager.java:477) ~[classes/:?]
	at org.apache.pinot.core.data.manager.realtime.RealtimeSegmentDataManager$PartitionConsumer.run(RealtimeSegmentDataManager.java:734) ~[classes/:?]
	at java.base/java.lang.Thread.run(Thread.java:1583) [?:?]

더 알아보기 (Learn more)