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

能否用ClickHouse实现高效Union-Find算法处理百亿级记录分组?

用ClickHouse处理百亿级Union-Find(不相交集合)分组问题

问题背景

需要对多文件中的百亿级记录执行Union-Find(不相交集合)分组操作,内存实现因数据量过大受限,希望借助ClickHouse的文件系统缓存能力完成任务。核心需求是基于item_id、from、to三列的图节点关系,生成包含id、group_id、item_id的分组结果,标记每条记录所属的不相交集合。

可复现示例

输入数据

item_id  from  to
0        101   102
1        102   103
2        104   105

预期结果

id  group_id  item_id
0   0         0
1   0         1
2   1         2

其中分组#0包含连通节点101->102->103,对应item_id 0和1;分组#1包含连通节点104->105,对应item_id 2。

ClickHouse解决方案思路

ClickHouse的列存引擎、磁盘友好的存储模型以及批量处理能力,适合处理超大规模的Union-Find场景。核心思路是通过迭代式的连通分量合并,结合ClickHouse的磁盘缓存避免内存溢出,具体步骤如下:

1. 数据导入与初始化

首先将多文件数据导入ClickHouse,使用MergeTree引擎存储边数据;同时初始化所有节点的父节点为自身(Union-Find的初始状态)。

-- 创建边数据表
CREATE TABLE edges (
    item_id UInt64,
    from_node UInt64,
    to_node UInt64
) ENGINE = MergeTree()
ORDER BY item_id;

-- 导入数据(示例:从本地文件导入)
INSERT INTO edges FORMAT CSVWithNames
FROM INFILE '/path/to/your/data.csv';

-- 创建节点表,初始化父节点为自身,秩为1
CREATE TABLE nodes (
    node UInt64,
    parent UInt64,
    rank UInt8
) ENGINE = ReplacingMergeTree()
ORDER BY node;

INSERT INTO nodes
SELECT DISTINCT node, node, 1
FROM (
    SELECT from_node AS node FROM edges
    UNION ALL
    SELECT to_node AS node FROM edges
);

2. 迭代合并连通分量

通过多次迭代执行合并操作,每次处理一批边的连通关系,更新节点的父节点与秩(按秩合并+路径压缩优化)。利用ReplacingMergeTree自动去重,保留每个节点的最新父节点信息。

-- 迭代合并脚本示例(可通过shell脚本循环执行,直到无更新)
CREATE TEMPORARY TABLE temp_merge AS
SELECT
    node,
    CASE
        WHEN root_from != root_to THEN least(root_from, root_to)
        ELSE parent
    END AS parent,
    CASE
        WHEN root_from != root_to THEN IF(rank_from >= rank_to, rank_from + 1, rank_to)
        ELSE rank
    END AS rank
FROM (
    SELECT
        n.node,
        n.parent,
        n.rank,
        -- 查找from节点的根
        (SELECT parent FROM nodes WHERE node = e.from_node) AS root_from,
        (SELECT rank FROM nodes WHERE node = e.from_node) AS rank_from,
        -- 查找to节点的根
        (SELECT parent FROM nodes WHERE node = e.to_node) AS root_to,
        (SELECT rank FROM nodes WHERE node = e.to_node) AS rank_to
    FROM nodes n
    LEFT JOIN edges e ON n.node IN (e.from_node, e.to_node)
)
WHERE root_from != root_to OR root_from IS NULL;

-- 合并更新节点表
INSERT INTO nodes SELECT * FROM temp_merge;

3. 生成最终分组结果

建立节点到group_id的映射,再关联原始边数据,得到每条item对应的分组标记。

-- 生成节点到group_id的映射(用根节点的稠密排名作为group_id)
CREATE TABLE node_group AS
SELECT
    node,
    dense_rank() OVER (ORDER BY root) AS group_id
FROM (
    SELECT
        node,
        -- 递归查找根节点(可通过多次迭代确保路径压缩完成)
        arrayReduce('first', arrayReverse(arrayAgg(parent))) AS root
    FROM nodes
    GROUP BY node
);

-- 关联原始边数据,生成最终结果
SELECT
    rowNumberInAllBlocks() AS id,
    ng.group_id,
    e.item_id
FROM edges e
JOIN node_group ng ON e.from_node = ng.node
ORDER BY id;

关键优化点

  • 引擎选择:使用ReplacingMergeTree处理节点的父节点更新,自动去重保留最新状态;
  • 分批处理:若数据量过大,将边数据拆分为多个分片,分批执行合并操作,降低单次查询的资源占用;
  • 分布式集群:借助ClickHouse分布式集群,将数据分片到多个节点并行处理,提升整体效率;
  • 路径压缩:定期执行路径压缩操作,将节点的父节点直接指向根节点,减少后续查找根节点的次数。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.06 13:01:22