PySpark中基于子串搜索拼接两个Dataframe的实现方案
PySpark 子串匹配关联两个DataFrame实现方案
PySpark的join不仅支持等值匹配,也支持非等值匹配(包括子串包含判断),可以直接通过自定义join条件实现需求,以下是两种不同数据量级的实现方案:
方案1:直接非等值Join(适合B表数据量较小的场景)
实现逻辑:直接用contains函数作为join条件,关联后分组聚合即可。
代码示例:
from pyspark.sql import functions as F # 假设df_a对应Dataframe A,df_b对应Dataframe B result_df = df_a.join( df_b, on=F.col("df_b.String").contains(F.col("df_a.String")), how="left" # 保留A中没有匹配到的记录 ).groupBy( "Current Accession", "String" ).agg( F.collect_list("Accession").alias("Mapped Accession") )
优缺点:代码简单易维护,但是非等值Join本质是先做全量笛卡尔积再过滤,若B表数据量较大(比如超10万条),性能会非常差。
方案2:N-Gram预过滤优化(适合A、B表均为大数据量的场景)
先通过N-Gram预匹配大幅减少候选集,再做子串校验,性能比直接Join提升10倍以上。
实现步骤:
- 第一步:预处理过滤,计算A表String字段的最小长度,过滤掉B表中String长度小于该值的记录,这部分记录不可能包含A表的任何子串
- 第二步:对A、B表的String字段,按A表最小子串长度切分N-Gram
- 第三步:基于N-Gram做等值Join,得到候选匹配集
- 第四步:在候选集上做子串包含校验,过滤出真实匹配的记录
- 第五步:分组聚合得到最终结果
代码示例:
from pyspark.sql import functions as F # 1. 计算A表字符串最小长度 min_len = df_a.select(F.min(F.length("String"))).head()[0] # 2. 过滤B表无效记录 df_b_filtered = df_b.filter(F.length("String") >= min_len) # 3. 自定义UDF生成指定长度的N-Gram列表 @F.udf("array<string>") def generate_ngram(s, n): if len(s) < n: return [] return [s[i:i+n] for i in range(len(s)-n+1)] # 4. 为A、B表生成N-Gram列 df_a_ngram = df_a.withColumn("ngram", generate_ngram("String", F.lit(min_len))) df_b_ngram = df_b_filtered.withColumn("ngram", generate_ngram("String", F.lit(min_len))) # 5. 炸开N-Gram列后做等值Join df_a_explode = df_a_ngram.withColumn("single_ngram", F.explode("ngram")) df_b_explode = df_b_ngram.withColumn("single_ngram", F.explode("ngram")) # 6. 候选集Join+去重后做子串校验 candidate_df = df_a_explode.join( df_b_explode, on="single_ngram", how="left" ).select( "Current Accession", df_a_ngram.String.alias("string_a"), "Accession", df_b_ngram.String.alias("string_b") ).dropDuplicates() # 去重避免同一对匹配重复计算 # 7. 子串校验+聚合 result_df = candidate_df.filter( F.col("string_b").contains(F.col("string_a")) | F.col("string_b").isNull() ).groupBy( "Current Accession", F.col("string_a").alias("String") ).agg( F.collect_list("Accession").alias("Mapped Accession") )
如果A表字符串长度差异很大,可以根据长度分段处理,进一步提升匹配效率。
内容的提问来源于stack exchange,提问作者ypriverol
相关产品推荐
相关产品推荐

