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
相关产品推荐
相关产品推荐

