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

如何在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.22 10:35:09