基于三列生成唯一标识符的实体解析技术方案咨询
基于多列关联生成唯一标识符的实现方案(适配千万级数据)
需求说明
需基于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
相关产品推荐
相关产品推荐

