Kusto中定时检测已完成批次并执行批次指标计算的实现方法
Kusto批次完成状态检测实现方案
这个需求完全可以在Kusto中实现,不需要硬编码批次号,也不是必须用定时调度。
核心检测逻辑
判定逻辑不需要手动逐行遍历,直接用Kusto内置的next()开窗函数就能实现:按时间戳全局升序排序后,取每条记录下一行的批次号,和当前行批次号不一致的,当前行所属批次即为已完成批次。
针对你给出的测试数据,可直接运行如下查询验证:
let data = datatable(BatchNumber: int,Timestamp:datetime, Power1:int, Power2: int, Speed1: int, Speed2: int, Enabled1: bool, Enabled2: bool) [ 1, datetime(2022-02-18 10:00:00 AM), 100, 200, 50, 80, false, true, 1, datetime(2022-02-18 10:01:00 AM), 100, 200, 50, 80, true, true, 1, datetime(2022-02-18 10:02:00 AM), 100, 200, 50, 80, false, true, 1, datetime(2022-02-18 10:03:00 AM), 100, 200, 50, 80, true, true, 1, datetime(2022-02-18 10:04:00 AM), 100, 200, 50, 80, false, true, 1, datetime(2022-02-18 10:05:00 AM), 100, 200, 50, 80, true, true, 1, datetime(2022-02-18 10:06:00 AM), 100, 200, 50, 80, false, true, 1, datetime(2022-02-18 10:07:00 AM), 100, 200, 50, 80, true, true, 2, datetime(2022-02-18 10:08:00 AM), 100, 200, 50, 80, false, true, 2, datetime(2022-02-18 10:09:00 AM), 100, 200, 50, 80, true, true, 2, datetime(2022-02-18 10:10:00 AM), 100, 200, 50, 80, false, true, 2, datetime(2022-02-18 10:11:00 AM), 100, 200, 50, 80, true, true, 2, datetime(2022-02-18 10:12:00 AM), 100, 200, 50, 80, false, true, 2, datetime(2022-02-18 10:13:00 AM), 100, 200, 50, 80, true, true, 2, datetime(2022-02-18 10:14:00 AM), 100, 200, 50, 80, false, true, 2, datetime(2022-02-18 10:15:00 AM), 100, 200, 50, 80, true, true, 2, datetime(2022-02-18 10:15:00 AM), 100, 200, 50, 80, false, true ]; data | order by Timestamp asc // 取时间顺序下一条记录的批次号 | extend NextBatch = next(BatchNumber, 1) // 筛选出批次切换的边界点,对应已完成的前序批次 | where NextBatch != BatchNumber | project CompletedBatch = BatchNumber, BatchFinishTime = Timestamp // 关联原表聚合该批次全量指标,可按需替换聚合逻辑 | join kind=inner (data) on $left.CompletedBatch == $right.BatchNumber | summarize BatchFinishTime = max(BatchFinishTime), TotalPower1 = sum(Power1), TotalPower2 = sum(Power2), AvgSpeed1 = avg(Speed1) by CompletedBatch
运行上述查询只会返回批次1的计算结果,完成时间为2022-02-18 10:07:00 AM;等后续批次3的第一条记录接入后,再次运行会自动返回批次2的聚合结果,全程不需要硬编码批次号。
实现方式选择
你提到的5分钟定时调度是可行方案,但不是必须的,可根据业务实时性要求选择:
- 若对批次完成的检测延迟容忍度在分钟级,直接用定时调度即可,实现成本最低。注意每次调度时只查询上次调度时间之后的增量数据,同时写入结果表前去重,避免重复计算已输出的批次。
- 若需要秒级延迟的实时检测,直接使用Kusto原生的物化视图(materialized view)即可,不需要自己维护调度任务。物化视图会在后台随着新数据接入自动增量计算已完成批次,查询时直接读取视图结果即可,效率远高于定时扫描。
物化视图的基础配置参考如下:
// 先创建结果存储表 .create table CompletedBatches (BatchNumber:int, BatchFinishTime:datetime, TotalPower1:long, TotalPower2:long, AvgSpeed1:double) // 创建自动更新的物化视图,lookback根据业务最大数据乱序时长设置即可 .create materialized-view with (autoUpdateSchema=true, lookback=1h) CompletedBatchesView on table 你的实时数据表名 { 你的实时数据表名 | order by Timestamp asc | extend NextBatch = next(BatchNumber, 1) | where NextBatch != BatchNumber | project CompletedBatch = BatchNumber, BatchFinishTime = Timestamp | join kind=inner (你的实时数据表名) on $left.CompletedBatch == $right.BatchNumber | summarize BatchFinishTime = max(BatchFinishTime), TotalPower1 = sum(Power1), TotalPower2 = sum(Power2), AvgSpeed1 = avg(Speed1) by CompletedBatch }
后续直接查询CompletedBatchesView就能拿到所有已完成批次的结果,新批次完成后结果会自动追加。
注意:如果你的数据存在超过lookback窗口的晚到乱序,需要定期补算对应时间段的结果,避免漏算。
内容的提问来源于stack exchange,提问作者Noobie2021
相关产品推荐
相关产品推荐

