news 2026/9/23 19:12:13

Apache Druid Moving Average Query 扩展详解:在 Druid 中原生实现移动平均与窗口函数

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
Apache Druid Moving Average Query 扩展详解:在 Druid 中原生实现移动平均与窗口函数
  • 数据库
  • OLAP
  • 大数据
  • 后端

【免费下载链接】druid

Apache Druid: a high performance real-time analytics database.

项目地址:https://gitcode.com/gh_mirrors/druid6/druid
点击查看免费下载

导读

本文聚焦 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 的输出进一步加工成带窗口语义的聚合结果。

该扩展带来的两个核心增强:

  1. 功能增强:为 Druid 查询引入窗口函数能力,支持MEANSUMMAXMIN等聚合窗口计算;
  2. 性能优化:通过一次 Segment 扫描 + Broker 侧窗口计算,消除多次查询拼接滚动窗口的重复扫描开销。

高层算法:两阶段流水线

从源码 MovingAverageQueryRunner.java 的类注释可以清晰看到整个执行流程。Moving Average Query 在内部封装了 groupBy 查询(无维度时退化为 timeseries 查询),以复用这两类成熟查询的能力,其执行分为两大阶段:

  1. 阶段一(内层查询):运行一个内层 groupBy(有维度)或 timeseries(无维度)查询,先计算出基础聚合值,例如"每天的编辑次数";
  2. 阶段二(Broker 侧窗口计算):在 Broker 上对聚合结果按时间桶(period bucket)滑动,计算 Averager,例如"每天编辑次数的 7 天移动平均"。

Runner 中更详细的步骤为:

  1. 取所有 Averager 中最大的buckets值,将查询区间起始时间向前回退buckets - 1个 period(见 MovingAverageQueryRunner.java),保证窗口有足够的回看数据;
  2. 按维度有无选择 groupBy 或 timeseries 内层查询;
  3. RowBucketIterable将各维度组合的行按 period 分桶(RowBucket);
  4. 交给MovingAverageIterable执行窗口计算、写入 Averager 结果列;
  5. 通过PostAveragerAggregatorCalculator应用 postAveragers;
  6. 过滤掉回看区间产生的、超出用户请求区间之外的行;
  7. 最后应用 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-processingdruid-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 查询,因此这两类查询的既有文档对该扩展同样适用。完整字段定义如下:

propertydescriptionrequired?
queryType固定为字符串"movingAverage",这是 Druid 判断如何解析查询的第一依据
dataSource定义待查询数据源的字符串或对象,类似关系数据库中的表,见 DataSource
dimensionsDimensionSpec 的 JSON 列表(注意:该属性为可选)
limitSpec见 LimitSpec
having见 Having
granularity周期粒度(Period Granularity),见 Period Granularities
filter见 Filters
aggregations聚合定义,作为 Averager 的输入,见 Aggregations
postAggregations仅支持以聚合结果为输入,见 Post Aggregations
intervalsISO-8601 时间区间的 JSON 对象,定义查询的时间范围
context附加 JSON 对象,用于指定某些查询标志
averagers定义移动平均函数,见下文 Averagers
postAveragers同时支持 averagers 与 aggregations 作为输入,语法与 postAggregations 一致(见 Post Aggregations)

从源码 MovingAverageQuery.java 可以看到,@JsonTypeName("movingAverage")将该查询类型注册为movingAverage,构造器会逐一校验这些字段,包括:

  • 必须指定 granularityPreconditions.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 共有的属性如下:

propertydescriptionrequired?
typeAverager 类型,见下文 Averager 类型
nameAverager 输出列名
fieldName输入字段名(必须是某个聚合的名字)
buckets回看桶(时间周期)数量,包含当前桶,必须 > 0
cycleSize周期大小,用于"星期几"这类周期内单桶计算,见 Cycle size(Day of Week);默认为 1

这些校验逻辑在源码 BaseAveragerFactory.java 的构造器中强制执行:

  • namefieldName非空;
  • cycleSize > 0numBuckets > 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(平均值)doubleMeanlongMean
MeanNoNulls(忽略空桶的平均值)doubleMeanNoNullslongMeanNoNulls
Sum(求和)doubleSumlongSum
Max(最大值)doubleMaxlongMax
Min(最小值)doubleMinlongMin

对应的实现类位于 extensions-contrib/moving-average-query/src/main/java/org/apache/druid/query/movingaverage/averagers/,例如DoubleMeanAveragerLongSumAveragerDoubleMaxAverager等,每个 Averager 都配套一个*Factory负责 Jackson 反序列化与参数校验,并有对应的单元测试(如DoubleMeanAveragerTestLongMeanNoNullAveragerTest)。

此外还有非常规类型constantConstantAveragerFactory),用于输出常量窗口值。

关于忽略空桶(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: 28
  • cycleSize: 7

此时对于每一个输出记录,averager 只会对以下桶做计算:当前桶(#0)、#7、#14、#21(即 4 个"相同星期几"的桶)。而如果不指定cycleSize,则会使用全部 28 个桶计算。

已知限制(Limitations)

根据官方文档,目前该扩展存在以下限制:

  • groupBy 属性缺失movingAverage不支持subtotalsSpecvirtualColumns
  • 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负责真正的滚动窗口计算,它们共同支撑bucketscycleSize语义;
  • Averager 体系:AveragerFactory.java 定义了工厂接口与全部类型注册,BaseAveragerFactory.java 完成通用参数校验,各*Averager类实现具体窗口计算;
  • 模块装配MovingAverageQueryModule负责把该查询类型与 Runner 注册进 Druid 的 Guice 依赖注入体系,MovingAverageQueryToolChest提供查询工具链支持;
  • 测试验证:MovingAverageQueryTest.java、MovingAverageIterableTestRowBucketIterableTestaveragers目录下各 Factory 测试覆盖了参数校验、窗口计算与查询结果行为,是理解边界条件的绝佳参考。

小结

Moving Average Query 通过"内层 groupBy/timeseries 聚合 + Broker 侧滑窗计算"的两阶段设计,为 Druid 补齐了移动平均与聚合窗口函数能力,同时借助区间自动回退与单次扫描避免了重复查询的性能损耗。在配置时请务必留意其限制:仅支持 PeriodGranularity、不支持subtotalsSpec/virtualColumns/descending,且必须保持druid.generic.useDefaultValueForNull为默认值。结合bucketscycleSize的组合,你可以灵活实现"滚动 N 期平均""星期几对比""每小时首桶均值"等丰富的时间序列分析场景。

  • 数据库
  • OLAP
  • 大数据
  • 后端

【免费下载链接】druid

Apache Druid: a high performance real-time analytics database.

项目地址:https://gitcode.com/gh_mirrors/druid6/druid
点击查看免费下载

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/23 19:12:15

后端性能优化:一文搞懂 irreversible 状态管理

后端性能优化:一文搞懂 irreversible 状态管理 很多开发者卡在“语法会、项目废”的泥潭里。代码能跑,但上线后高并发下响应时间飙升,甚至直接雪崩。这时候,你需要的不是背更多 API,而是一篇能直接指导你 一文搞懂 系统级不可逆操作(irreversible…

作者头像 李华
网站建设 2026/9/23 19:12:00

搞定四季教案源码:附完整示例与避坑指南

搞定四季教案源码:附完整示例与避坑指南 刚把网上扒来的“四季教案”Demo复制进IDE,点运行直接报错,心里那叫一个慌?别急,这种“代码跑不通、报错看不懂、改哪都不对”的情况,老鸟当年也经历过。很多教程只给结果,不给过程,导致你拿着“完整示例”却像拿着天书。…

作者头像 李华
网站建设 2026/9/23 19:11:35

传话机制手写实现:高频面试题背后的分布式一致性陷阱

传话机制手写实现:高频面试题背后的分布式一致性陷阱 面试被问原理答不上来,这大概是很多后端开发者最尴尬的时刻。特别是当面试官抛出“如何实现一个可靠的传话机制”时,很多人只能背出“TCP三次握手”,却对底层的丢包重传、幂等性处理一无所知。这不仅是高频面试题,更是检验你是否真正理解网络编程与并发控制的试…

作者头像 李华
网站建设 2026/9/23 19:11:05

百联集团实战项目揭秘:版本升级API变更下的底层逻辑与避坑指南

百联集团实战项目揭秘:版本升级API变更下的底层逻辑与避坑指南 版本升级后 API 全变了,这种崩溃感在接手【百联集团】相关的 实战项目 时尤为强烈。很多开发者面对百联集团这类大型零售企业的数字化系统重构,往往陷入“代码跑不通”的死循环,却忽略了底层协议映射的核心变化。别急着抱怨,我们先拆解这背后的…

作者头像 李华
网站建设 2026/9/23 19:11:00

泽洛斯避坑指南:版本升级API变更应对与面试高频考点解析

泽洛斯避坑指南:版本升级API变更应对与面试高频考点解析 版本升级后 API 全变了,代码跑不起来,报错信息满屏红,这是无数开发者在接手老项目或升级依赖时的噩梦。如果你正在为泽洛斯(Zeus)相关框架的接口变动而头疼,或者准备面试被问倒,这篇避坑指南就是为你准备的。我们不讲虚的,直接拆解版本差异、给…

作者头像 李华
网站建设 2026/9/23 19:10:12

左手螺旋定则与性能优化:3个细节搞定面试原理难题

左手螺旋定则与性能优化:3个细节搞定面试原理难题 面试被问电机控制底层原理,你卡壳了吗? 很多后端或嵌入式工程师在复盘 性能优化 方案时,发现瓶颈不在代码,而在对物理底层逻辑的误判。 今天用3个代码实例,讲透 左手螺旋定则 在工程中的映射,帮你把面试答得漂亮。 一、 定位差异:物理直觉 vs…

作者头像 李华