Moving Average Query

Moving Average Query (이동 평균 질의)

druid-moving-average-query 확장은 Druid 질의에서 이동 평균(Moving Average)과 그 외 집계 윈도우 함수(Aggregate Window Functions)를 지원해요. 이 확장은 표준 Druid 집계기를 소비해서 Averager라는 추가적인 윈도우 집계를 출력해요.

출처: 문서

본문

Moving Average Query는 Druid 질의에서 Moving Average와 그 외 Aggregate Window Functions를 지원하는 확장이에요.

이 Aggregate Window Functions는 표준 Druid Aggregator를 소비하고, Averager라고 불리는 추가적인 windowed aggregate를 출력해요.

High level algorithm

Moving Average는 groupBy 질의(dimension이 없으면 timeseries)를 캡슐화해서 이 질의 타입들의 성숙함에 기대요. 질의를 두 가지 주요 단계로 실행해요.

  1. 내부 groupBy 또는 timeseries 질의를 실행해서 Aggregator(예: 일별 이벤트 수)를 계산해요.
  2. Broker에서 집계 결과를 넘겨 Averager(예: 일별 수의 7일 이동 평균)를 계산해요.

이 확장이 제공하는 주요 개선점:

  • 기능(Functionality): Druid 질의 기능 확장 (즉 Window Functions의 초기 도입).
  • 성능(Performance): 여러 세그먼트 스캔을 제거해 이러한 이동 집계의 성능을 개선.

Further reading

Operations

Installation

Druid와 함께 제공되는 pull-deps 도구로 모든 Druid broker와 router 노드에 이 확장을 설치해 주세요.

java -classpath "<your_druid_dir>/lib/*" org.apache.druid.cli.Main tools pull-deps -c org.apache.druid.extensions.contrib:druid-moving-average-query:{VERSION}

Enabling

설치 후 이 확장을 활성화하려면 broker와 router의 runtime.properties 파일의 druid.extensions.loadList에 druid-moving-average-query를 추가하고 broker와 router 노드를 재시작해요.

예를 들어:

druid.extensions.loadList=["druid-moving-average-query"]

Configuration

현재 Moving Average에 특화된 구성 속성은 없어요.

Limitations

  • movingAverage는 다음 groupBy 속성들을 지원하지 않아요: subtotalsSpec, virtualColumns.
  • movingAverage는 다음 timeseries 속성을 지원하지 않아요: descending.
  • movingAverage의 averager는 달리 명시되지 않는 한 빈 버킷과 null 집계 값을 0으로 간주해요.

Query spec

query spec의 대부분의 속성은 groupBy query / timeseries에서 파생돼요. 이 질의 타입들의 문서를 참고해 주세요.

속성 설명 필수
queryType 이 String은 항상 "movingAverage"여야 해요. Druid가 질의를 어떻게 해석할지 가장 먼저 확인하는 값이에요. yes
dataSource 질의할 데이터 소스를 정의하는 String 또는 Object. 관계형 DB의 테이블과 매우 유사해요. DataSource 참고. yes
dimensions DimensionSpec의 JSON 리스트 (이 속성은 선택 항목이에요). no
limitSpec LimitSpec 참고. no
having Having 참고. no
granularity period granularity. Period Granularities 참고. yes
filter Filters 참고. no
aggregations Aggregations는 Averager의 입력이 돼요. Aggregations 참고. yes
postAggregations 입력으로 aggregations만 지원해요. Post Aggregations 참고. no
intervals ISO-8601 Intervals를 나타내는 JSON Object. 질의를 실행할 시간 범위를 정의해요. yes
context 특정 플래그를 지정할 수 있는 추가 JSON Object. no
averagers 이동 평균 함수를 정의해요. Averagers 참고. yes
postAveragers averagers와 aggregations 둘 다를 입력으로 지원해요. 문법은 postAggregations와 동일해요 (Post Aggregations 참고). no

Averagers

Averager는 Moving-Average 함수를 정의하는 데 사용돼요. Averager는 평균에만 국한되지 않아요. MAX()/MIN() 같은 다른 타입의 윈도우 함수도 제공할 수 있어요.

Common properties

모든 Averager에 공통인 속성들:

속성 설명 필수
type Averager 타입. Averager types 참고. yes
name Averager 이름. yes
fieldName 입력 이름 (집계 이름). yes
buckets 현재 버킷을 포함한 lookback 버킷(시간 기간) 수. 반드시 >0. yes
cycleSize Cycle 크기. 요일(day-of-week) 옵션 계산에 사용돼요. Cycle size (Day of Week) 참고. no, 기본값 1

Averager types

표준 averager:

  • doubleMean
  • doubleMeanNoNulls
  • doubleSum
  • doubleMax
  • doubleMin
  • longMean
  • longMeanNoNulls
  • longSum
  • longMax
  • longMin

Standard averagers

이 averager들은 네 가지 함수를 제공해요.

  • Mean (Average)
  • MeanNoNulls (빈 버킷 무시).
  • Sum
  • Max
  • Min

Null 무시하기: interval이 데이터셋 시작 시점에서 시작할 때는 MeanNoNulls averager를 사용하는 것이 유용해요. 그 경우 첫 기록들은 누락된 버킷을 무시하고 평균이 인위적으로 낮아지지 않아요. 다만 희소(sparse) 데이터셋의 빈 날짜들도 무시된다는 뜻이에요.

사용 예시:

{ "type" : "doubleMean", "name" : <output_name>, "fieldName": <input_name> }

Cycle size (Day of Week)

이 선택 파라미터는 모든 버킷 대신 각 사이클 안의 단일 버킷에 대해 계산하는 데 사용돼요. 대표적인 예가 주간(weekly) 버킷으로, 요일(Day of Week) 계산이 돼요. (다른 예: 월중 일, 시간 중 시각.)

즉, 다음 파라미터들을 사용하면:

  • granularity: period=P1D (일별)
  • buckets: 28
  • cycleSize: 7

각 출력 레코드 안에서 averager는 다음 버킷들에 대해 결과를 계산해요: 현재(#0), #7, #14, #21. 반면 cycleSize를 지정하지 않으면 28개 버킷 전체에 대해 계산했을 거예요.

Examples

모든 예시는 Druid 튜토리얼에서 제공하는 Wikipedia 데이터셋을 기반으로 해요.

Basic example

Wikipedia 편집 delta에 대한 7-버킷 이동 평균 계산.

질의 문법:

{
  "queryType": "movingAverage",
  "dataSource": "wikipedia",
  "granularity": {
    "type": "period",
    "period": "PT30M"
  },
  "intervals": [
    "2015-09-12T00:00:00Z/2015-09-13T00:00:00Z"
  ],
  "aggregations": [
    {
      "name": "delta30Min",
      "fieldName": "delta",
      "type": "longSum"
    }
  ],
  "averagers": [
    {
      "name": "trailing30MinChanges",
      "fieldName": "delta30Min",
      "type": "longMean",
      "buckets": 7
    }
  ]
}

결과:

[ {
   "version" : "v1",
   "timestamp" : "2015-09-12T00:30:00.000Z",
   "event" : {
     "delta30Min" : 30490,
     "trailing30MinChanges" : 4355.714285714285
   }
 }, {
   "version" : "v1",
   "timestamp" : "2015-09-12T01:00:00.000Z",
   "event" : {
     "delta30Min" : 96526,
     "trailing30MinChanges" : 18145.14285714286
   }
 }, {
...
...
...
}, {
  "version" : "v1",
  "timestamp" : "2015-09-12T23:00:00.000Z",
  "event" : {
    "delta30Min" : 119100,
    "trailing30MinChanges" : 198697.2857142857
  }
}, {
  "version" : "v1",
  "timestamp" : "2015-09-12T23:30:00.000Z",
  "event" : {
    "delta30Min" : 177882,
    "trailing30MinChanges" : 193890.0
  }
}

Post averager example

Wikipedia 편집 delta에 대한 7-버킷 이동 평균과, 현재 기간과 이동 평균 사이의 비율 계산.

질의 문법:

{
  "queryType": "movingAverage",
  "dataSource": "wikipedia",
  "granularity": {
    "type": "period",
    "period": "PT30M"
  },
  "intervals": [
    "2015-09-12T22:00:00Z/2015-09-13T00:00:00Z"
  ],
  "aggregations": [
    {
      "name": "delta30Min",
      "fieldName": "delta",
      "type": "longSum"
    }
  ],
  "averagers": [
    {
      "name": "trailing30MinChanges",
      "fieldName": "delta30Min",
      "type": "longMean",
      "buckets": 7
    }
  ],
  "postAveragers" : [
    {
      "name": "ratioTrailing30MinChanges",
      "type": "arithmetic",
      "fn": "/",
      "fields": [
        {
          "type": "fieldAccess",
          "fieldName": "delta30Min"
        },
        {
          "type": "fieldAccess",
          "fieldName": "trailing30MinChanges"
        }
      ]
    }
  ]
}

결과:

[ {
  "version" : "v1",
  "timestamp" : "2015-09-12T22:00:00.000Z",
  "event" : {
    "delta30Min" : 144269,
    "trailing30MinChanges" : 204088.14285714287,
    "ratioTrailing30MinChanges" : 0.7068955500319539
  }
}, {
  "version" : "v1",
  "timestamp" : "2015-09-12T22:30:00.000Z",
  "event" : {
    "delta30Min" : 242860,
    "trailing30MinChanges" : 214031.57142857142,
    "ratioTrailing30MinChanges" : 1.134692411867141
  }
}, {
  "version" : "v1",
  "timestamp" : "2015-09-12T23:00:00.000Z",
  "event" : {
    "delta30Min" : 119100,
    "trailing30MinChanges" : 198697.2857142857,
    "ratioTrailing30MinChanges" : 0.5994042624782422
  }
}, {
  "version" : "v1",
  "timestamp" : "2015-09-12T23:30:00.000Z",
  "event" : {
    "delta30Min" : 177882,
    "trailing30MinChanges" : 193890.0,
    "ratioTrailing30MinChanges" : 0.9174377224199288
  }
} ]

Cycle size example

지난 3시간의 매 첫 10분의 평균 계산.

질의 문법:

{
  "queryType": "movingAverage",
  "dataSource": "wikipedia",
  "granularity": {
    "type": "period",
    "period": "PT10M"
  },
  "intervals": [
    "2015-09-12T00:00:00Z/2015-09-13T00:00:00Z"
  ],
  "aggregations": [
    {
      "name": "delta10Min",
      "fieldName": "delta",
      "type": "doubleSum"
    }
  ],
  "averagers": [
    {
      "name": "trailing10MinPerHourChanges",
      "fieldName": "delta10Min",
      "type": "doubleMeanNoNulls",
      "buckets": 18,
      "cycleSize": 6
    }
  ]
}

더 알아보기 (Learn more)