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

如何基于另一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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.11 16:53:16