- 数据库
- OLAP
- 大数据
- 后端
【免费下载链接】druid
Apache Druid: a high performance real-time analytics database.
导读
本文聚焦 Apache Druid 社区扩展druid-moving-average-query(Moving Average Query),它让 Druid 原生查询首次具备移动平均(Moving Average)及其他聚合型窗口函数(Aggregate Window Functions)的能力。文章将以扩展官方文档与源码为骨架,讲解其两阶段执行算法、安装启用方式、完整 Query Spec 字段、Averager 类型体系与cycleSize周期语义,并结合仓库源码剖析其内部实现细节与已知限制,帮助你在 Broker 侧无需多次扫描 Segment 即可完成滚动窗口聚合分析。
扩展概述:是什么,解决什么问题
Moving Average Query 是一个 Druid 扩展,位于仓库 extensions-contrib/moving-average-query 模块,其官方说明见 docs/development/extensions-contrib/moving-average-query.md。
它解决的问题非常具体:Druid 原生的 groupBy / timeseries 查询擅长对时间桶内的数据做聚合,但无法跨时间桶计算滚动窗口(例如"过去 7 天的平均编辑量")。该扩展通过引入Averager(Averager,窗口聚合器)概念,把标准 Druid Aggregator 的输出进一步加工成带窗口语义的聚合结果。
该扩展带来的两个核心增强:
- 功能增强:为 Druid 查询引入窗口函数能力,支持
MEAN、SUM、MAX、MIN等聚合窗口计算; - 性能优化:通过一次 Segment 扫描 + Broker 侧窗口计算,消除多次查询拼接滚动窗口的重复扫描开销。
高层算法:两阶段流水线
从源码 MovingAverageQueryRunner.java 的类注释可以清晰看到整个执行流程。Moving Average Query 在内部封装了 groupBy 查询(无维度时退化为 timeseries 查询),以复用这两类成熟查询的能力,其执行分为两大阶段:
- 阶段一(内层查询):运行一个内层 groupBy(有维度)或 timeseries(无维度)查询,先计算出基础聚合值,例如"每天的编辑次数";
- 阶段二(Broker 侧窗口计算):在 Broker 上对聚合结果按时间桶(period bucket)滑动,计算 Averager,例如"每天编辑次数的 7 天移动平均"。
Runner 中更详细的步骤为:
- 取所有 Averager 中最大的
buckets值,将查询区间起始时间向前回退buckets - 1个 period(见 MovingAverageQueryRunner.java),保证窗口有足够的回看数据; - 按维度有无选择 groupBy 或 timeseries 内层查询;
- 用
RowBucketIterable将各维度组合的行按 period 分桶(RowBucket); - 交给
MovingAverageIterable执行窗口计算、写入 Averager 结果列; - 通过
PostAveragerAggregatorCalculator应用 postAveragers; - 过滤掉回看区间产生的、超出用户请求区间之外的行;
- 最后应用 having / 排序 / limit(
applyLimit)。
安装与启用
安装(Installation)
使用 Druid 自带的 pull-deps 工具在所有Broker 和 Router 节点上安装该社区扩展(Community Extension),命令如下:
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}其中{VERSION}需替换为与你的 Druid 版本一致的扩展版本号。该扩展的 Maven 坐标为org.apache.druid.extensions.contrib:druid-moving-average-query,模块定义见 extensions-contrib/moving-average-query/pom.xml(当前仓库版本为31.0.0-SNAPSHOT),它依赖druid-processing、druid-server等核心模块。
启用(Enabling)
安装完成后,在 Broker 和 Router 节点的runtime.properties中把druid-moving-average-query加入druid.extensions.loadList,然后重启 Broker 与 Router 节点:
druid.extensions.loadList=["druid-moving-average-query"]关于社区扩展的加载机制可参考 extensions 文档。
配置(Configuration)
目前Moving Average 没有任何专属配置属性,其行为完全由查询 JSON 本身驱动。
查询规范(Query Spec)
Moving Average Query 的大部分属性继承自 groupBy 查询 / timeseries 查询,因此这两类查询的既有文档对该扩展同样适用。完整字段定义如下:
| property | description | required? |
|---|---|---|
| queryType | 固定为字符串"movingAverage",这是 Druid 判断如何解析查询的第一依据 | 是 |
| dataSource | 定义待查询数据源的字符串或对象,类似关系数据库中的表,见 DataSource | 是 |
| dimensions | DimensionSpec 的 JSON 列表(注意:该属性为可选) | 否 |
| limitSpec | 见 LimitSpec | 否 |
| having | 见 Having | 否 |
| granularity | 周期粒度(Period Granularity),见 Period Granularities | 是 |
| filter | 见 Filters | 否 |
| aggregations | 聚合定义,作为 Averager 的输入,见 Aggregations | 是 |
| postAggregations | 仅支持以聚合结果为输入,见 Post Aggregations | 否 |
| intervals | ISO-8601 时间区间的 JSON 对象,定义查询的时间范围 | 是 |
| context | 附加 JSON 对象,用于指定某些查询标志 | 否 |
| averagers | 定义移动平均函数,见下文 Averagers | 是 |
| postAveragers | 同时支持 averagers 与 aggregations 作为输入,语法与 postAggregations 一致(见 Post Aggregations) | 否 |
从源码 MovingAverageQuery.java 可以看到,@JsonTypeName("movingAverage")将该查询类型注册为movingAverage,构造器会逐一校验这些字段,包括:
- 必须指定 granularity(
Preconditions.checkNotNull(this.granularity, "Must specify a granularity")); - 输出名列不得重名:
verifyOutputNames会对 dimensions、aggregations、postAggregations 的输出名做去重校验,重复即抛Duplicate output name[...]异常; - 仅支持
PeriodGranularity:Runner 中若 granularity 不是 PeriodGranularity,直接抛Only PeriodGranulaity is supported for movingAverage queries(见 MovingAverageQueryRunner.java)。
此外,查询内部会把 averagers 包装成AveragerFactoryWrapper与原始聚合合并,构造一个用于应用 having / limit 的内部 groupBy 查询(groupByQueryForLimitSpec)。
无维度场景
当查询没有指定 dimensions时,Runner 会自动将内层查询优化为 timeseries 查询(源码 MovingAverageQueryRunner.java),因此单指标时间序列上的移动平均无需额外维度也能高效执行。
Averagers 详解
Averager 用于定义移动平均(窗口)函数,且并不局限于平均——它同样可以提供MAX()/MIN()等其他窗口函数。Averager 的输入是内层查询聚合出的字段(aggregations 的输出),输出是写入结果事件中的新列。
通用属性
所有 Averager 共有的属性如下:
| property | description | required? |
|---|---|---|
| type | Averager 类型,见下文 Averager 类型 | 是 |
| name | Averager 输出列名 | 是 |
| fieldName | 输入字段名(必须是某个聚合的名字) | 是 |
| buckets | 回看桶(时间周期)数量,包含当前桶,必须 > 0 | 是 |
| cycleSize | 周期大小,用于"星期几"这类周期内单桶计算,见 Cycle size(Day of Week);默认为 1 | 否 |
这些校验逻辑在源码 BaseAveragerFactory.java 的构造器中强制执行:
name、fieldName非空;cycleSize > 0,numBuckets > 0;cycleSize <= numBuckets;numBuckets必须能被cycleSize整除(numBuckets % cycleSize == 0),否则构造直接抛异常。
这意味着当你配置buckets: 28, cycleSize: 7时合法(28 % 7 == 0),而buckets: 10, cycleSize: 3会在查询解析期就被拒绝。
Averager 类型
所有类型在 AveragerFactory.java 的@JsonSubTypes中注册,分为标准类型与常量类型:
标准 averagers(Standard averagers),提供五种函数:
| 函数 | double 版本 | long 版本 |
|---|---|---|
| Mean(平均值) | doubleMean | longMean |
| MeanNoNulls(忽略空桶的平均值) | doubleMeanNoNulls | longMeanNoNulls |
| Sum(求和) | doubleSum | longSum |
| Max(最大值) | doubleMax | longMax |
| Min(最小值) | doubleMin | longMin |
对应的实现类位于 extensions-contrib/moving-average-query/src/main/java/org/apache/druid/query/movingaverage/averagers/,例如DoubleMeanAverager、LongSumAverager、DoubleMaxAverager等,每个 Averager 都配套一个*Factory负责 Jackson 反序列化与参数校验,并有对应的单元测试(如DoubleMeanAveragerTest、LongMeanNoNullAveragerTest)。
此外还有非常规类型constant(ConstantAveragerFactory),用于输出常量窗口值。
关于忽略空桶(Ignoring nulls):
使用MeanNoNulls类 averager 在查询区间起始于数据集开头时非常有用——此时首批记录会忽略缺失的桶,平均值不会被"人为拉低"。但反过来,如果数据集本身稀疏、存在空天,这些空桶同样会被忽略,平均值可能偏高。选择Mean还是MeanNoNulls需要根据数据稀疏度权衡。
示例用法:
{ "type" : "doubleMean", "name" : "<输出名>", "fieldName": "<输入聚合名>" }Cycle size(Day of Week)
cycleSize是可选参数,用于在每个周期内只取单一桶参与计算,而不是取全部桶。最典型的场景是:当桶粒度为"天"(period=P1D)、cycleSize=7时,就得到了"星期几"(Day of Week)的计算语义;同理可推广到"月内第几号""一天内第几小时"等场景。
官方文档给出的示例:
granularity: period=P1D(按天)buckets: 28cycleSize: 7
此时对于每一个输出记录,averager 只会对以下桶做计算:当前桶(#0)、#7、#14、#21(即 4 个"相同星期几"的桶)。而如果不指定cycleSize,则会使用全部 28 个桶计算。
已知限制(Limitations)
根据官方文档,目前该扩展存在以下限制:
- groupBy 属性缺失:
movingAverage不支持subtotalsSpec、virtualColumns; - timeseries 属性缺失:
movingAverage不支持descending; - 空值处理不兼容:
movingAverage不支持 SQL 兼容的空值处理(SQL-compatible null handling),因此设置druid.generic.useDefaultValueForNull=false会直接报错。
这条限制在源码中有明确印证:MovingAverageQuery构造器在开头就断言NullHandling.replaceWithDefault(),否则抛"movingAverage does not support druid.generic.useDefaultValueForNull=false"(见 MovingAverageQuery.java)。
实战示例
以下示例均基于 Druid tutorials 中提供的 Wikipedia 数据集。
基础示例:7 桶移动平均
计算 Wikipedia 编辑增量(delta)的 7 桶移动平均,桶粒度为 30 分钟:
{ "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:30:00.000Z", "event" : { "delta30Min" : 177882, "trailing30MinChanges" : 193890.0 } } ]注意每个输出的trailing30MinChanges等于当前桶及此前 6 个桶的delta30Min之和除以 7——这就是buckets: 7的滚动窗口语义,由MovingAverageIterable在 Broker 侧逐桶滑窗计算。
Post Averager 示例:当前值与移动平均的比率
在上一示例基础上,用postAveragers计算"当前周期值 / 移动平均值"的比率。postAveragers的语法与 postAggregations 完全相同,但输入同时支持聚合字段和 averager 输出字段:
{ "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-12T23:30:00.000Z", "event" : { "delta30Min" : 177882, "trailing30MinChanges" : 193890.0, "ratioTrailing30MinChanges" : 0.9174377224199288 } } ]可以看到ratioTrailing30MinChanges = delta30Min / trailing30MinChanges,该值在窗口内实现了"当前桶相对滚动均值"的归一化比较,可用于异常检测或环比分析。这一阶段由 PostAveragerAggregatorCalculator.java 实现。
Cycle size 示例:过去 3 小时中每个小时的第一个 10 分钟
计算"过去 3 小时内每个小时头 10 分钟的平均值",桶粒度为 10 分钟,每小时有 6 个桶,所以buckets: 18(3 小时 × 6 桶)、cycleSize: 6(每小时一个周期),每个输出只取当前桶及 6、12 号桶(即各小时的第 1 个桶):
{ "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 } ] }此处使用doubleMeanNoNulls而非doubleMean,目的是在窗口内某些 10 分钟桶无数据(空桶)时忽略它们,避免拉低平均值。
源码级延伸:理解内部实现
如果你希望深入理解该扩展,以下几个源码入口非常关键:
- 查询对象:MovingAverageQuery.java 定义了全部查询字段、输出名去重校验、空值处理断言,以及内部 groupBy(
groupByQueryForLimitSpec)与 having/limit 应用逻辑(applyLimit); - 执行引擎:MovingAverageQueryRunner.java 展示了完整的五步流水线:区间回退 → 内层 groupBy/timeseries →
RowBucketIterable分桶 →MovingAverageIterable滑窗 → postAveragers 与后处理; - 分桶与迭代:
RowBucketIterable/RowBucket负责把聚合行按 period 归入时间桶,MovingAverageIterable负责真正的滚动窗口计算,它们共同支撑buckets与cycleSize语义; - Averager 体系:AveragerFactory.java 定义了工厂接口与全部类型注册,BaseAveragerFactory.java 完成通用参数校验,各
*Averager类实现具体窗口计算; - 模块装配:
MovingAverageQueryModule负责把该查询类型与 Runner 注册进 Druid 的 Guice 依赖注入体系,MovingAverageQueryToolChest提供查询工具链支持; - 测试验证:MovingAverageQueryTest.java、
MovingAverageIterableTest、RowBucketIterableTest及averagers目录下各 Factory 测试覆盖了参数校验、窗口计算与查询结果行为,是理解边界条件的绝佳参考。
小结
Moving Average Query 通过"内层 groupBy/timeseries 聚合 + Broker 侧滑窗计算"的两阶段设计,为 Druid 补齐了移动平均与聚合窗口函数能力,同时借助区间自动回退与单次扫描避免了重复查询的性能损耗。在配置时请务必留意其限制:仅支持 PeriodGranularity、不支持subtotalsSpec/virtualColumns/descending,且必须保持druid.generic.useDefaultValueForNull为默认值。结合buckets与cycleSize的组合,你可以灵活实现"滚动 N 期平均""星期几对比""每小时首桶均值"等丰富的时间序列分析场景。
- 数据库
- OLAP
- 大数据
- 后端
【免费下载链接】druid
Apache Druid: a high performance real-time analytics database.
相关推荐
Apache Druid Moving Average 查询扩展:在 Druid 中原生实现移动平均与聚合窗口函数
Apache Druid Moving Average 查询扩展:在 Druid 中原生实现移动平均与聚合窗口函数 导读 本文基于 Apache Druid 开
数据库OLAP大数据后端Moving Average from Data Stream(数据流中的移动平均值):三种滑动窗口实现与复杂度深度剖析
Moving Average from Data Stream(数据流中的移动平均值):三种滑动窗口实现与复杂度深度剖析 本文基于本仓库 articles/mo
示例工程教程Apache Druid 的 Kerberos 认证扩展(druid-kerberos)配置与原理详解
Apache Druid 的 Kerberos 认证扩展(druid kerberos)配置与原理详解 本文以 druid kerberos 官方文档 http
数据库OLAP大数据后端
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考