You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Kusto大数据集:替代scan算子的高效前向数据填充与查询方案寻求

大规模数据集下Kusto前向填充的高效优化方案

针对百万级属性、大量事件的数据集,原scan算子逐行迭代处理导致超时,以下提供两种更高效的实现方案,用于完成字符串/布尔值的前向填充需求。

测试数据集

let T = datatable(PropertyId:int, Tenant:string, Owner:string, NoisyNeighbour:bool , PropertyTitle:string, EventDate:datetime )
[
   1, "", "", bool(0),"",datetime(2022-08-01 00:00),
   1, "", "abc", bool(null),"",datetime(2022-08-01 01:00),
   1, "X","", bool(null),"Title updated",datetime(2022-08-02 00:00),
   1, "X", "cde",bool(null),"",datetime(2022-08-03 00:00),
   1, "A1", "",bool(null),"",datetime(2022-08-03 00:00),
   1, "A2", "",bool(null),"",datetime(2022-08-03 02:00),
   1, "A2", "def",bool(null),"",datetime(2022-08-03 03:00),
   1, "B", "", bool(null),"",datetime(2022-08-04 00:00),
   1, "C","", bool(1),"",datetime(2022-08-05 00:00),
   1, "D", "xyz",bool(null),"",datetime(2022-08-06 00:00),
];
T

期望填充效果(中文)

填充后需实现:

  • 租户字段:后续空值/空字符串被最近的非空租户值覆盖
  • 业主字段:后续空值/空字符串被最近的非空业主值覆盖
  • 是否为噪音租户字段:后续null值被最近的有效布尔值填充
  • 房产标题字段:后续空值/空字符串被最近的非空标题值覆盖

方案一:使用原生forward_fill()函数(推荐)

forward_fill()是Kusto引擎原生优化的矢量化操作,相比scan的逐行处理,性能提升显著,且代码简洁。

实现代码

let T = datatable(PropertyId:int, Tenant:string, Owner:string, NoisyNeighbour:bool , PropertyTitle:string, EventDate:datetime )
[
   1, "", "", bool(0),"",datetime(2022-08-01 00:00),
   1, "", "abc", bool(null),"",datetime(2022-08-01 01:00),
   1, "X","", bool(null),"Title updated",datetime(2022-08-02 00:00),
   1, "X", "cde",bool(null),"",datetime(2022-08-03 00:00),
   1, "A1", "",bool(null),"",datetime(2022-08-03 00:00),
   1, "A2", "",bool(null),"",datetime(2022-08-03 02:00),
   1, "A2", "def",bool(null),"",datetime(2022-08-03 03:00),
   1, "B", "", bool(null),"",datetime(2022-08-04 00:00),
   1, "C","", bool(1),"",datetime(2022-08-05 00:00),
   1, "D", "xyz",bool(null),"",datetime(2022-08-06 00:00),
];
T
// 将空字符串转为null,适配forward_fill的处理逻辑
| extend 
    Tenant = iff(Tenant == "", dynamic(null), Tenant),
    Owner = iff(Owner == "", dynamic(null), Owner),
    PropertyTitle = iff(PropertyTitle == "", dynamic(null), PropertyTitle)
// 按房产ID分区,按事件时间排序后执行前向填充
| partition by PropertyId (
    sort by EventDate asc
    | extend 
        Tenant = forward_fill(Tenant),
        Owner = forward_fill(Owner),
        NoisyNeighbour = forward_fill(NoisyNeighbour),
        PropertyTitle = forward_fill(PropertyTitle)
)
// 将null转回空字符串(保留原始空值的展示逻辑)
| extend 
    Tenant = coalesce(Tenant, ""),
    Owner = coalesce(Owner, ""),
    PropertyTitle = coalesce(PropertyTitle, "")

性能优势

  • 原生矢量化操作,避免逐行迭代的性能开销
  • partition by实现并行分区处理,充分利用集群计算资源,适合大规模数据场景

方案二:基于row_number()与max_by的兼容方案

若环境不支持forward_fill(),可采用此方案,同样基于矢量化计算,性能优于scan。

实现代码

let T = datatable(PropertyId:int, Tenant:string, Owner:string, NoisyNeighbour:bool , PropertyTitle:string, EventDate:datetime )
[
   1, "", "", bool(0),"",datetime(2022-08-01 00:00),
   1, "", "abc", bool(null),"",datetime(2022-08-01 01:00),
   1, "X","", bool(null),"Title updated",datetime(2022-08-02 00:00),
   1, "X", "cde",bool(null),"",datetime(2022-08-03 00:00),
   1, "A1", "",bool(null),"",datetime(2022-08-03 00:00),
   1, "A2", "",bool(null),"",datetime(2022-08-03 02:00),
   1, "A2", "def",bool(null),"",datetime(2022-08-03 03:00),
   1, "B", "", bool(null),"",datetime(2022-08-04 00:00),
   1, "C","", bool(1),"",datetime(2022-08-05 00:00),
   1, "D", "xyz",bool(null),"",datetime(2022-08-06 00:00),
];
T
// 按房产ID分区,生成行号
| extend rn = row_number() over (partition by PropertyId order by EventDate asc)
// 分区后按行号排序,利用max_by实现前向填充
| partition by PropertyId (
    sort by rn asc
    | extend 
        Tenant = max_by(Tenant, rn) over (partition by PropertyId, iff(Tenant != "", rn, 0) order by rn asc rows between unbounded preceding and current row),
        Owner = max_by(Owner, rn) over (partition by PropertyId, iff(Owner != "", rn, 0) order by rn asc rows between unbounded preceding and current row),
        NoisyNeighbour = max_by(NoisyNeighbour, rn) over (partition by PropertyId, iff(isnotnull(NoisyNeighbour), rn, 0) order by rn asc rows between unbounded preceding and current row),
        PropertyTitle = max_by(PropertyTitle, rn) over (partition by PropertyId, iff(PropertyTitle != "", rn, 0) order by rn asc rows between unbounded preceding and current row)
)
// 移除临时行号字段
| project-away rn

实现逻辑

通过row_number()标记每行的顺序,以非空值的行号为分区点,利用max_by在分区范围内取最新的有效值,实现前向填充效果。


内容的提问来源于stack exchange,提问作者Sahil Raj

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.22 10:27:46