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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.15 04:19:49