ARTICLE DETAIL

资讯详情

深耕网站建设、视觉设计与SEO优化的一线实战洞察。

Apache Druid druid-stats 扩展指南:方差与标准差聚合的摄入预聚合与查询实战

Apache Druid druid-stats 扩展指南:方差与标准差聚合的摄入预聚合与查询实战 Apache Druid druid-stats 扩展指南方差与标准差聚合的摄入预聚合与查询实战【免费下载链接】druidApache Druid: a high performance real-time analytics database.项目地址: https://gitcode.com/gh_mirrors/druid7/druiddruid-stats 是 Apache Druid 的核心扩展core extension为 Druid 引入方差variance与标准差standard deviation两类统计聚合能力覆盖摄入阶段预聚合与查询阶段合并的全链路。读完本文你将掌握variance聚合器与varianceFold折叠聚合器的 JSON 配置、inputType/estimator参数的语义、stddev后聚合器的用法并能结合源码理解其数值稳定算法的底层实现直接在自己的 Timeseries / TopN / GroupBy 查询中落地。扩展概览与加载方式druid-stats位于仓库的 extensions-core/stats 模块其 Maven artifactId 为druid-stats见 pom.xml。它属于 Druid 打包在发行版中的核心扩展因此无需额外下载 jar只需在common.runtime.properties的druid.extensions.loadList中加入扩展名即可启用druid.extensions.loadList[druid-stats]完整的扩展加载说明见 including-extensions.md核心扩展清单见 extensions.md其中对 druid-stats 的定位描述为Statistics related module including variance and standard deviation。启用后扩展模块 DruidStatsModule.java 会通过 Guice 与 Jackson 完成三件事注册variance、varianceFold、stddev三个 JSON 子类型并为类型名为variance的复杂列注册序列化器VarianceSerde从而打通摄入、存储与查询三个阶段。Variance 聚合器与数值稳定算法variance聚合器计算一组数值的方差其算法与 Apache Hive 的GenericUDAFVariance完全一致源自 Chan、Golub 与 LeVeque 发表于The American Statistician1983, 37: 242–247的论文Algorithms for computing the sample variance: analysis and recommendations。该算法是**增量式incremental**的不要求一次性持有全部原始值而是每来一条数据就更新一组统计量因而天然适合 Druid 的流式/批式聚合场景。其核心合并公式为variance variance1 variance2 n/(m*(mn)) * pow(((m/n)*t1 - t2), 2)各符号含义variance表示sum[x-avg^2]即n 倍方差非样本方差本身每一步都会被更新n为 chunk1 的元素个数m为 chunk2 的元素个数t1为 chunk1 的元素之和t2为 chunk2 的元素之和。该算法的数值稳定性由 J.L. Barlow 在Error analysis of a pairwise summation algorithm to compute sample varianceNumer. Math, 58 (1991) pp. 583–590中证明可有效避免朴素两遍算法在数据量大、均值接近时出现的灾难性抵消问题。源码中的中间状态与增量更新源码 VarianceAggregatorCollector.java 用三个字段承载中间状态countlong元素个数sumdouble元素之和nvariancedoublesum[x-avg^2]即 n 倍方差。单条数据的增量更新add(float/long)在 VarianceAggregatorCollector.java 中实现count; sum v; if (count 1) { double t count * v - sum; nvariance (t * t) / ((double) count * (count - 1)); }两个 chunk 的合并逻辑combineValues则直接对应上文论文公式位于 VarianceAggregatorCollector.java。这正是 Druid 在 Broker、Historical 节点上做跨段segment合并时反复调用的路径也是整个扩展正确性的基石。中间状态大小与序列化由于聚合中间态只有count sum nvariance三个数值其固定大小为Longs.BYTES Doubles.BYTES Doubles.BYTES 24字节见 VarianceAggregatorCollector.java并以紧凑的putLong putDouble putDouble二进制格式落盘toByteBuffer因此相比保留全量原始值再计算内存与磁盘开销非常可控。摄入时预聚合方差Pre-aggregation at Ingestion使用该特性的前提是在索引摄入阶段就必须把variance聚合器写进摄入任务。摄入时的聚合器只能作用于数值类型的指标列若某输入行缺失该指标值会被视为取值0参与计算。摄入阶段variance聚合器的 JSON 结构如下{ type : variance, name : output_name, fieldName : metric_name, inputType : input_type, estimator : string }字段说明属性说明默认值type固定为variance无name聚合结果在输出中的列名无fieldName参与计算的指标列名无inputType期望的输入类型可取float、long、variancefloatestimator设为population时输出总体方差variance_pop否则为样本方差variance_samplenullinputType的默认值与三态分支在源码 VarianceAggregatorFactory.java 与 factorize 中体现float使用FloatVarianceAggregator按 float 列逐行累加long使用LongVarianceAggregator按 long 列逐行累加variance输入本身就是上一步聚合出的VarianceAggregatorCollector对象使用ObjectVarianceAggregator直接做 chunk 合并此时factorize通过makeObjectColumnSelector读取复杂列其余值会抛出IAE异常expected a float, long or variance。与Aggregator对象式并列还有基于ByteBuffer的 VarianceBufferAggregator.java它把count、sum、nvariance按偏移量0 / 8 / 16写入聚合缓冲区用于 Druid 的行式聚合执行路径两种实现共用同一套增量与合并公式。摄入段构建时VarianceSerde.java 的 extractor 会把输入行的原始值解析进VarianceAggregatorCollector支持直接接收 collector 对象也支持把多值维度逐值Float.parseFloat后逐个add最终以复杂列complex column形式存入 segment供后续查询期合并。查询期折叠varianceFold 聚合器如果摄入时已经预聚合了variance即以variance复杂列落盘那么查询时必须用variance类型的聚合器去合并这些中间态。文档给出的推荐写法有两种等价选择在查询中使用inputType为variance的variance聚合器或直接使用简化的varianceFold聚合器{ type : varianceFold, name : output_name, fieldName : metric_name, estimator : string }varianceFold在源码中是 VarianceFoldingAggregatorFactory.java它继承VarianceAggregatorFactory并强制把inputType固定为variance从声明层面保证只做中间态合并、不再解析原始数值。这也解释了为何摄入用variance、查询用varianceFold是一对标准组合——VarianceAggregatorFactory.getCombiningFactory()源码内部返回的正是new VarianceFoldingAggregatorFactory(name, name, estimator)。estimator 参数总体方差与样本方差estimator同时存在于variance、varianceFold与后文的stddev中用于选择方差口径取值含义计算公式源码 getVariancepopulation不区分大小写总体方差variance_popnvariance / count其他值或null默认样本方差variance_samplenvariance / (count - 1)判定逻辑只有一行estimator ! null estimator.equalsIgnoreCase(population)见 VarianceAggregatorCollector.java因此只要不显式传population一律按样本方差处理。建议在摄入与查询两侧保持一致的estimator设置。stddev 后聚合器由方差求标准差要从已聚合的方差结果获得标准差使用stddevpost-aggregator后聚合器{ type: stddev, name: output_name, fieldName: aggregator_name, estimator: string }其实现位于 StandardDeviationPostAggregator.javacompute逻辑即对fieldName指向的方差聚合结果开平方return Math.sqrt(((VarianceAggregatorCollector) combinedAggregators.get(fieldName)).getVariance(isVariancePop));注意fieldName必须指向查询aggregations中某个variance聚合器的name且该聚合器的estimator口径应与 post-aggregator 的estimator一致否则输出的是另一种口径的标准差。查询实战示例以下三类查询示例均来自官方文档可直接在启用了druid-stats的集群上运行数据源名为testing指标为index/index_var。Timeseries 查询按天聚合方差{ queryType: timeseries, dataSource: testing, granularity: day, aggregations: [ { type: variance, name: index_var, fieldName: index_var } ], intervals: [ 2016-03-01T00:00:00.000/2013-03-20T00:00:00.000 ] }TopN 查询按维度 TopN 并附带标准差{ queryType: topN, dataSource: testing, dimensions: [alias], threshold: 5, granularity: all, aggregations: [ { type: variance, name: index_var, fieldName: index } ], postAggregations: [ { type: stddev, name: index_stddev, fieldName: index_var } ], intervals: [ 2016-03-06T00:00:00/2016-03-06T23:59:59 ] }GroupBy 查询按维度分组输出方差与标准差{ queryType: groupBy, dataSource: testing, dimensions: [alias], granularity: all, aggregations: [ { type: variance, name: index_var, fieldName: index } ], postAggregations: [ { type: stddev, name: index_stddev, fieldName: index_var } ], intervals: [ 2016-03-06T00:00:00/2016-03-06T23:59:59 ] }三个示例演示了两种典型搭配直接对原始数值列fieldName: index用variance聚合出方差再通过postAggregations中的stddev派生出标准差index_stddev。若数据已在摄入期预聚合则把查询期聚合器换成varianceFold即可。边界行为与测试验证聚合器的边界行为在源码中有明确规定并有对应测试用例佐证空结果集当count 0时getVariance会抛出IllegalStateException(should not be empty holder)源码注释指出SQL 标准下应返回 null但 Druid 中不应出现该场景单元素count 1时方差按定义返回0d递增校验测试 VarianceAggregatorTest.java 验证了逐条喂入1.1 → 2.7 → 3.5 → 1.3后count/sum/nvariance的递推值以及总体/样本两种口径的结果testCombine还验证了两组 chunk 合并结果与一次性累加的等价性。仓库中另有 VarianceGroupByQueryTest.java、VarianceTimeseriesQueryTest.java、VarianceTopNQueryTest.java、VarianceSerdeTest.java 等端到端测试覆盖三类查询与序列化路径可作为理解或复现本扩展行为的参考。小结druid-stats以 24 字节的紧凑中间态、论文级的数值稳定合并公式为 Druid 补齐了方差与标准差这两类常用统计指标摄入阶段用varianceinputType取float/long预聚合查询阶段用varianceFold合并中间态再以stddev后聚合器输出标准差estimator: population可在总体与样本口径间切换。若你的业务需要按维度求波动率监控指标离散程度等统计型分析这套组合即可无缝嵌入 Druid 的 Timeseries、TopN 与 GroupBy 查询。【免费下载链接】druidApache Druid: a high performance real-time analytics database.项目地址: https://gitcode.com/gh_mirrors/druid7/druid创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表