基于至少1个共同元素合并Hive/SparkSQL表记录
实现连通分量合并:合并共享key2元素的Hive/SparkSQL记录
这个需求本质是要解决连通分量合并的问题——把共享至少一个key2元素的记录归为同一组,再合并对应的key1集合和去重后的key2集合。我分别针对你提到的两种表结构,给出Hive/SparkSQL的具体实现方案:
方案一:输入表为int key1, int key2(推荐,处理更简便)
这种扁平化的结构更适合处理连通分量,我们可以用递归CTE来找出所有连通的key2节点,再分组合并结果。
测试数据准备
先创建测试表并插入示例数据:
CREATE TABLE IF NOT EXISTS test_key_pair ( key1 INT, key2 INT ); INSERT INTO test_key_pair VALUES (10, 1), (10, 3), (10, 5), (12, 1), (12, 2), (12, 4);
核心实现SQL
-- Hive需要先开启递归CTE支持 SET hive.recursiveCTE = true; SET hive.support.quoted.identifiers = none; WITH RECURSIVE cte AS ( -- 初始步骤:每个key2先和自己关联,同时关联同一key1下的其他key2 SELECT key2 AS start_key, key2 AS end_key FROM test_key_pair UNION -- 递归步骤:不断拓展连通的key2节点 SELECT c.start_key, t.key2 AS end_key FROM cte c JOIN test_key_pair t ON c.end_key = t.key2 JOIN test_key_pair t2 ON t.key1 = t2.key1 WHERE c.start_key != t.key2 ), -- 给每个连通分量分配唯一的组ID(用分量中最小的key2作为标识) key_groups AS ( SELECT end_key, MIN(start_key) AS group_id FROM cte GROUP BY end_key ) -- 按组ID合并key1集合和去重后的key2集合 SELECT COLLECT_SET(t.key1) AS merged_key1, COLLECT_SET(t.key2) AS merged_key2 FROM test_key_pair t JOIN key_groups kg ON t.key2 = kg.end_key GROUP BY kg.group_id;
逻辑说明
- 递归CTE
cte会找出所有连通的key2节点(比如1和3、5连通,1又和2、4连通,最终所有key2都属于同一分量); key_groups给每个key2节点分配所属分量的组ID,确保同一分量的节点组ID一致;- 最后按组ID分组,用
COLLECT_SET收集去重后的key1和key2集合,得到合并结果。
方案二:输入表为int key1, array<int> key2_list
这种结构需要先将数组展开为key1-key2的扁平结构,再复用方案一的逻辑:
测试数据准备
CREATE TABLE IF NOT EXISTS test_key_array ( key1 INT, key2_list ARRAY<INT> ); INSERT INTO test_key_array VALUES (10, array(1, 3, 5)), (12, array(1, 2, 4));
核心实现SQL
-- Hive开启递归CTE支持(SparkSQL无需额外配置) SET hive.recursiveCTE = true; WITH exploded AS ( -- 将数组展开为key1-key2的扁平结构 SELECT key1, explode(key2_list) AS key2 FROM test_key_array ), -- 以下逻辑和方案一完全一致 recursive cte AS ( SELECT key2 AS start_key, key2 AS end_key FROM exploded UNION SELECT c.start_key, t.key2 AS end_key FROM cte c JOIN exploded t ON c.end_key = t.key2 JOIN exploded t2 ON t.key1 = t2.key1 WHERE c.start_key != t.key2 ), key_groups AS ( SELECT end_key, MIN(start_key) AS group_id FROM cte GROUP BY end_key ), grouped AS ( SELECT kg.group_id, COLLECT_SET(e.key1) AS merged_key1, COLLECT_SET(e.key2) AS merged_key2 FROM exploded e JOIN key_groups kg ON e.key2 = kg.end_key GROUP BY kg.group_id ) -- 输出合并结果,若只需要key2集合可只选merged_key2 SELECT merged_key1, merged_key2 FROM grouped;
额外注意事项
- Hive版本要求:递归CTE需要Hive 2.1.0及以上版本,若版本较低,可通过多次自连接模拟递归(适合数据量较小的场景);
- SparkSQL支持:Spark 2.4.0及以上版本原生支持递归CTE,无需额外配置;
- 结果定制:如果不需要保留key1集合,直接去掉
merged_key1字段即可,只输出merged_key2。
内容的提问来源于stack exchange,提问作者kee
相关产品推荐
相关产品推荐

