如何计算PySpark DataFrame列元素与指定列表的Jaccard相似度
PySpark计算固定列表与DataFrame数组列Jaccard相似度的解决方案
报错原因
你遇到的属性缺失报错,核心原因是df和df2是两个无关联的独立DataFrame,未做关联操作就直接跨表引用列进行计算,Spark无法解析两个无血缘关系表的行对齐规则,因此抛出解析错误。
方案1:直接使用字面量数组(单目标对比场景推荐)
不需要额外创建df2存储对比列表,直接将Python列表转换为Spark数组字面量参与计算即可,代码如下:
from pyspark.sql import functions as F # 原始数据 df = spark.createDataFrame( [ (1, ['112','333']), (2, ['112','223']) ], ["id", "minhash"] ) minhash_sig = ['112', '223'] # 将Python列表转换为Spark Array类型字面量 target_sig = F.array(*[F.lit(item) for item in minhash_sig]) # 计算Jaccard相似度 = 交集元素个数 / 并集元素个数 df = df.withColumn('jaccard_sim', F.size(F.array_intersect(target_sig, F.col('minhash'))) / F.size(F.array_union(target_sig, F.col('minhash'))) ) df.show()
运行后输出结果:
+---+----------+------------------+ | id| minhash| jaccard_sim| +---+----------+------------------+ | 1|[112, 333]|0.3333333333333333| | 2|[112, 223]| 1.0| +---+----------+------------------+
方案2:跨DataFrame关联计算(多目标对比场景适用)
如果确实需要用另一个DataFrame存储对比签名(比如需要同时对比多个目标签名的场景),先通过交叉连接将两个DataFrame的行合并,再进行相似度计算:
from pyspark.sql import functions as F from pyspark.sql import Row df = spark.createDataFrame( [ (1, ['112','333']), (2, ['112','223']) ], ["id", "minhash"] ) minhash_sig = ['112', '223'] df2 = spark.createDataFrame([Row(c1=minhash_sig)]) # 交叉连接将df2的行拼接到df的每一行 df_join = df.crossJoin(df2) # 计算相似度 df_res = df_join.withColumn('jaccard_sim', F.size(F.array_intersect(F.col('c1'), F.col('minhash'))) / F.size(F.array_union(F.col('c1'), F.col('minhash'))) ) df_res.show()
可选优化
如果数组中存在重复元素,可先对两个数组做去重处理后再计算,避免重复元素干扰相似度结果:
# 去重后计算示例 df = df.withColumn('jaccard_sim', F.size(F.array_intersect(F.array_distinct(target_sig), F.array_distinct(F.col('minhash')))) / F.size(F.array_union(F.array_distinct(target_sig), F.array_distinct(F.col('minhash')))) )
内容的提问来源于stack exchange,提问作者coderboi
相关产品推荐
相关产品推荐

