如何高效实现PySpark DataFrame与单词列表的匹配及相似度检测?
PySpark DataFrame 过滤:精确匹配 + 文本相似度检测
没问题,我来帮你搞定这个需求!咱们分步骤来,先搞定基础的精确匹配,再把SequenceMatcher的相似度检测加进去,最后整合两种逻辑,保证高效又满足你的要求。
1. 先实现基础的精确匹配
首先咱们先准备测试数据和目标单词列表,方便后续测试:
from pyspark.sql import SparkSession from pyspark.sql.functions import col, lower spark = SparkSession.builder.appName("TextFilter").getOrCreate() # 测试数据 data = [("hello foo",), ("bar world",), ("foobar test",), ("hello world",), ("baz foo",), ("f00",), ("b@r",)] df = spark.createDataFrame(data, ["text"]) # 目标单词列表 target_words = ["foo", "bar"]
精确匹配分两种常见场景,你可以根据需求选:
场景A:文本列完全等于目标列表中的某个单词
这种适合你要筛选那些文本内容就是目标单词的行:
exact_full_match_df = df.filter(col("text").isin(target_words)) exact_full_match_df.show()
场景B:文本列包含目标列表中的独立单词
这种适合筛选文本里出现目标单词的行(比如"hello foo"会被选中,但"foobar"不会,因为是连在一起的):
# 生成匹配独立单词的正则表达式,\b是单词边界 regex_pattern = r"\b(" + "|".join(target_words) + r")\b" exact_contains_match_df = df.filter(col("text").rlike(regex_pattern)) exact_contains_match_df.show()
2. 加入SequenceMatcher相似度检测
接下来咱们把difflib.SequenceMatcher集成进来,实现模糊匹配。因为PySpark是分布式计算,咱们需要把相似度逻辑包装成**UDF(用户自定义函数)**来使用。
步骤1:定义相似度判断函数
这里我写了两种逻辑,你可以根据需求切换:
from difflib import SequenceMatcher from functools import partial from pyspark.sql.functions import udf from pyspark.sql.types import BooleanType def has_similar_word(text, target_list, threshold=0.7): if not text: return False text_lower = text.lower() for word in target_list: word_lower = word.lower() # 逻辑1:整个文本和目标单词的相似度(适合短文本匹配) similarity = SequenceMatcher(None, text_lower, word_lower).ratio() if similarity >= threshold: return True # 逻辑2:文本中的任意单个单词和目标单词的相似度(适合长文本,只要有一个相似词就匹配) # for token in text_lower.split(): # sim = SequenceMatcher(None, token, word_lower).ratio() # if sim >= threshold: # return True return False # 用partial固定目标列表和阈值,方便注册UDF similar_udf = udf(partial(has_similar_word, target_list=target_words, threshold=0.7), BooleanType())
步骤2:用UDF过滤相似度达标的行
现在就可以用这个UDF筛选出和目标单词相似的行了:
similar_match_df = df.filter(similar_udf(col("text"))) similar_match_df.show()
比如测试数据里的"f00"(和foo相似度很高)、"b@r"(和bar相似度高)都会被筛选出来。
3. 整合精确匹配 + 相似度匹配
最后咱们把两种条件用OR结合,这样既保留精确匹配的行,也保留相似度达标的行:
# 整合条件:精确包含匹配 OR 相似度匹配 combined_df = df.filter( col("text").rlike(regex_pattern) | similar_udf(col("text")) ) combined_df.show()
小优化:提升大数据量下的性能
如果你的目标单词列表很大,建议把它广播出去,减少分布式任务之间的数据传输开销:
# 广播目标单词列表 broadcast_target = spark.sparkContext.broadcast(target_words) # 修改相似度函数,引用广播变量 def has_similar_word_broadcast(text, threshold=0.7): target_list = broadcast_target.value if not text: return False text_lower = text.lower() for word in target_list: word_lower = word.lower() similarity = SequenceMatcher(None, text_lower, word_lower).ratio() if similarity >= threshold: return True return False similar_udf_broadcast = udf(partial(has_similar_word_broadcast, threshold=0.7), BooleanType())
内容的提问来源于stack exchange,提问作者Kishintai
相关产品推荐
相关产品推荐

