커스텀 집계 함수 작성

커스텀 집계 함수 작성 (Writing Custom Aggregation Function)

Pinot에서 커스텀 집계 함수(Aggregation Function)를 작성하는 방법을 다루는 문서예요. Pinot에는 MIN, MAX, SUM, AVG 등 많은 내장 집계 함수가 있어요. 집계 함수 목록은 Pinot 쿼리하기 페이지를 보세요.

출처: 문서

본문

Pinot에는 MIN, MAX, SUM, AVG 등 많은 내장 집계 함수가 있어요. 집계 함수 목록은 Pinot 쿼리하기 페이지를 보세요.

새 AggregationFunction을 추가하려면 두 가지가 필요해요:

  • AggregationFunction 인터페이스를 구현하고 클래스패스의 일부로 사용 가능하게 만들기
  • AggregationFunctionFactory에 함수 등록. 현재로선 Pinot 내 코드 변경이 필요하지만, Pinot 코드를 변경하지 않고 함수를 플러그인할 수 있는 기능을 추가할 계획이에요.

전체적인 아이디어를 얻으려면 MAX 집계 함수 구현을 보세요. 다른 모든 구현은 여기에서 찾을 수 있어요.

AggregationFunction에서 구현해야 할 주요 메서드를 살펴봐요:

interface AggregationFunction {

  AggregationResultHolder createAggregationResultHolder();

  GroupByResultHolder createGroupByResultHolder(int initialCapacity, int maxCapacity);

  void aggregate(int length, AggregationResultHolder aggregationResultHolder, Map<String, BlockValSet> blockValSetMap);

  void aggregateGroupBySV(int length, int[] groupKeyArray, GroupByResultHolder groupByResultHolder,
      Map<String, BlockValSet> blockValSets);

  void aggregateGroupByMV(int length, int[][] groupKeysArray, GroupByResultHolder groupByResultHolder,
      Map<String, BlockValSet> blockValSets);

  IntermediateResult extractAggregationResult(AggregationResultHolder aggregationResultHolder);

  IntermediateResult extractGroupByResult(GroupByResultHolder groupByResultHolder, int groupKey);

  IntermediateResult merge(IntermediateResult intermediateResult1, IntermediateResult intermediateResult2);

  FinalResult extractFinalResult(IntermediateResult intermediateResult);

}

구현에 들어가기 전에 Pinot에서 Aggregation이 어떻게 작동하는지 이해하는 것이 중요해요.

이것은 고급 주제이며 Pinot 개념을 알고 있다고 가정해요. Pinot의 모든 데이터는 여러 노드의 세그먼트에 저장돼요. 쿼리 계획은 높은 수준에서 3가지 단계로 구성돼요.

1. Map 단계 (Map phase)

이 단계는 Pinot의 개별 세그먼트에서 작동해요.

  • 초기화: 쿼리 타입에 따라 결과 홀더를 설정하기 위해 다음 메서드가 호출돼요. 서로 다른 메서드와 반환 타입을 가지는 것이 복잡성을 더하지만 성능에 도움이 돼요.
  • 콜백: 쿼리의 필터 조건과 일치하는 각 레코드에 대해 다음 메서드 중 하나가 queryType(집계 vs group by)과 columnType(단일 값 vs 다중 값)에 따라 호출돼요. 성능상의 이유로 매 행마다가 아니라 레코드 배치에 대해 이 메서드를 호출하고, 가능하면 JVM이 실행의 일부를 벡터화할 수 있게 해준다는 점에 주의해요.
    • AGGREGATION: aggregate(int length, AggregationResultHolder aggregationResultHolder, Map<String,BlockValSet> blockValSetMap)
      • length: 블록의 길이를 나타냄. 일반적으로 < 10k
      • aggregationResultHolder: createAggregationResultHolder에서 반환된 객체
      • blockValSetMap: AggFunction의 인자에 따른 blockValSet의 맵
    • Group By 단일 값: aggregateGroupBySV(int length, int[] groupKeyArray, GroupByResultHolder groupByResultHolder, Map blockValSets)
      • length: 블록의 길이를 나타냄. 일반적으로 < 10k
      • groupKeyArray: Pinot는 내부적으로 값-정수 매핑을 유지하며 이 groupKeyArray는 내부 매핑에 매핑됨. 이 값들은 함께 고유 키를 형성.
      • groupByResultHolder: createGroupByResultHolder에서 반환된 객체
      • blockValSetMap: AggFunction의 인자에 따른 blockValSet의 맵
    • Group By 다중 값: aggregateGroupBySV(int length, int[] groupKeyArray, GroupByResultHolder groupByResultHolder, Map blockValSets)
      • length: 블록의 길이를 나타냄. 일반적으로 < 10k
      • groupKeyArray: Pinot는 내부적으로 값-정수 매핑을 유지하며 이 groupKeyArray는 내부 매핑에 매핑됨. 이 값들은 함께 고유 키를 형성.
      • groupByResultHolder: createGroupByResultHolder에서 반환된 객체
      • blockValSetMap: AggFunction의 인자에 따른 blockValSet의 맵

2. Combine 단계 (Combine phase)

이 단계에서 단일 pinot 서버 내의 모든 세그먼트 결과가 IntermediateResult로 결합돼요. IntermediateResult의 타입은 AggregationFunction 구현에서 정의된 Generic 타입에 기반해요.

public interface AggregationFunction<IntermediateResult, FinalResult extends Comparable> {

  IntermediateResult merge(IntermediateResult intermediateResult1, IntermediateResult intermediateResult2);

}

3. Reduce 단계 (Reduce phase)

Reduce 단계에는 두 단계가 있어요:

  • merge 함수를 사용해 다양한 서버의 모든 IntermediateResult 병합
  • extractFinalResult 메서드를 호출해 최종 결과 추출. 대부분의 경우 FinalResult는 IntermediateResult와 같은 타입이에요. AverageAggregationFunction은 IntermediateResult(AvgPair)가 FinalResult(Double)와 다른 경우의 예시예요.
  FinalResult extractFinalResult(IntermediateResult intermediateResult);

더 알아보기 (Learn more)