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)를 캡슐화해서 이 질의 타입들의 성숙함에 기대요. 질의를 두 가지 주요 단계로 실행해요.
- 내부 groupBy 또는 timeseries 질의를 실행해서 Aggregator(예: 일별 이벤트 수)를 계산해요.
- Broker에서 집계 결과를 넘겨 Averager(예: 일별 수의 7일 이동 평균)를 계산해요.
이 확장이 제공하는 주요 개선점:
- 기능(Functionality): Druid 질의 기능 확장 (즉 Window Functions의 초기 도입).
- 성능(Performance): 여러 세그먼트 스캔을 제거해 이러한 이동 집계의 성능을 개선.
Further reading
- Moving Average — 이동 평균 개념.
- Window Functions — 윈도우 함수.
- Analytic Functions — 분석 함수.
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:
doubleMeandoubleMeanNoNullsdoubleSumdoubleMaxdoubleMinlongMeanlongMeanNoNullslongSumlongMaxlongMin
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: 28cycleSize: 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)
- Moving Average와 Window Functions 문서를 참고해 주세요.
- GroupBy 질의와 Timeseries 질의 문서를 함께 보면 도움이 돼요.