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

Pyspark如何基于前缀匹配合并并展开两个DataFrame?

核心问题原因

你原方案不可行的根本原因是:UDF执行在Executor进程中,仅支持传入行级别的基础类型参数,DataFrame是Driver端维护的分布式数据集抽象,无法序列化后传递到Executor中执行。

最优解决方案(无需UDF,性能最优)

直接使用PySpark原生支持的非等值join即可实现前缀匹配需求,代码如下:

from pyspark.sql.functions import col

df_final = df_1.join(
    df_2,
    on=col("complete_code").startswith(col("partial_code")),
    how="inner"
).select("name", col("complete_code").alias("code"))

大数据量场景优化方案

如果你的数据集规模很大,非等值join性能不佳,且partial_code长度固定(比如示例中都是3位),可以提前为df2计算前缀字段,改用等值join提升性能:

# 前缀长度可根据实际场景调整
prefix_length = 3
df_2_with_prefix = df_2.withColumn("prefix", col("complete_code").substr(1, prefix_length))

df_final = df_1.join(
    df_2_with_prefix,
    on=df_1["partial_code"] == df_2_with_prefix["prefix"],
    how="inner"
).select("name", col("complete_code").alias("code"))

两种方案运行后都可以得到你预期的输出结果。


内容的提问来源于stack exchange,提问作者DawnZHANG

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.04 04:39:03