ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

ES管道聚合实战:从buckets_path到环比累计

ES管道聚合实战:从buckets_path到环比累计 做数据分析的人应该都遇到过这种场景业务方跑过来问这个月环比上个月涨了多少今年累计做到了多少去掉单月波动后整体趋势是往上还是往下。拿ElasticSearch来说这些问题的答案都不是简单跑一个sum或avg就能出来的——它们是对聚合结果的再加工。我第一次面对这类需求时就是在Kibana Dev Tools里一层层套聚合直到同事提醒我这不就是Pipeline聚合吗当时那个原来如此的感觉我到现在还记得。ElasticSearch的聚合查询基础玩法是metric聚合和bucket聚合前者对字段算指标后者把文档分组。但如果你想在已经分好组的桶上面继续做计算——比如求各桶指标的最大值、最小值、平均值或者算桶与桶之间的环比差异、累计求和、移动平均——这两种聚合就不够用了。Pipeline聚合管道聚合就是专门干这个的它不直接对文档操作而是把另一个聚合的输出结果作为自己的输入再进行二次计算。用个不太严谨的类比普通聚合是把肉切成块码进盘子Pipeline聚合是淋汁、摆盘、加装饰的那一步。这篇文章面向两类人一是已经在用ES做聚合统计但对Pipeline聚合只听说过名字、没系统性跑过的人二是被类似月度环比累计值Top桶筛选需求纠缠过、想用一个查询而不是把结果拉回程序里慢慢算的人。我会用一份模拟的销售订单数据串起全程把buckets_path寻址规则、兄弟管道、父管道、脚本管道和常见坑一次性讲清楚。文章里所有的DSL你都可以直接复制到Kibana Dev Tools里跑前提是环境里已经启动了ES如果还没装好先启动elasticsearch和kibana很多奇怪的问题都是环境没就绪引起的误判。1. Pipeline聚合为什么不能缺席普通聚合的天花板1.1 普通聚合的一层能力边界先看一个具体的场景。假设我在sales索引里存了近一年的订单数据我用一个date_histogram按月聚合再加一个sum求每月销售额GET /sales/_search { size: 0, aggs: { monthly_sales: { date_histogram: { field: date, calendar_interval: month }, aggs: { total_amount: { sum: { field: amount } } } } } }返回结果里你得到的是12个桶每个桶有自己的total_amount值。到这里普通聚合的能力就触顶了它能告诉你每个月卖了多少钱但回答不了下面这些看上去很自然的问题这12个月里哪个月的销售额最高哪个最低全年的月平均销售额是多少每个月相比上个月增长了多少金额从年初到现在累计卖了多少去掉单月大促的波动整体趋势的平滑曲线长什么样这几个问题分别对应求极值、求均值、求差值、求累计、求移动平均。它们共同的特点都是需要跨桶计算或者对桶与桶的关系做运算。而metric聚合是单桶内的计算bucket聚合是分组两者都没有桶间运算的能力。1.2 父管道与兄弟管道两种输入输出模型Pipeline聚合能补上这一块但它也分两类理解起来要抓住一个核心它的输入从哪来输出往哪去。兄弟管道Sibling Pipeline Aggregation这种管道聚合和桶聚合放在同一层它把某个桶聚合产出的整组桶作为输入计算出一个统计值或者一个新的桶。输出的结果作为当前查询的一个独立聚合字段返回不会塞回原来的桶里。典型有min_bucket、max_bucket、avg_bucket、sum_bucket、stats_bucket、percentiles_bucket。父管道Parent Pipeline Aggregation这种管道聚合写在某个桶聚合的内部它针对桶聚合中的每个桶依次计算并把计算结果写到当前桶上成为该桶的一个新指标。典型有derivative、cumulative_sum、moving_function、bucket_script、bucket_selector、bucket_sort。我习惯用一个比喻来记桶聚合是流水线上分拣出来的一个个筐metric聚合在筐里给商品称重。兄弟管道把十几个筐放在一起比较哪筐最重、平均每筐多重父管道则沿着传送带一个筐一个筐地记录重量变化、计算到目前为止的累计重量。一个偏整体评估一个偏趋势演进。1.3 常用Pipeline聚合能力地图下面这张表是我在实际项目里用得比较高频的管道聚合可以当速查表类型聚合名称输入输出典型场景兄弟管道avg_bucket / sum_bucket某桶聚合内指定指标单个数值各月销售额的平均值、总和兄弟管道min_bucket / max_bucket某桶聚合内指定指标一个桶对象带key找出销量最高/最低的月份兄弟管道stats_bucket / percentiles_bucket某桶聚合内指定指标统计对象看所有桶的分布情况父管道derivative桶内指标追加到每个桶环比增量、变化速率父管道cumulative_sum桶内指标追加到每个桶累计销售额、累计活跃用户父管道moving_function桶内指标追加到每个桶移动平均、平滑趋势父管道bucket_script桶内多个指标追加到每个桶客单价、利润率等自定义计算父管道bucket_selector桶内指标过滤桶只保留满足条件的桶父管道bucket_sort桶内指标排序并截断桶按指标排序取TopN在开始写案例之前必须先搞定寻址规则。因为Pipeline聚合报错十有八九是buckets_path没写对。2. buckets_path寻址规则Pipeline聚合的路径坐标2.1 为什么buckets_path是命门Pipeline聚合自己不去索引里找数据它必须通过buckets_path参数告诉ES去哪个聚合的输出里拿数据。这个参数就是一个路径字符串有点像文件系统的目录定位。如果路径写错了ES会直接报一个no such aggregation [xxx]或者Unable to find a value for path [xxx]的错误而且很多时候报错信息里的聚合名你一眼看不出哪里有问题排查起来特别郁闷。搞清楚buckets_path的语法其实只需要记住三个符号、.、_count。2.2 三个关键语法元素首先是它用来分隔聚合层级的路径。ES里的这些聚合是有嵌套关系的buckets_path从你当前所在的位置出发用一层一层往下找。例如monthly_salestotal_amount表示先找到名为monthly_sales的桶聚合再找它内部的total_amount指标聚合。如果嵌套了多层桶聚合比如outer_bucketinner_bucketthe_metric路径就要一级一级写全。其次是.它用来定位多值指标里的某个具体值。sum、avg、min、max、value_count这些指标聚合是单值的直接写聚合名就能引用。但stats聚合、percentiles聚合返回的是多个值你不能直接拿stats这个名字去给管道用需要写成statsavg、statsmax、percentiles99.0这种形式用点号把具体要哪个值指出来。然后是特殊的_count和_key。_count代表桶的文档数量它不是某个子聚合而是每个桶自带的信息。比如想基于每个月的订单量而不是销售额做Pipeline计算路径就写monthly_sales_count。_key代表桶的键值比如date_histogram的桶key对应时间戳。2.3 父管道和兄弟管道的路径差异最容易搞混这一点是初学者踩坑重灾区。父管道写在桶聚合内部它的buckets_path是相对路径不需要包含当前父桶聚合的名字。比如你在monthly_sales的aggs里写derivative路径直接写buckets_path: total_amount就行不用写monthly_salestotal_amount。兄弟管道写在桶聚合外部它的buckets_path是全局路径必须从桶聚合的名字开始写。比如avg_bucket和monthly_sales平级路径必须写buckets_path: monthly_salestotal_amount。如果这里只写total_amountES在当前上下文找不到这个聚合直接就报错。2.4 三个典型的寻址错误我整理了自己调错时遇到的三种高频情况供你对照排查错误写法正确写法问题原因兄弟管道内写total_amountmonthly_salestotal_amount缺少父桶聚合名父管道内写monthly_salestotal_amounttotal_amount父管道在当前桶上下文中不能多写层级对stats聚合写price_statsprice_statsavg多值指标必须指定具体子值另外如果你的桶聚合嵌套了两层比如先按category分桶、再按date_histogram分桶那么管道路径也要跟着加一层。比如在category_bucketdate_bucketsales_amount这个三级结构上做兄弟管道路径就得写全。记住一个原则buckets_path是从当前聚合所在位置出发沿聚合树走到目标叶子节点的完整指引少一层、多一层都不行。3. 兄弟管道实操跨桶求极值与整体统计3.1 准备一份模拟数据为了把后面的例子串起来我建了一个sales索引字段如下order_id订单号、date下单时间、category品类、amount金额、quantity件数。数据量不用很大几百条就够关键是日期要覆盖一整年并且有分布差异方便看到有涨有跌。PUT /sales { mappings: { properties: { order_id: { type: keyword }, date: { type: date }, category: { type: keyword }, amount: { type: double }, quantity: { type: integer } } } }批量写入时可以写一个小脚本循环生成也可以手工构造几十条典型数据。我偷懒的方式是在程序里生成CSV后通过_bulk接口灌进去字段值保持自然分布就行不需要太精确。这一步的目的是让你能有一条可以反复跑DSL的本地数据链路。3.2 一次查询同时拿到月平均、最高月、最低月假设现在需要回答这一年里月平均销售额是多少哪个月卖得最高哪个月最低常规做法是先把12个月的数据拉出来在程序里排序取最大最小值。用兄弟管道聚合一个查询直接搞定GET /sales/_search { size: 0, aggs: { monthly_sales: { date_histogram: { field: date, calendar_interval: month }, aggs: { total_amount: { sum: { field: amount } } } }, avg_monthly: { avg_bucket: { buckets_path: monthly_salestotal_amount } }, max_month: { max_bucket: { buckets_path: monthly_salestotal_amount } }, min_month: { min_bucket: { buckets_path: monthly_salestotal_amount } } } }响应里的关键部分是这三个兄弟管道聚合的结果。avg_monthly返回一个数值表示这12个月的total_amount平均值。max_month和min_month的返回值不是简单的数值而是带桶标识的对象比如max_month: { value: 26800.0, keys: [2024-05-01T00:00:00.000Z] }这个keys字段很重要它告诉你最高值出现在哪个桶。在程序里做二次加工时我通常会同时读取value和keys这样不仅知道卖出26800元还能定位到这个月是哪个月。到这里有一个细节想提醒你sum_bucket和avg_bucket的输出是纯数值而min_bucket/max_bucket的输出是带着桶坐标的数值。解析响应的代码要区别处理否则可能拿到的不是你想用的那个结构。3.3 用stats_bucket看整体分布如果你不想分别写四个管道聚合可以直接用stats_bucket它会一次性返回所有桶指标的count、min、max、avg、summonthly_stats: { stats_bucket: { buckets_path: monthly_salestotal_amount } }对应的返回结构大概是monthly_stats: { count: 12, min: 5200.0, max: 26800.0, avg: 12350.0, sum: 148200.0 }这个聚合在做业务复盘时特别顺手一眼看出全年总盘子和各月波动范围不用拼好几个单独的管道。如果数据量更大、桶更多还可以用percentiles_bucket看分位数比如90%的月份销售额都低于某个值这对设定业绩目标很有参考价值。4. 父管道实操环比、累计和移动平均的建模过程4.1 derivative求环比第一个桶为什么没有值业务上最常见的环比需求是这个月比上个月多卖了多少。在ES里父管道derivative可以做这件事。它本质是一阶差分对时间序列里相邻两个桶的指标值做差。它的实现逻辑决定了第一个桶一定没有值——因为第一个桶前面没有上个月可以减。DSL写法如下GET /sales/_search { size: 0, aggs: { monthly_sales: { date_histogram: { field: date, calendar_interval: month }, aggs: { total_amount: { sum: { field: amount } }, sales_deriv: { derivative: { buckets_path: total_amount } } } } } }注意这里sales_deriv写在monthly_sales的aggs里所以它的buckets_path直接写total_amount就够了不需要带上monthly_sales。返回结果中2月桶的sales_deriv是1月到2月的差额3月桶的sales_deriv是2月到3月的差额而1月桶没有这个值。还有一个实用参数unit。如果你的时间桶是月但业务口径需要看日均变化可以在derivative里设置unit: dayES会把相邻两个月的总差额除以两个桶之间的天数得到平均每天比上月多多少的速率。这个参数在时间间隔不固定时尤其有用。4.2 cumulative_sum求累计gap_policy决定累计口径累计需求也很常见从年初到现在一共卖了多少。cumulative_sum做的事情是对桶内某个指标按桶顺序累加sales_cumsum: { cumulative_sum: { buckets_path: total_amount } }第3个桶的sales_cumsum就是第1、2、3个月的total_amount之和。这个聚合本身不复杂复杂的是它和缺失桶的关系。date_histogram默认min_doc_count是0也就是说没有订单的月份也会生成一个空桶total_amount为空。此时cumulative_sum要用gap_policy决定怎么处理。默认值是skip遇到缺失值直接跳过累计结果会像漏了一级台阶一样跨空设置gap_policy: insert_zeros时空桶被当作0参与累加累计曲线连续。选哪个取决于业务口径如果没有订单和订单金额为0在业务意义上是不同的那就要和生产对数过口径。4.3 moving_function移动平均window和model怎么选移动平均用来平滑单月波动看趋势方向。老版本的moving_avg在新版ES里已经废弃官方推荐用moving_function。最基础的写法是model: simple也就是简单平均sales_movavg: { moving_function: { buckets_path: total_amount, window: 3, model: simple } }window: 3 表示取当前桶以及前两个桶的均值。比如3月的sales_movavg是1、2、3三个月销售额的平均值。model可选simple、linear、ewma、holt、holt_winters。实际使用中simple最直观但它的滞后效应明显——如果上个月是大促月简单移动平均会把下个月的正常水平拉高。linear对更近的数据给了更高权重ewma适合短期波动大但希望快速响应的场景holt系列则适合有明显趋势或季节性的业务数据。还有一个shift参数容易被忽略。shift: 1表示窗口整体往后移一个桶例如window: 1、shift: 1窗口就只包含上一个桶的值。用这个特性可以很方便地取到上个月的销售额再配合bucket_script就能算同比增长率之类的指标了。4.4 把三种父管道叠加在同一个查询里实际项目中我很少只用一个管道更多是把derivative、cumulative_sum、moving_function全部铺在同一个date_histogram里一次拿到原始值、环比值、累计值、移动平均值四组数据GET /sales/_search { size: 0, aggs: { monthly_sales: { date_histogram: { field: date, calendar_interval: month }, aggs: { total_amount: { sum: { field: amount } }, sales_deriv: { derivative: { buckets_path: total_amount } }, sales_cumsum: { cumulative_sum: { buckets_path: total_amount, gap_policy: insert_zeros } }, sales_movavg: { moving_function: { buckets_path: total_amount, window: 3, model: simple } } } } } }返回的每个桶大致长这样{ key_as_string: 2024-03-01, key: 1709251200000, doc_count: 28, total_amount: { value: 15000.0 }, sales_deriv: { value: 3000.0 }, sales_cumsum: { value: 36000.0 }, sales_movavg: { value: 13333.333333333334 } }读这个返回值有一个技巧先看sales_deriv从哪个月开始有值有值前说明桶的起点再看sales_movavg的取值确认窗口是否符合预期最后用sales_cumsum倒推验证——把当年最后一个累计值减去倒数第二个月的累计值应该等于最后一个月自身的total_amount。我习惯把这个当作自检手段能快速发现有没有管道配置错位。5. 进阶组合bucket_script、bucket_selector与bucket_sort5.1 bucket_script用桶内多个指标做自定义计算运营经常问这个月的客单价是多少客单价 总销售额 / 总订单量。如果只存了销售总金额没有现成的客单价指标就得在ES里现算。bucket_script可以接收桶内多个指标做任意脚本计算。它和derivative这类固定逻辑的管道不同完全由你传入的脚本决定结果。GET /sales/_search { size: 0, aggs: { monthly_sales: { date_histogram: { field: date, calendar_interval: month }, aggs: { total_amount: { sum: { field: amount } }, total_quantity: { sum: { field: quantity } }, avg_unit_price: { bucket_script: { buckets_path: { amount: total_amount, quantity: total_quantity }, script: params.amount / params.quantity } } } } } }这里buckets_path是一个对象键名是你在脚本里引用的变量名值是具体的子聚合路径。脚本用宽松的params.xxx语法引用。PAINLESS脚本的要求比较严格写完最好先在Kibana的Dev Tools里跑一遍确认没有语法错误。5.2 bucket_selector把不达标的桶直接过滤掉从聚合结果里挑出需要的桶最直接的手段是bucket_selector。它的原理是对每个桶执行脚本返回true的桶保留false的桶丢弃。典型场景是只看销售额超过10000元的月份sales_filter: { bucket_selector: { buckets_path: { amount: total_amount }, script: params.amount 10000 } }要注意的是bucket_selector是在协调节点上对聚合结果做过滤的。它的好处是网络传输量变小但代价是计算发生在查询阶段复杂脚本会拖慢响应。我建议过滤条件里只用简单的比较运算复杂的打分逻辑放到程序侧做。还有一个常见的误解有人觉得用了bucket_selector之后查询的size参数就能控制桶的数量了。这里明确一下size控制的是文档数不是桶数控制桶的数量要用bucket_sort。5.3 bucket_sort管道阶段的排序与截断bucket_sort可以让你在管道层直接对桶做排序、分页和截断。比如按销售额降序取前三month_sort: { bucket_sort: { sort: [ { total_amount: { order: desc } } ], size: 3 } }size: 3 表示只保留前三个桶from参数还可以配合做分页比如取第4到第6名。为什么单独把它拎出来讲因为它是改变桶集合的管道对后续管道的影响很大。如果你还想在排序后再做别的计算要注意顺序bucket_sort一般放在管道链的末尾否则截断后的桶跑去喂后续管道数据就不完整了。5.4 组合案例筛选高销售额高客单价月份把上面三个串起来做一个真实点的报表需求找出销售额大于10000元且客单价大于50元的月份按销售额降序取前3个。GET /sales/_search { size: 0, aggs: { monthly_sales: { date_histogram: { field: date, calendar_interval: month }, aggs: { total_amount: { sum: { field: amount } }, total_quantity: { sum: { field: quantity } }, avg_unit_price: { bucket_script: { buckets_path: { amount: total_amount, quantity: total_quantity }, script: params.amount / params.quantity } }, sales_filter: { bucket_selector: { buckets_path: { amount: total_amount, price: avg_unit_price }, script: params.amount 10000 params.price 50 } }, month_sort: { bucket_sort: { sort: [ { total_amount: { order: desc } } ], size: 3 } } } } } }这个查询是一个桶聚合 三个metric聚合 三个管道聚合的组合。bucket_script先算出每月的客单价bucket_selector再用销售总额和客单价两个条件筛桶最后bucket_sort排序截断。整个过程一次HTTP请求完成数据不需要离开ES。运行之后你会发现返回的桶数量已经从一年12个被压到3个响应体也轻了很多。写这种组合查询时我的顺序习惯是先写date_histogram和两个sum跑通再加bucket_script确认avg_unit_price数值合理再加bucket_selector看过滤结果最后加bucket_sort截断。每一步单独验证出问题能立刻定位到是哪一层。6. 实战中容易踩的坑从报错到业务失真6.1 多值指标直接传给管道的报错这是新手报错的第一大来源。stats、percentiles、terms这类聚合返回的不是单个值而是多个值。直接拿price_stats当buckets_path传ES会报类似required a single-valued aggregation的错误。解决办法就是加后缀取单值buckets_path: price_statsavg取stats聚合的平均值buckets_path: price_statssum取总和buckets_path: price_percentiles99.0取P99分位另外terms聚合也是多值的Pipeline聚合不能直接以它作为数据源。如果需要按品类分组后的指标再做管道计算一般要先在terms桶内做单值metric聚合管道数据源指向那个metric。6.2 空桶和min_doc_count0的干扰date_histogram默认min_doc_count为0意味着没有数据的月份也会生成空桶。在一年12个月里只要有一两个月完全没有订单管道聚合的计算就会被空桶打乱。比如derivative在有空桶的情况下会把2月有数据、3月无数据、4月又有数据当成相邻两个有值桶做差算出来的环比其实是跨月比较非常容易误导业务。处理方式是根据口径分两类如果业务上允许把空月当0看设置min_doc_count: 0并配合gap_policy: insert_zeros让管道把空桶当0如果业务上要求严格对齐没有订单的月份不参与计算就把date_histogram的min_doc_count设为1或者使用gap_policy: skip。没有统一答案必须要对数据特征和生产口径。6.3 gap_policy选错导致平均值失真gap_policy直接影响业务数值最典型的例子在avg_bucket上。假设一年中有两个月没有订单avg_bucket默认skip会跳过空桶它的分母是10而不是12月均销售额就会被高估。如果改成insert_zeros空桶补0均值马上被拉低。哪个是对的取决于业务口径你计算的到底是在有销售的月份里的平均水平还是全自然年的月均水平。这个选择请务必和生产方确认不要在文档里自己猜。又一个相关坑是cumulative_sum。如果使用skip累计值会在跨过空桶时跳过该月但总额还是对得上如果使用insert_zeros中间月份会多出一个值为0的桶累计曲线更平滑。两者最终值一样但中间过程的形状不同。后续如果还接了moving_function或derivative这两个口径下的结果差异就会变大所以要全局统一。6.4 版本差异和环境问题Pipeline聚合在不同ES版本里的表现有差异。最明显的是moving_avg在新版已废弃老教程里用moving_avg的DSL直接粘到新版本会报错要改成moving_function。响应结构也有变化比如max_bucket的返回字段在不同大版本里可能是keys也可能带上value_as_string写程序解析时要注意兼容。环境方面我遇到过不少人在本地把ES和Kibana版本配错然后跑DSL报各种类型错误最后发现是Kibana连的ES版本和写法不匹配。如果是刚接触ES准备本地验证建议先把ES启动起来再起Kibana确认两个版本一致或兼容后再开始跑案例。6.5 性能边界Pipeline不是在ES里无限算的Pipeline聚合是在协调节点上完成的它不访问磁盘但会把外层桶聚合的所有桶都加载进内存再逐个计算。如果桶的数量到几万、上十万内存消耗会非常明显查询响应也会变慢。我处理过最夸张的场景是一个terms聚合产生了5万个桶再叠一个bucket_script协调节点直接把GC时间拉满。如果做全局统计或看整体的极值可以适当提高批量查询的size并用composite分页拉数据或者允许的话在应用侧对部分流程做后处理。Pipeline聚合适合中等规模、计算逻辑清晰的报表场景不要把所有复杂度都堆给它。结尾一点个人使用经验在Kibana Dev Tools里调Pipeline聚合我养成一个习惯先单独跑一遍桶聚合看清桶结构、确认指标名和返回类型再逐个添加管道。每次只加一个管道确认输出符合预期再加下一个。这样真正部署到生产环境的报错概率小很多。另外发现奇怪的报错先别急着改脚本把buckets_path从头拆开重新检查一遍。我至少有一半的排错时间都花在寻址路径上而路径错误的规律是高度可预测的父管道多写了外层、兄弟管道少写了外层、多值指标没加后缀就这三板斧。最后再提醒一句Pipeline聚合解决的是聚合结果的再加工问题如果某个业务场景需要处理几十万桶的超大规模聚合优先考虑把原始聚合结果导出或分页取回再在程序侧做二次计算有时候性能和可维护性都更好。
返回列表