首页/新闻资讯/正文详情

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

发布时间:2026/9/23 19:12:18 来源:云帆数科 栏目:资讯中心
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),仅供参考

相关推荐

CTF杂项工具箱:维吉尼亚与曼彻斯特解码实战指南
CTF杂项工具箱:维吉尼亚与曼彻斯特解码实战指南

简介:面向CTF爱好者与安全研究者的专用工具箱,集合编码解码、进制转换、维吉尼亚密码暴力破解、USB流量识别、曼彻斯编码分析、CRC32爆破、ZIP伪加密破解等高频功能,专门解决比赛里常见加密与隐写题目的快速解题需求。压缩包共92个文件&#… · 2026/9/23 19:12:18

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

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

BPSK匹配滤波实战:根升余弦成形与匹配滤波联合设计
BPSK匹配滤波实战:根升余弦成形与匹配滤波联合设计

简介:本资源是一份面向通信工程专业本科生与数字信号处理初学者的MATLAB仿真实验包,聚焦BPSK调制系统中匹配滤波与根升余弦脉冲成形的核心原理验证。资源通过完整闭环仿真,解决数字通信接收端如何在加性高斯白噪声环境下提升信噪比、抑制码间… · 2026/9/23 19:12:04

C#进销存系统源码:WinForms+ADO.NET可直接部署运行
C#进销存系统源码:WinForms+ADO.NET可直接部署运行

简介:这是一套基于C#开发的完整进销存管理系统源码,面向.NET初学者与中小型企业管理软件开发者,解决企业采购、销售、库存等核心业务环节的数字化管理需求。资源包含109个文件,以49个C#业务逻辑文件(如frmMain、frmJhG… · 2026/9/23 19:51:04

i.MX6嵌入式核心板开发与工业应用实战
i.MX6嵌入式核心板开发与工业应用实战

1. 项目概述:解密"dragonballz_e210-1"的硬件基因第一次看到"dragonballz_e210-1"这个型号时,我下意识摸了摸手边的开发板——这串字符背后往往藏着工程师才能懂的密码。经过拆解验证,确认这是基于Freescale(… · 2026/9/23 19:50:57

三星手机刷机实战:Odin3救砖与固件刷写全指南
三星手机刷机实战:Odin3救砖与固件刷写全指南

很多人把刷机想象成高风险手术,动一下就可能变砖。实际玩过三星设备的人会告诉你,只要手里有Odin3和一套匹配的官方固件,刷机这事儿反而比你在系统里瞎折腾要稳得多。我说这话是有底气的:手头这台卡在开机LOGO的S10,就… · 2026/9/23 19:50:50

电脑连接打印机速查手册:5种方案横向对比与避坑指南
电脑连接打印机速查手册:5种方案横向对比与避坑指南

电脑连接打印机速查手册:5种方案横向对比与避坑指南 刚把网上复制的驱动安装脚本扔进终端,结果报错代码一闪而过,系统托盘里打印机图标灰着不动?这种“复制粘贴即崩溃”的绝望感,我懂。很多人以为连打印机就是插上线、点两下鼠标的事,但在实际运维或开… · 2026/9/23 19:50:36

Kustomize 结构化数据内嵌 JSON/YAML 的定向替换与合并提案(22-03)深度解析
Kustomize 结构化数据内嵌 JSON/YAML 的定向替换与合并提案(22-03)深度解析

CLI开发工具云原生 【免费下载链接】kustomize Customization of kubernetes YAML configurations 项目地址: https://gitcode.com/gh_mirrors/ku/kustomize 点击查看 免费下载 本文档基于仓库 proposals/22-03-value-in-the-structured-data.md 展开,并… · 2026/9/23 19:50:23

AI Agent技能管理实战:从散装工具到可维护技能体系
AI Agent技能管理实战:从散装工具到可维护技能体系

写Agent技能管理这个话题,得从一次真实踩坑说起。三个月前,我给自己搭的自动化助手塞了十几个API调用,结果没过两周就乱成一锅粥——有的工具参数格式过时了,有的技能描述写得模糊让模型选错函数,还有几个技能互相冲突… · 2026/9/23 19:50:17

3招搞定手机怎么下载微信面试难题实战项目解析
3招搞定手机怎么下载微信面试难题实战项目解析

3招搞定手机怎么下载微信面试难题实战项目解析 面试被问“手机怎么下载微信”背后的原理,90%的人答不上来。别笑,这看似弱智的问题,实则是考察你对移动应用分发机制、安全校验及网络协议理解的试金石。我带过不少校招新人,他们背了八股文,却连一个A… · 2026/9/23 0:00:03

你有新短消息请注意查收:3个新手避坑指南搞定消息系统选型
你有新短消息请注意查收:3个新手避坑指南搞定消息系统选型

你有新短消息请注意查收:3个新手避坑指南搞定消息系统选型 面试被问“高并发下如何保证消息不丢失”,你张口就是“用Redis”,结果面试官追问“如果Redis宕机了怎么办”,你瞬间卡壳。这种场景太常见了,很多新手在背八股文时,只记住了技术名词… · 2026/9/23 0:00:29

Win7无线热点配置工具源码解析:解决API失效的3个实战技巧
Win7无线热点配置工具源码解析:解决API失效的3个实战技巧

Win7无线热点配置工具源码解析:解决API失效的3个实战技巧 Win7无线热点配置工具在Win10/11上跑不动?不是你的问题,是版本升级后 API 全变了。很多老项目里的 netsh wlan… · 2026/9/23 0:00:36

了解更多?预约专属演示

我们的顾问将为您一对一讲解产品与方案

企业微信二维码