PySpark:如何合并两个行数相同的不同DataFrame/RDD
在PySpark中合并行数一致的两个DataFrame
因为PySpark的DataFrame是分布式存储的,没有像Pandas那样的行索引,所以不能直接通过列赋值的方式合并,需要借助临时关联键来实现,以下是两种可靠的方法:
方法1:使用monotonically_increasing_id()生成临时ID
这种方法会生成全局唯一的递增ID,适合行顺序严格对应的场景:
- 导入所需函数:
from pyspark.sql.functions import monotonically_increasing_id
- 给两个DataFrame添加临时ID列:
# 给user_ids添加临时ID user_ids_with_id = user_ids.withColumn("temp_id", monotonically_increasing_id()) # 给prediction添加临时ID prediction_with_id = prediction.withColumn("temp_id", monotonically_increasing_id())
- 通过临时ID关联并删除临时列:
combined_df = user_ids_with_id.join(prediction_with_id, on="temp_id", how="inner").drop("temp_id")
方法2:使用窗口函数生成连续行号
如果需要更严格的连续行号(避免分布式环境下ID不连续的情况),可以用row_number()窗口函数:
- 导入所需模块和函数:
from pyspark.sql.window import Window from pyspark.sql.functions import row_number, monotonically_increasing_id
- 定义窗口(用
monotonically_increasing_id()保证排序稳定性):
window_spec = Window.orderBy(monotonically_increasing_id())
- 给两个DataFrame添加行号列:
user_ids_with_id = user_ids.withColumn("temp_id", row_number().over(window_spec)) prediction_with_id = prediction.withColumn("temp_id", row_number().over(window_spec))
- 关联并清理临时列:
combined_df = user_ids_with_id.join(prediction_with_id, on="temp_id", how="inner").drop("temp_id")
注意事项
- 两种方法都要求两个DataFrame的行顺序完全对应,否则关联后的结果会出现列值不匹配的情况。
- 合并后的
combined_df包含user_id和prediction两列,后续可通过combined_df.rdd转为RDD。
内容的提问来源于stack exchange,提问作者Carlos Antunes
相关产品推荐
相关产品推荐

