Spark传递相等映射场景下的ID重分配技术实现问询
基于传递性关联重新分配实体ID的实现方案
业务场景与需求
- 数据按
Comp_Key分组,每组内包含多个Comp_Name,每个Comp_Name关联若干Comp_Value,且ID对Comp_Name-Comp_Value组合唯一 - 需遵循传递相等规则重新分配ID:若A与B共享某Comp_Value,B与C共享某Comp_Value,则A、B、C视为同一实体,分配相同ID(可选用组内任意原有ID)
- 无关联的
Comp_Name保持原有ID不变
输入数据
| Comp_KEY | Comp_NAME | Comp_Value | ID |
|---|---|---|---|
| A | abc.com | 123 | 1 |
| A | abc.com | 1 | |
| A | MB inc | 356 | 2 |
| A | Collin.com | 123 | 3 |
| A | Collin.com | 356 | 3 |
| A | Boyd.com | 123 | 4 |
| A | Boyd.com | 790 | 4 |
| A | clad.com | 123 | 5 |
| A | clad.com | 2324 | 5 |
| A | Denton.com | 555 | 6 |
| A | Denton.com | 666 | 6 |
| A | Micro.com | 435 | 7 |
| A | Micro.com | 987 | 7 |
| A | XYZ.com | 435 | 8 |
| A | XYZ.com | 334 | 8 |
| A | Insta.com | 777 | 9 |
期望输出
| Comp_KEY | Comp_NAME | Comp_Value | ID |
|---|---|---|---|
| A | abc.com | 123 | 1 |
| A | abc.com | 1 | |
| A | MB inc | 356 | 1 |
| A | Collin.com | 123 | 1 |
| A | Collin.com | 356 | 1 |
| A | Boyd.com | 123 | 1 |
| A | Boyd.com | 790 | 1 |
| A | clad.com | 123 | 1 |
| A | clad.com | 2324 | 1 |
| A | Denton.com | 555 | 6 |
| A | Denton.com | 666 | 6 |
| A | Micro.com | 435 | 7 |
| A | Micro.com | 987 | 7 |
| A | XYZ.com | 435 | 7 |
| A | XYZ.com | 334 | 7 |
| A | Insta.com | 777 | 9 |
实现方案
方案1:SQL(基于递归CTE处理连通分量)
问题本质是识别连通分量:Comp_Value为连接点,共享同一Comp_Value的Comp_Name属于同一连通分量,再通过传递性合并所有关联分量。
步骤说明
- 筛选非空Comp_Value数据,建立Comp_Name与Comp_Value的关联(去重)
- 用递归CTE遍历所有连通的Comp_Name,为每个分量分配基准ID(比如分量内最小的ID)
- 将原始数据与基准ID映射关联,替换原有ID;空Comp_Value行直接沿用所属Comp_Name的基准ID
示例代码(PostgreSQL)
WITH name_value AS ( -- 整理非空Comp_Value的关联数据,去重 SELECT DISTINCT Comp_KEY, Comp_NAME, Comp_Value, ID FROM your_table WHERE Comp_Value IS NOT NULL AND Comp_Value <> '' ), connected_components AS ( -- 初始节点:每个Comp_Name的最小原始ID SELECT Comp_KEY, Comp_NAME, MIN(ID) AS base_id FROM name_value GROUP BY Comp_KEY, Comp_NAME UNION ALL -- 递归关联:通过共享Comp_Value合并连通的Comp_Name,统一基准ID SELECT nv.Comp_KEY, nv.Comp_NAME, cc.base_id FROM name_value nv JOIN name_value nv2 ON nv.Comp_KEY = nv2.Comp_KEY AND nv.Comp_Value = nv2.Comp_Value AND nv.Comp_NAME <> nv2.Comp_NAME JOIN connected_components cc ON nv2.Comp_KEY = cc.Comp_KEY AND nv2.Comp_NAME = cc.Comp_NAME ), final_base_ids AS ( -- 确定每个Comp_Name的最终基准ID(取分量内最小的ID) SELECT Comp_KEY, Comp_NAME, MIN(base_id) AS final_id FROM connected_components GROUP BY Comp_KEY, Comp_NAME ) -- 关联原始数据,替换ID SELECT t.Comp_KEY, t.Comp_NAME, t.Comp_Value, COALESCE(f.final_id, t.ID) AS ID FROM your_table t LEFT JOIN final_base_ids f ON t.Comp_KEY = f.Comp_KEY AND t.Comp_NAME = f.Comp_NAME ORDER BY t.Comp_KEY, t.Comp_NAME, t.Comp_Value;
方案2:Python Pandas(图论连通分量算法)
利用NetworkX库处理图结构:Comp_Name为节点,共享Comp_Value的节点间连边,识别连通分量后分配基准ID。
步骤说明
- 加载数据,筛选非空Comp_Value行,构建Comp_Name与Comp_Value的映射
- 生成图:同一Comp_Value对应的所有Comp_Name两两连边
- 识别连通分量,为每个分量分配该组内最小的原始ID
- 替换原始数据中的ID,空Comp_Value行直接沿用基准ID
示例代码
import pandas as pd import networkx as nx # 加载数据(可替换为直接读取数据库或Excel) df = pd.DataFrame([ ["A", "abc.com", "123", 1], ["A", "abc.com", "", 1], ["A", "MB inc", "356", 2], ["A", "Collin.com", "123", 3], ["A", "Collin.com", "356", 3], ["A", "Boyd.com", "123", 4], ["A", "Boyd.com", "790", 4], ["A", "clad.com", "123", 5], ["A", "clad.com", "2324", 5], ["A", "Denton.com", "555", 6], ["A", "Denton.com", "666", 6], ["A", "Micro.com", "435", 7], ["A", "Micro.com", "987", 7], ["A", "XYZ.com", "435", 8], ["A", "XYZ.com", "334", 8], ["A", "Insta.com", "777", 9] ], columns=["Comp_KEY", "Comp_NAME", "Comp_Value", "ID"]) # 处理非空Comp_Value的关联,生成边关系 non_empty = df[df['Comp_Value'].notna() & (df['Comp_Value'] != '')] edges = [] for _, group in non_empty.groupby(['Comp_KEY', 'Comp_Value']): names = group['Comp_NAME'].unique() if len(names) >= 2: # 生成组内Comp_Name的两两连接边 for i in range(len(names)): for j in range(i+1, len(names)): edges.append((names[i], names[j])) # 构建图并识别连通分量 G = nx.Graph() G.add_nodes_from(df['Comp_NAME'].unique()) G.add_edges_from(edges) # 为每个连通分量分配基准ID(取分量内最小的原始ID) component_id_map = {} for component in nx.connected_components(G): min_id = df[df['Comp_NAME'].isin(component)]['ID'].min() for name in component: component_id_map[name] = min_id # 替换原始数据的ID df['ID'] = df['Comp_NAME'].map(component_id_map).fillna(df['ID']) # 输出结果 print(df.to_markdown(index=False))
关键说明
- 空Comp_Value行:直接沿用所属Comp_Name的基准ID,因为属于同一实体
- 基准ID选择:可根据需求调整,比如取分量内最早/最小的ID,或生成全新ID
- 多Comp_KEY场景:上述方案均支持按Comp_KEY分组处理,不同分组的分量互不干扰
内容的提问来源于stack exchange,提问作者marc
相关产品推荐
相关产品推荐

