如何在Spark DataFrame中按行匹配添加列而非生成笛卡尔积
原生PySpark实现按行顺序匹配关联
要实现你需要的按分组内行顺序一一匹配的效果,不能直接用仅按col1的join(会产生笛卡尔积),需要通过添加分组内的行号,再基于行号和col1进行关联,具体步骤如下:
步骤1:导入所需函数
from pyspark.sql import Window from pyspark.sql.functions import row_number, monotonically_increasing_id
步骤2:给两个DataFrame添加分组内的行号
通过窗口函数row_number(),按col1分组后给每行分配唯一行号。若需严格遵循数据的原始输入顺序,可使用monotonically_increasing_id()作为排序依据;如果业务上有明确的排序字段(比如示例中的col2/col3),也可以直接指定:
# 处理df1,添加行号rn window_df1 = Window.partitionBy("col1").orderBy(monotonically_increasing_id()) df1_rn = df1.withColumn("rn", row_number().over(window_df1)) # 处理df2,添加行号rn window_df2 = Window.partitionBy("col1").orderBy(monotonically_increasing_id()) df2_rn = df2.withColumn("rn", row_number().over(window_df2))
步骤3:基于col1和行号进行全外连接
通过full_outer连接,让两个DataFrame中同一分组、同一行号的行匹配,多余的行自动补null,最后整理列并排序:
df3 = df1_rn.join(df2_rn, on=["col1", "rn"], how="full_outer") \ .select("col1", "col2", "col3") \ .orderBy("col1", "rn")
验证结果
执行后df3的输出与你期望的一致:
+----+----+----+ |col1|col2|col3| +----+----+----+ | 1| a| k1| | 1| b| k2| | 1| c| k3| | 1|null| k4| +----+----+----+
关于转Pandas的说明
如果数据量较大,不建议转Pandas处理:Pandas是单机内存计算,会把所有数据拉到Driver节点,极易出现内存溢出;而原生PySpark方法是分布式处理,更适合大数据场景,性能和稳定性更优。
内容的提问来源于stack exchange,提问作者Prasanna Josium
相关产品推荐
相关产品推荐

