如何在PySpark中为DataFrame添加指定列表值的新列?
在Spark DataFrame中添加新列的简单方法
Spark和Pandas的核心区别在于Spark是分布式计算框架,无法像Pandas那样直接将Python列表赋值为新列(Pandas是单机内存中的有序数据)。要实现需求,需要通过行号关联的方式确保列表元素和DataFrame行一一对应,以下是两种常用的简单方法:
方法一:行号关联法(兼容所有Spark版本)
- 导入依赖模块
from pyspark.sql import SparkSession from pyspark.sql.functions import row_number, monotonically_increasing_id from pyspark.sql.window import Window
- 为原DataFrame添加行号列
# 假设你的原DataFrame名为df window_spec = Window.orderBy(monotonically_increasing_id()) df_with_row = df.withColumn("row_id", row_number().over(window_spec))
- 将Salary列表转为带行号的DataFrame
# 生成带行号的薪资数据 salary_data = [(i+1, s) for i, s in enumerate(Salary)] salary_df = spark.createDataFrame(salary_data, schema=["row_id", "Salary"])
- 关联两个DataFrame并移除行号列
result_df = df_with_row.join(salary_df, on="row_id").drop("row_id")
方法二:数组索引法(Spark 3.0+适用)
这种写法更简洁,但需要Spark版本在3.0及以上:
- 导入依赖模块
from pyspark.sql.functions import array, lit, element_at, row_number from pyspark.sql.window import Window
- 将Salary列表转为Spark数组,通过行号匹配对应元素
window_spec = Window.orderBy(monotonically_increasing_id()) # 将薪资列表转为Spark数组对象 salary_array = array(*[lit(sal) for sal in Salary]) # 添加薪资列,用行号取数组中对应位置的元素 result_df = df.withColumn("Salary", element_at(salary_array, row_number().over(window_spec)))
关键注意事项
- 必须通过
monotonically_increasing_id()或其他有序字段生成行号,确保原DataFrame行顺序和Salary列表顺序完全匹配,否则会出现数据对应错误。 - Spark本身不保证DataFrame的默认行顺序,绝对不能依赖输出显示的顺序来关联数据。
内容的提问来源于stack exchange,提问作者Anjum Hassan
相关产品推荐
相关产品推荐

