PySpark Levenshtein Join报错求助及Fuzzywuzzy使用方法咨询
问题分析与解决方案
一、先解决当前的TypeError报错
你碰到的TypeError: 'DataFrame' object is not callable是个典型的Spark语法错误——在Spark里不能用括号直接调用DataFrame取列(比如df7_ct_map("description")这种写法是错的),正确的列引用方式有两种:
- 使用点符号:
df7_ct_map.description - 使用
pyspark.sql.functions.col()函数:col("description")
另外还要注意你代码里的列名对应问题:你描述的是Data表有Disease description列,df7_ct_map表有Disease Indication列,但你的代码里写反了列的归属,得先修正这个对应关系。
修正后的代码示例:
from pyspark.sql.functions import levenshtein, col # 注意列名要和实际表结构对应,这里按你描述的表列名调整 joinedDF = df7_ct_map.join( Data, levenshtein(col("Disease Indication"), col("Disease description")) < 3 ) joinedDF.show(10)
⚠️ 重要性能提醒:直接用Levenshtein距离做全表join会触发笛卡尔积(15000*20000=3亿次文本比较),性能会极差甚至耗尽集群资源。建议先做预处理过滤:
- 比如提取文本关键词,用
like或分词交集过滤,只保留有共同关键词的行对再计算距离; - 或者将较小的表通过
broadcast()广播,减少数据传输开销。
二、使用Fuzzywuzzy实现需求
当然可以用Fuzzywuzzy来实现你的需求(优先完全匹配,其次选共同词数最多的匹配项),不过Fuzzywuzzy是Python本地库,需要结合Spark的UDF来适配分布式场景,具体步骤如下:
1. 安装依赖包
在你的运行环境或Spark集群所有节点上安装:
pip install fuzzywuzzy python-Levenshtein
2. 实现匹配逻辑的UDF
我们可以先广播较小的表减少数据传输,再通过UDF为每个疾病描述找到最优匹配项:
from pyspark.sql.functions import broadcast, udf, collect_list from pyspark.sql.types import StringType from fuzzywuzzy import fuzz # 广播较小的表(这里假设df7_ct_map规模稍大,也可以根据实际大小选择广播Data) broadcasted_ct_map = broadcast(df7_ct_map.select("Disease Indication").withColumnRenamed("Disease Indication", "indication")) # 定义UDF:从候选适应症中找出最优匹配项 def find_best_match(description, candidates): if not description or not candidates: return None # 优先找完全匹配 exact_matches = [c for c in candidates if c.strip().lower() == description.strip().lower()] if exact_matches: return exact_matches[0] # 无完全匹配时,选共同词数最多的(用token_set_ratio更适合词集合匹配) best_match = max(candidates, key=lambda x: fuzz.token_set_ratio(description, x)) return best_match match_udf = udf(find_best_match, StringType()) # 收集所有候选适应症,为每个疾病描述匹配最优项 result_df = Data.crossJoin(broadcasted_ct_map) \ .groupBy("Disease description") \ .agg(collect_list("indication").alias("candidates")) \ .withColumn("best_matched_indication", match_udf(col("Disease description"), col("candidates"))) \ .drop("candidates") # 关联回df7_ct_map的其他列(如果需要) final_result = result_df.join( df7_ct_map, result_df.best_matched_indication == df7_ct_map["Disease Indication"], "left" ) final_result.show(10)
补充说明:
fuzz.token_set_ratio会忽略词的顺序,计算两个文本的词集合相似度,完美契合你“共同词数最多”的需求;- 如果表数据量极大,收集所有候选适应症到每个节点会占用较多内存,这时可以先按关键词(比如疾病首字母、核心病种词)分组,只在组内找匹配项,缩小候选集规模;
- 也可以改用pandas UDF进一步优化性能,适合处理超大规模数据。
内容的提问来源于stack exchange,提问作者Lizou
相关产品推荐
相关产品推荐

