如何基于另一PySpark DataFrame的列值拆分PySpark DataFrame?
PySpark高效拆分DataFrame:根据另一DataFrame的列值拆分
要避免对原DataFrame执行两次过滤带来的性能损耗,你可以先给原DataFrame添加一个标记列,标记每条记录是否匹配目标类别,再基于这个标记列拆分出两个DataFrame,这样只需要扫描原DataFrame一次。以下是具体实现方法:
准备工作
假设你的原DataFrame名为main_df,存储类别信息的DataFrame名为category_df,先统一列名方便后续操作:
from pyspark.sql import functions as F # 重命名类别DataFrame的列 category_df = category_df.withColumnRenamed("Product Categories", "category")
方法一:左连接+标记拆分
通过左连接匹配类别,标记是否存在匹配项,再拆分:
# 左连接并添加匹配标记 main_df_with_flag = main_df.join( category_df, main_df["Category"] == category_df["category"], how="left" ).withColumn( "is_in_category", F.when(F.col("category").isNotNull(), F.lit(True)).otherwise(F.lit(False)) ).drop("category") # 移除临时列 # 拆分出两个DataFrame df1 = main_df_with_flag.filter(F.col("is_in_category")).drop("is_in_category") df2 = main_df_with_flag.filter(~F.col("is_in_category")).drop("is_in_category")
方法二:exists子查询+标记拆分
利用exists子查询直接判断类别是否存在,再拆分:
# 添加是否匹配的标记列 main_df_with_flag = main_df.withColumn( "is_in_category", F.exists( category_df, lambda c: c["category"] == main_df["Category"] ) ) # 拆分DataFrame df1 = main_df_with_flag.filter(F.col("is_in_category")).drop("is_in_category") df2 = main_df_with_flag.filter(~F.col("is_in_category")).drop("is_in_category")
注意事项
如果你的数据存在格式不一致(比如示例中leisure goods.和leisure goods的标点差异),可以先清理数据再匹配:
# 清理Category列末尾的标点 main_df = main_df.withColumn( "Category", F.regexp_replace(F.col("Category"), r"\.$", "") )
这种方式的优势在于:仅对原DataFrame执行一次扫描和计算,相比两次独立过滤,能大幅减少大数据量下的运行时间。
内容的提问来源于stack exchange,提问作者Gerrit
相关产品推荐
相关产品推荐

