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

基于三列生成唯一标识符的实体解析技术方案咨询

基于多列关联生成唯一标识符的实现方案(适配千万级数据)

需求说明

需基于user_id、universal_id、session_id三列生成唯一标识符(expected_result为预期输出),最终构建universal_id:unique_id映射表。支持使用Snowflake、Databricks工具,可采用Python/PySpark实现。

核心规则

  • 当user_id缺失时,通过universal_id和session_id的关联关系生成唯一ID;
  • 若user_id不匹配但universal_id匹配,需视为不同唯一ID(仅当存在其他属性关联链路时才归为同一ID);
  • 按id列的入库顺序处理,新行若与已有行的任意一列匹配,需沿用已有唯一ID(本质是将所有共享任意属性的记录归为同一实体组)。

列间关系

  • user_id与universal_id为1:N或N:1(N:1时每个N需对应唯一ID,即共享同一universal_id的不同user_id需归为不同组,除非有其他属性关联);
  • user_id与session_id为1:N;
  • universal_id与session_id为1:N或N:1。

测试数据集

import pandas as pd

data = [
    [1, 1, 'apple', 'fiat', 1],
    [2, 1, 'pear', 'bmw', 1],
    [3, 2, 'bananna', 'citroen', 2],
    [4, 3, 'bananna', 'kia', 3],
    [5, 4, 'blueberry', 'peugeot', 4],
    [6, None, 'blueberry', 'peugeot', 4],
    [7, None, 'blueberry', 'yamaha', 4],
    [8, 5, 'plum', 'ford', 5],
    [9, None, 'watermelon', 'ford', 5],
    [10, None, 'raspberry', 'honda', 6],
    [11, None, 'raspberry', 'toyota', 6],
    [12, None, 'avocado', 'mercedes', 7],
    [13, None, 'cherry', 'mercedes', 7],
    [14, None, 'apricot', 'volkswagen', 2],
    [15, 2, 'apricot', 'volkswagen', 2],
    [16, 6, 'blueberry', 'audi', 8],
    [17, None, 'blackberry', 'bmw', 1],
    [18, 7, 'plum', 'porsche', 9]
]

df = pd.DataFrame(data, columns=['id', 'user_id', 'universal_id', 'session_id', 'expected_result'])

实现方案

1. Pandas版本(小数据测试用)

思路:按id顺序遍历,维护各属性到unique_id的映射,同时处理关联合并——当新记录的任一属性已存在映射时复用对应ID;若多个属性映射到不同ID,则合并映射关系。

def generate_unique_id(df):
    # 初始化属性到unique_id的映射字典
    user_map = {}
    universal_map = {}
    session_map = {}
    # 用于连通分量合并:存储每个ID对应的根ID
    parent = {}
    
    def find(u):
        # 路径压缩查找根ID
        while parent[u] != u:
            parent[u] = parent[parent[u]]
            u = parent[u]
        return u
    
    def union(u, v):
        # 合并两个连通分量,根取较小值(保证对应最早入库的记录)
        u_root = find(u)
        v_root = find(v)
        if u_root != v_root:
            if u_root < v_root:
                parent[v_root] = u_root
            else:
                parent[u_root] = v_root
    
    result = []
    current_max_id = 0
    
    # 按id顺序遍历处理每条记录
    for _, row in df.sort_values('id').iterrows():
        user_id = row['user_id']
        universal_id = row['universal_id']
        session_id = row['session_id']
        
        # 收集当前记录所有已关联的unique_id
        existing_ids = set()
        if user_id is not None and user_id in user_map:
            existing_ids.add(find(user_map[user_id]))
        if universal_id in universal_map:
            existing_ids.add(find(universal_map[universal_id]))
        if session_id in session_map:
            existing_ids.add(find(session_map[session_id]))
        
        if existing_ids:
            # 取连通分量的根ID作为当前记录的unique_id
            unique_id = min(existing_ids)
            # 合并所有关联ID到同一根
            for e_id in existing_ids:
                union(unique_id, e_id)
        else:
            # 生成新的唯一ID
            current_max_id += 1
            unique_id = current_max_id
            parent[unique_id] = unique_id
        
        # 更新映射字典
        if user_id is not None:
            user_map[user_id] = unique_id
        universal_map[universal_id] = unique_id
        session_map[session_id] = unique_id
        
        result.append(unique_id)
    
    # 修正所有ID为对应连通分量的根ID
    final_result = [find(id_) for id_ in result]
    df['generated_unique_id'] = final_result
    return df

# 测试执行
result_df = generate_unique_id(df)
print(result_df[['id', 'expected_result', 'generated_unique_id']])

2. PySpark版本(千万级大数据适配)

对于千万级数据,需使用分布式图计算框架处理连通分量,推荐使用GraphFrames(Databricks原生支持),步骤如下:

核心思路

将每条记录视为图的节点,同一属性(user_id/universal_id/session_id)关联的记录之间创建边,通过计算图的连通分量,将同一分量内的记录分配同一unique_id。

代码实现

from pyspark.sql import SparkSession
from pyspark.sql.functions import col, explode, collect_list, dense_rank
from pyspark.sql.window import Window
from graphframes import GraphFrame

# 初始化SparkSession
spark = SparkSession.builder.appName("UniqueIDGenerator").getOrCreate()

# 加载测试数据(实际可从Snowflake/Databricks表读取)
data = [
    [1, 1, 'apple', 'fiat', 1],
    [2, 1, 'pear', 'bmw', 1],
    [3, 2, 'bananna', 'citroen', 2],
    [4, 3, 'bananna', 'kia', 3],
    [5, 4, 'blueberry', 'peugeot', 4],
    [6, None, 'blueberry', 'peugeot', 4],
    [7, None, 'blueberry', 'yamaha', 4],
    [8, 5, 'plum', 'ford', 5],
    [9, None, 'watermelon', 'ford', 5],
    [10, None, 'raspberry', 'honda', 6],
    [11, None, 'raspberry', 'toyota', 6],
    [12, None, 'avocado', 'mercedes', 7],
    [13, None, 'cherry', 'mercedes', 7],
    [14, None, 'apricot', 'volkswagen', 2],
    [15, 2, 'apricot', 'volkswagen', 2],
    [16, 6, 'blueberry', 'audi', 8],
    [17, None, 'blackberry', 'bmw', 1],
    [18, 7, 'plum', 'porsche', 9]
]
df = spark.createDataFrame(data, ['id', 'user_id', 'universal_id', 'session_id', 'expected_result'])

# 1. 创建节点表:用记录id作为节点ID
vertices = df.select(col("id").alias("id")).distinct()

# 2. 创建边表:为每个属性关联的记录生成边
# 基于user_id生成边
user_edges = df.filter(col("user_id").isNotNull()) \
    .groupBy("user_id") \
    .agg(collect_list("id").alias("ids")) \
    .withColumn("src", explode(col("ids"))) \
    .withColumn("dst", explode(col("ids"))) \
    .filter(col("src") != col("dst")) \
    .select("src", "dst")

# 基于universal_id生成边
universal_edges = df.groupBy("universal_id") \
    .agg(collect_list("id").alias("ids")) \
    .withColumn("src", explode(col("ids"))) \
    .withColumn("dst", explode(col("ids"))) \
    .filter(col("src") != col("dst")) \
    .select("src", "dst")

# 基于session_id生成边
session_edges = df.groupBy("session_id") \
    .agg(collect_list("id").alias("ids")) \
    .withColumn("src", explode(col("ids"))) \
    .withColumn("dst", explode(col("ids"))) \
    .filter(col("src") != col("dst")) \
    .select("src", "dst")

# 合并所有边并去重
edges = user_edges.union(universal_edges).union(session_edges).distinct()

# 3. 构建图并计算连通分量
g = GraphFrame(vertices, edges)
connected_components = g.connectedComponents()

# 4. 关联原数据,将分量ID映射为连续的unique_id
window = Window.orderBy("min_id")
result_df = df.join(connected_components, on="id", how="left") \
    .groupBy("component") \
    .agg(min("id").alias("min_id")) \
    .join(df.join(connected_components, on="id", how="left"), on="component", how="right") \
    .withColumn("unique_id", dense_rank().over(window)) \
    .select("id", "user_id", "universal_id", "session_id", "expected_result", "unique_id")

# 构建universal_id到unique_id的映射表
universal_map = result_df.select("universal_id", "unique_id").distinct()

# 展示结果
result_df.show()
universal_map.show()

性能优化建议

  • 虚拟节点优化:为每个属性值创建虚拟节点(如user_1),将同一属性的所有记录连接到该虚拟节点,大幅减少边的数量;
  • 数仓原生计算:使用Snowflake的GRAPH_TABLE函数直接在数仓内计算连通分量,避免跨平台数据移动;
  • 缓存优化:在Databricks中对节点、边表启用缓存,加速图计算过程。

相关研究方向

  • 实体解析(Entity Resolution):该问题本质是实体解析中的共指消解,可参考基于属性关联的连通分量算法;
  • 分布式图计算:针对千万级数据,研究高效的连通分量计算实现(如分布式Union-Find、Louvain算法等);
  • 增量更新策略:针对增量入库的数据,研究如何在已有映射表基础上高效更新唯一ID,避免全量计算。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.22 13:52:07