PySpark:为含非统一重复ID的数据集进行ID标准化处理
统一用户标识的数据标准化方案
给定包含Date、Id2、id3、Name的数据集,其中:
Id2:部门可分配给任意用户的ID,同一用户可拥有多个id3:用户专属ID,但会因遗忘变更Name:存在重名情况
需要将同一用户的不同Id2和id3统一为同一标识,输入输出示例如下:
输入数据集
| Date | Id2 | id3 | Name |
|---|---|---|---|
| 2018-01-01 | 001 | AA | Name1 |
| 2018-01-02 | 001 | AB | Name1 |
| 2018-01-01 | 001 | AC | Name1 |
| 2018-01-04 | 003 | AA | Name1 |
| 2018-01-01 | 004 | AA | Name1 |
| 2018-01-01 | 001 | AD | Name2 |
| 2018-01-04 | 002 | AE | Name3 |
| 2018-01-04 | 005 | AG | Name1 |
期望输出数据集
| Date | Id2 | id3 | Name |
|---|---|---|---|
| 2018-01-01 | 001 | AA | Name1 |
| 2018-01-02 | 001 | AA | Name1 |
| 2018-01-01 | 001 | AA | Name1 |
| 2018-01-04 | 001 | AA | Name1 |
| 2018-01-01 | 001 | AA | Name1 |
| 2018-01-01 | 001 | AD | Name2 |
| 2018-01-04 | 002 | AE | Name3 |
| 2018-01-04 | 005 | AG | Name1 |
解决方案思路
核心是通过关联关系识别同一用户:
- 同一
Name下,若Id2和id3存在交叉关联(如Id2=001关联id3=AA,id3=AA又关联Id2=003),则这些ID属于同一用户 - 为每个用户组选择一个基准标识(示例中选最早出现的
Id2和对应的id3) - 批量替换原数据中的对应ID
代码实现(Python Pandas)
import pandas as pd from itertools import chain # 加载输入数据 df = pd.DataFrame([ ["2018-01-01", "001", "AA", "Name1"], ["2018-01-02", "001", "AB", "Name1"], ["2018-01-01", "001", "AC", "Name1"], ["2018-01-04", "003", "AA", "Name1"], ["2018-01-01", "004", "AA", "Name1"], ["2018-01-01", "001", "AD", "Name2"], ["2018-01-04", "002", "AE", "Name3"], ["2018-01-04", "005", "AG", "Name1"] ], columns=["Date", "Id2", "id3", "Name"]) # 按Name分组处理 def standardize_group(group): # 构建Id2和id3的关联图,用并查集合并关联ID def find_root(node, parent): while parent[node] != node: parent[node] = parent[parent[node]] node = parent[node] return node # 收集所有关联对并初始化父节点 pairs = list(zip(group["Id2"], group["id3"])) all_ids = list(set(chain.from_iterable(pairs))) parent = {id_: id_ for id_ in all_ids} # 合并关联的ID组 for id2, id3 in pairs: root2 = find_root(id2, parent) root3 = find_root(id3, parent) if root2 != root3: parent[root3] = root2 # 确定每个组的基准标识(最早出现的Id2及对应id3) id2_first_id3 = group.groupby("Id2")["id3"].first().to_dict() root_to_base = {} for id_ in all_ids: root = find_root(id_, parent) if root not in root_to_base: if root in id2_first_id3: base_id2, base_id3 = root, id2_first_id3[root] else: base_id2 = next(k for k, v in parent.items() if v == root and k in id2_first_id3) base_id3 = id2_first_id3[base_id2] root_to_base[root] = (base_id2, base_id3) # 替换当前组的Id2和id3 def replace_ids(row): root = find_root(row["Id2"], parent) return pd.Series(root_to_base[root]) group[["Id2", "id3"]] = group.apply(replace_ids, axis=1) return group # 应用分组处理并输出结果 result_df = df.groupby("Name", group_keys=False).apply(standardize_group) print(result_df.to_markdown(index=False))
代码说明
- 使用**并查集(Union-Find)**算法自动合并所有属于同一用户的
Id2和id3,构建关联网络 - 按
Name分组处理,避免重名用户被误合并 - 为每个用户组选择最早出现的
Id2和对应id3作为统一标识,与示例输出匹配 - 批量替换原数据中的ID,完成标准化处理
内容的提问来源于stack exchange,提问作者The Data Scientist
相关产品推荐
相关产品推荐

