如何在PySpark中基于不同列的不同条件关联两个数据集?
PySpark多列关联替换实现方案
需求是将df1中的rc1、rc2、rc3字段分别与df2的Key匹配,替换为对应的description,并调整部分字段的大小写,得到目标数据集。
原始数据集
df1
| rc1 | rc2 | rc3 | resp |
|---|---|---|---|
| AB2 | AB1 | AB6 | jean |
| AB4 | AB3 | AB7 | shein |
| AB9 | AB5 | AB8 | patrick |
df2
| Key | description |
|---|---|
| AB1 | Normal |
| AB4 | Expand |
| AB3 | small |
| AB6 | Big |
| AB8 | First |
| AB2 | Dock |
| AB7 | Missing |
| AB9 | Package |
| AB5 | Wrong |
实现代码
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()
代码说明
- 广播小表优化:使用
broadcast(df2)将小表df2广播到所有节点,避免关联时的大量数据shuffle,提升性能。 - 分步关联替换:依次将df1的rc1、rc2、rc3与df2的Key关联,替换为对应的description,每次关联后清理冗余列。
- 大小写调整:通过
initcap函数将rc2的"small"转为"Small"、resp的"patrick"转为"Patrick",匹配目标格式。 - 列名整理:关联后通过
withColumnRenamed和drop调整列名,确保最终结果的列名与需求一致。
最终输出
+-------+-------+-------+-------+ | rc1| rc2| rc3| resp| +-------+-------+-------+-------+ | Dock| Normal| Big| Jean| | Expand| Small| Missing| Shein| |Package| Wrong| First|Patrick| +-------+-------+-------+-------+
内容的提问来源于stack exchange,提问作者elokema
相关产品推荐
相关产品推荐

