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

如何在PySpark中基于不同列的不同条件关联两个数据集?

PySpark多列关联替换实现方案

需求是将df1中的rc1、rc2、rc3字段分别与df2的Key匹配,替换为对应的description,并调整部分字段的大小写,得到目标数据集。

原始数据集

df1

rc1rc2rc3resp
AB2AB1AB6jean
AB4AB3AB7shein
AB9AB5AB8patrick

df2

Keydescription
AB1Normal
AB4Expand
AB3small
AB6Big
AB8First
AB2Dock
AB7Missing
AB9Package
AB5Wrong

实现代码

from pyspark.sql import SparkSession
from pyspark.sql.functions import col, initcap, broadcast

# 初始化SparkSession
spark = SparkSession.builder.appName("MultiColumnMapping").getOrCreate()

# 创建df1 DataFrame
data_df1 = [
    ("AB2", "AB1", "AB6", "jean"),
    ("AB4", "AB3", "AB7", "shein"),
    ("AB9", "AB5", "AB8", "patrick")
]
df1 = spark.createDataFrame(data_df1, ["rc1", "rc2", "rc3", "resp"])

# 创建df2 DataFrame
data_df2 = [
    ("AB1", "Normal"),
    ("AB4", "Expand"),
    ("AB3", "small"),
    ("AB6", "Big"),
    ("AB8", "First"),
    ("AB2", "Dock"),
    ("AB7", "Missing"),
    ("AB9", "Package"),
    ("AB5", "Wrong")
]
df2 = spark.createDataFrame(data_df2, ["Key", "description"])

# 分步关联替换各列
# 替换rc1
df_temp = df1.join(broadcast(df2), df1.rc1 == df2.Key, "left") \
    .withColumnRenamed("description", "rc1") \
    .drop("Key", df1.rc1)

# 替换rc2并处理首字母大写
df_temp = df_temp.join(broadcast(df2), df_temp.rc2 == df2.Key, "left") \
    .withColumn("rc2", initcap(col("description"))) \
    .drop("Key", df_temp.rc2, "description")

# 替换rc3
df_temp = df_temp.join(broadcast(df2), df_temp.rc3 == df2.Key, "left") \
    .withColumnRenamed("description", "rc3") \
    .drop("Key", df_temp.rc3)

# 处理resp首字母大写并得到最终结果
final_df = df_temp.withColumn("resp", initcap(col("resp")))

# 展示结果
final_df.show()

代码说明

  1. 广播小表优化:使用broadcast(df2)将小表df2广播到所有节点,避免关联时的大量数据shuffle,提升性能。
  2. 分步关联替换:依次将df1的rc1、rc2、rc3与df2的Key关联,替换为对应的description,每次关联后清理冗余列。
  3. 大小写调整:通过initcap函数将rc2的"small"转为"Small"、resp的"patrick"转为"Patrick",匹配目标格式。
  4. 列名整理:关联后通过withColumnRenamed和drop调整列名,确保最终结果的列名与需求一致。

最终输出

+-------+-------+-------+-------+
|    rc1|    rc2|    rc3|   resp|
+-------+-------+-------+-------+
|   Dock| Normal|    Big|   Jean|
| Expand|  Small| Missing|  Shein|
|Package|  Wrong|  First|Patrick|
+-------+-------+-------+-------+

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.25 13:35:20