如何在Kusto中基于累积和阈值实现数据分组?
实现Kusto中按累积和阈值分组的需求
假设我们有如下Kusto数据集(category字段唯一):
// assumptions:- category is unique let threshold=100; datatable(category:string,measure:int) [ 'cat1',10, 'cat2',20, 'cat3',30, 'cat4',40, 'cat5',45, 'cat6',45, 'cat7',49, 'cat8',50, 'cat9',50 ] | order by measure asc
需求是按measure升序计算,将累积和≤阈值的category划分为不同组,期望输出如下:
group1= {cat1,cat2,cat3,cat4} its cumulative sum is 100 group2= {cat5,cat6} its cumulative sum is 90 group3= {cat7,cat8} its cumulative sum is 99 group4= {cat9} its cumulative sum is 50
之前尝试的逻辑因Kusto不支持在row_cumsum中引用自身计算出的前一行值而失效:
// assumptions:- category is unique let threshold=100; datatable(category:string,measure:int) [ 'cat1',10, 'cat2',20, 'cat3',30, 'cat4',40, 'cat5',45, 'cat6',45, 'cat7',49, 'cat8',50, 'cat9',50 ] | order by measure asc | serialize cumsum=row_cumsum(measure,prev(cumsum)+measure>threshold)
解决方案:使用scan运算符实现状态跟踪
Kusto的scan运算符可以处理这种需要持续跟踪状态的累积计算,具体实现如下:
let threshold=100; datatable(category:string,measure:int) [ 'cat1',10, 'cat2',20, 'cat3',30, 'cat4',40, 'cat5',45, 'cat6',45, 'cat7',49, 'cat8',50, 'cat9',50 ] | order by measure asc | scan declare (group_id:long=1, current_sum:int=0) with ( step s: true => group_id = case(s.current_sum + measure > threshold, s.group_id + 1, s.group_id), current_sum = case(s.current_sum + measure > threshold, measure, s.current_sum + measure); ) | summarize category_list = strcat_array(make_list(category), ","), cumulative_sum = sum(measure) by group_id | extend result = strcat("group", group_id, "= {", category_list, "} its cumulative sum is ", cumulative_sum) | project result
代码说明
scan运算符:声明group_id(组编号,初始为1)和current_sum(当前组累积和,初始为0)两个状态变量。- 分组逻辑:每一行判断当前组累积和加上当前行的
measure是否超过阈值:- 超过则开启新组,
group_id加1,current_sum重置为当前行的measure - 未超过则继续累加,
group_id保持不变,current_sum加上当前行的measure
- 超过则开启新组,
- 聚合输出:按
group_id聚合,将每个组的category拼接成列表,计算组的总累积和,最后格式化为期望的输出字符串。
内容的提问来源于stack exchange,提问作者Dhiraj
相关产品推荐
相关产品推荐

