尧图网络科技YAOTU DIGITAL 获取报价
获取报价
首页 / 资讯中心 / 文章详情

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

发布时间:2026/9/23 19:12:18

资讯中心
01
ARTICLE

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

Apache Druid Moving Average Query 扩展详解:在 Druid 中原生实现移动平均与窗口函数
数据库OLAP大数据后端【免费下载链接】druidApache Druid: a high performance real-time analytics database.项目地址https://gitcode.com/gh_mirrors/druid6/druid点击查看免费下载导读本文聚焦 Apache Druid 社区扩展druid-moving-average-queryMoving 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 天的平均编辑量。该扩展通过引入AveragerAverager窗口聚合器概念把标准 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 / 排序 / limitapplyLimit。安装与启用安装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 SpecMoving 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[...]异常仅支持PeriodGranularityRunner 中若 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 sizeDay of Week默认为 1否这些校验逻辑在源码 BaseAveragerFactory.java 的构造器中强制执行name、fieldName非空cycleSize 0numBuckets 0cycleSize numBucketsnumBuckets必须能被cycleSize整除numBuckets % cycleSize 0否则构造直接抛异常。这意味着当你配置buckets: 28, cycleSize: 7时合法28 % 7 0而buckets: 10, cycleSize: 3会在查询解析期就被拒绝。Averager 类型所有类型在 AveragerFactory.java 的JsonSubTypes中注册分为标准类型与常量类型标准 averagersStandard averagers提供五种函数函数double 版本long 版本Mean平均值doubleMeanlongMeanMeanNoNulls忽略空桶的平均值doubleMeanNoNullslongMeanNoNullsSum求和doubleSumlongSumMax最大值doubleMaxlongMaxMin最小值doubleMinlongMin对应的实现类位于 extensions-contrib/moving-average-query/src/main/java/org/apache/druid/query/movingaverage/averagers/例如DoubleMeanAverager、LongSumAverager、DoubleMaxAverager等每个 Averager 都配套一个*Factory负责 Jackson 反序列化与参数校验并有对应的单元测试如DoubleMeanAveragerTest、LongMeanNoNullAveragerTest。此外还有非常规类型constantConstantAveragerFactory用于输出常量窗口值。关于忽略空桶Ignoring nulls使用MeanNoNulls类 averager 在查询区间起始于数据集开头时非常有用——此时首批记录会忽略缺失的桶平均值不会被人为拉低。但反过来如果数据集本身稀疏、存在空天这些空桶同样会被忽略平均值可能偏高。选择Mean还是MeanNoNulls需要根据数据稀疏度权衡。示例用法{ type : doubleMean, name : 输出名, fieldName: 输入聚合名 }Cycle sizeDay of WeekcycleSize是可选参数用于在每个周期内只取单一桶参与计算而不是取全部桶。最典型的场景是当桶粒度为天periodP1D、cycleSize7时就得到了星期几Day of Week的计算语义同理可推广到月内第几号一天内第几小时等场景。官方文档给出的示例granularity: periodP1D按天buckets: 28cycleSize: 7此时对于每一个输出记录averager 只会对以下桶做计算当前桶#0、#7、#14、#21即 4 个相同星期几的桶。而如果不指定cycleSize则会使用全部 28 个桶计算。已知限制Limitations根据官方文档目前该扩展存在以下限制groupBy 属性缺失movingAverage不支持subtotalsSpec、virtualColumnstimeseries 属性缺失movingAverage不支持descending空值处理不兼容movingAverage不支持 SQL 兼容的空值处理SQL-compatible null handling因此设置druid.generic.useDefaultValueForNullfalse会直接报错。这条限制在源码中有明确印证MovingAverageQuery构造器在开头就断言NullHandling.replaceWithDefault()否则抛movingAverage does not support druid.generic.useDefaultValueForNullfalse见 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: 183 小时 × 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 定义了全部查询字段、输出名去重校验、空值处理断言以及内部 groupBygroupByQueryForLimitSpec与 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大数据后端【免费下载链接】druidApache Druid: a high performance real-time analytics database.项目地址https://gitcode.com/gh_mirrors/druid6/druid点击查看免费下载相关推荐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),仅供参考
02
RELATED NEWS

相关资讯

更多网站建设与数字化升级内容

03
WHY YAOTU

想打造同款高转化官网?

懂行业、懂生意,从建站到增长一站式陪跑

场景化定制

不做模板站,围绕你的业务场景量身设计,小众不撞款。

营销型架构

以转化目标组织内容与路径,让官网真正带来询盘。

全周期服务

设计、开发、运营、运维一体,上线只是开始。

免费获取你的建站方案

留下需求,专属顾问 24 小时内为你输出方案建议。