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

Kusto百万行数据场景下Python插件计算连通分量性能优化咨询

Kusto百万级边表连通分量计算优化方案

问题说明

边列表数据存储在Kusto表中,需要划分连通分量,为每条边分配所属组号。当前实现通过Kusto Python插件运行DFS算法完成计算,百万行数据下脚本耗时普遍在2分钟以上,且受沙箱环境资源限制容易运行失败,需要适配超大规模行集的稳定优化方案。

业务规则:记录人员出借、借入资金的关系表中,存在直接/间接资金往来的人员属于同一连通分组。例如A借钱给B、C,C借钱给D,则A、B、C、D同属一个分组。

现有实现代码

// Payment表字段:出借人姓名LenderName、出借人手机号LenderPhone、借款人姓名BorrowerName、借款人手机号BorrowerPhone
Payment
| evaluate python(
    typeof(*,GroupNo:int),
    ```
    class Node:
        def __init__(self,name,phone):
            self.name=name
            self.phone=phone
        def __eq__(self,other):
            return self.name==other.name and self.phone==other.phone
        def __ne__(self,other):
            return not self.__eq__(other)
        def __hash__(self):
            return hash((self.name,self.phone))
        
    def get_adjacency_list(connections):
        adjacency_list=dict() # 键为Node对象,值为关联Node集合
        for index, row in connections.iterrows():
            node1=Node(row["LenderName"],row["LenderPhone"])
            node2=Node(row["BorrowerName"],row["BorrowerPhone"])
            
            if(node1 not in adjacency_list):
                adjacency_list[node1]=set()
            if(node2 not in adjacency_list):
                adjacency_list[node2]=set()
            
            adjacency_list[node1].add(node2)
            adjacency_list[node2].add(node1)
        return adjacency_list

    def run_dfs(node,group_no,visited,graph_list,groups):
        visited.add(node)
        groups[node]=group_no
        for next_node in graph_list[node]:
            if next_node not in visited:
                visited.add(next_node)
                run_dfs(next_node,group_no,visited,graph_list,groups)

                
    def get_groups(graph_list):
        groups=dict()
        visited=set()
        group_no=0
        for node in graph_list.keys():
            print(node)
            if node not in visited:
                group_no+=1
                run_dfs(node,group_no,visited,graph_list,groups)
        return groups

    graph_list=get_adjacency_list(df)
    groups=get_groups(graph_list)
    result=df
    result["GroupNo"]=df.apply(lambda row: groups.get(Node(row["LenderName"],row["LenderPhone"]),-1),axis=1)
    ```
)

优化策略

代码逻辑层面优化(仅改Python代码即可将耗时降低80%以上)

  • 替换核心算法为路径压缩+按秩合并的并查集:现有递归DFS实现存在两个致命问题,一是递归深度超过Python默认递归上限(1000)时会直接栈溢出崩溃,二是时间复杂度高于并查集。并查集是连通分量计算的最优算法,时间复杂度接近线性,无递归风险,实现简单,纯Python写的并查集处理百万边仅需数百毫秒。
  • 废弃自定义Node类:直接将(姓名,手机号)用特殊分隔符拼接为唯一字符串ID(比如f"{name}|{phone}"),省去自定义类实例化、哈希计算、属性访问的额外开销,比自定义对象操作快3倍以上,内存占用也更低。
  • 彻底弃用pandas低效逐行API:现有代码用的iterrows()、df.apply(axis=1)是pandas性能最差的逐行操作,百万行下仅遍历就要耗时数十秒。直接取DataFrame的底层numpy值数组做遍历,或者用itertuples()替代iterrows(),遍历速度可提升10~100倍;最终赋值GroupNo时用pandas的map()批量映射替代逐行apply,赋值速度可提升50倍以上。
  • 删除无意义的调试输出:代码中遍历节点时的print(node)语句在百万节点下会产生巨量IO开销,沙箱环境的stdout输出本身有严格的资源限制,多余的print会大幅拖慢运行速度,直接删除即可。
  • 省去邻接表构建环节:现有逻辑先构建完整邻接表再跑DFS,内存占用高且多了一次全量遍历。用并查集实现时,读取每条边直接对两个节点做union操作即可,不需要存完整邻接表,内存占用可降低60%以上。

Kusto沙箱适配优化

  • 进入Python前先做数据裁剪:在调用python插件前,先在Kusto层对边做去重,剔除完全重复的出借人-借款人关系,重复边不会改变连通分量结果,只会增加计算量;不需要参与计算的字段提前在Kusto层project掉,减少Python插件加载的数据量。
  • 控制单次计算的数据规模:如果单批次数据量超过200万行,先按业务维度(比如时间、地区)做分片,保证同一连通分量的所有边落在同一分片后,再分片调用Python插件计算,避免沙箱内存超限导致任务失败。
  • 不要引入额外第三方依赖:Kusto Python沙箱默认预装的库有限,不要尝试引入networkx等第三方图计算库,安装依赖会大幅增加脚本启动时间,还容易触发沙箱权限限制,纯Python实现的并查集性能够用。

架构层面优化(性能提升最大,可将耗时从分钟级降到秒级)

  • 优先用Kusto原生能力替代Python插件:Kusto自带原生图计算算子,支持make-graph构建图、graph-components直接计算连通分量,原生算子是编译执行的,性能比Python沙箱高1~2个数量级,不存在沙箱资源限制,千万级边也能在数秒内完成计算,稳定性远高于自定义Python脚本,是这类场景的首选方案。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.27 20:27:20