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

PySpark:为含非统一重复ID的数据集进行ID标准化处理

统一用户标识的数据标准化方案

给定包含Date、Id2、id3、Name的数据集,其中:

  • Id2:部门可分配给任意用户的ID,同一用户可拥有多个
  • id3:用户专属ID,但会因遗忘变更
  • Name:存在重名情况

需要将同一用户的不同Id2和id3统一为同一标识,输入输出示例如下:

输入数据集

DateId2id3Name
2018-01-01001AAName1
2018-01-02001ABName1
2018-01-01001ACName1
2018-01-04003AAName1
2018-01-01004AAName1
2018-01-01001ADName2
2018-01-04002AEName3
2018-01-04005AGName1

期望输出数据集

DateId2id3Name
2018-01-01001AAName1
2018-01-02001AAName1
2018-01-01001AAName1
2018-01-04001AAName1
2018-01-01001AAName1
2018-01-01001ADName2
2018-01-04002AEName3
2018-01-04005AGName1

解决方案思路

核心是通过关联关系识别同一用户:

  1. 同一Name下,若Id2和id3存在交叉关联(如Id2=001关联id3=AA,id3=AA又关联Id2=003),则这些ID属于同一用户
  2. 为每个用户组选择一个基准标识(示例中选最早出现的Id2和对应的id3)
  3. 批量替换原数据中的对应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))

代码说明

  1. 使用**并查集(Union-Find)**算法自动合并所有属于同一用户的Id2和id3,构建关联网络
  2. 按Name分组处理,避免重名用户被误合并
  3. 为每个用户组选择最早出现的Id2和对应id3作为统一标识,与示例输出匹配
  4. 批量替换原数据中的ID,完成标准化处理

内容的提问来源于stack exchange,提问作者The Data Scientist

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.20 01:30:58