如何在Spark DataFrame中使用fuzzywuzzy实现模糊字符串匹配筛选行
实现方案
你可以将fuzzywuzzy的匹配逻辑封装为Spark UDF(用户自定义函数),然后对目标列执行筛选即可,具体操作如下:
操作步骤
- 环境依赖确认
确保所有Spark worker节点都安装了所需依赖,可选安装python-Levenshtein加速匹配运算:
pip install fuzzywuzzy python-Levenshtein
- 封装模糊匹配UDF
首先定义匹配逻辑:只要输入的字符串和任意一个目标词的partial_ratio得分超过设定阈值,就判定为匹配成功。
from pyspark.sql.functions import udf from pyspark.sql.types import BooleanType from fuzzywuzzy import fuzz # 待匹配的目标词列表 TARGET_WORDS = ["apple", "orange"] # 匹配阈值 *可根据实际匹配效果调整,示例中appel和apple的得分为83,设置为80即可命中* MATCH_THRESHOLD = 80 def do_fuzzy_match(input_str): for target in TARGET_WORDS: # 统一转小写避免大小写影响匹配结果 if fuzz.partial_ratio(input_str.lower(), target.lower()) >= MATCH_THRESHOLD: return True return False # 注册为Spark UDF fuzzy_match_udf = udf(do_fuzzy_match, BooleanType())
- 筛选DataFrame
你可以选择DataFrame API或者Spark SQL两种方式执行筛选:- 方式1:DataFrame API 调用
# 假设你的原始DataFrame名为df result_df = df.filter(fuzzy_match_udf(df["fruit"])) # 输出验证结果 result_df.show()
- 方式2:Spark SQL 调用
如果习惯用SQL操作,先将UDF注册到SparkSession中即可:
# 注册UDF到Spark SQL环境 spark.udf.register("fuzzy_match", do_fuzzy_match, BooleanType()) # 执行SQL查询 result_df = spark.sql("SELECT * FROM your_table WHERE fuzzy_match(fruit)") # 输出验证结果 result_df.show()
优化提示:如果目标词列表很长,可以将TARGET_WORDS设置为Spark广播变量,减少worker节点之间的重复数据传输开销。
按上述配置执行后,会返回id为0、1、2的三行数据,包含拼写错误的appel记录。
内容的提问来源于stack exchange,提问作者John Constantine
相关产品推荐
相关产品推荐

