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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.05 02:24:03