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

如何将二维数组元素逐一插入PySpark DataFrame的每行?

PySpark实现DataFrame新增对应数组列的方法

下面提供几种可行的实现方式,核心是保证数组元素与DataFrame行的顺序一一对应:

方法一:添加行索引后关联

通过给原DataFrame和数组构造的临时DataFrame添加索引,再进行关联匹配:

from pyspark.sql import SparkSession
from pyspark.sql.functions import row_number
from pyspark.sql.window import Window

# 初始化Spark会话与原DataFrame
spark = SparkSession.builder.appName("AddArrayColumn").getOrCreate()
df = spark.createDataFrame([(1, "foo"), (2, "bar")], ["id", "name"])

# 目标数组
arrays = [[1, 2, 3], [4, 5, 6]]

# 给原DataFrame添加行索引(按id排序保证顺序对应)
window = Window.orderBy("id")
df_with_index = df.withColumn("index", row_number().over(window))

# 将数组转为带索引的DataFrame
array_df = spark.createDataFrame(enumerate(arrays), ["index", "numbers"])

# 关联并移除索引列
result_df = df_with_index.join(array_df, on="index", how="inner").drop("index")
result_df.show()

方法二:使用自定义UDF

通过添加索引列,结合UDF根据索引提取对应数组元素:

from pyspark.sql import SparkSession
from pyspark.sql.functions import udf, row_number
from pyspark.sql.types import ArrayType, IntegerType
from pyspark.sql.window import Window

spark = SparkSession.builder.appName("AddArrayColumn").getOrCreate()
df = spark.createDataFrame([(1, "foo"), (2, "bar")], ["id", "name"])
arrays = [[1, 2, 3], [4, 5, 6]]

# 添加从0开始的索引列
window = Window.orderBy("id")
df_with_index = df.withColumn("index", row_number().over(window) - 1)

# 定义UDF根据索引获取对应数组
get_array_udf = udf(lambda idx: arrays[idx], ArrayType(IntegerType()))

# 生成目标列并移除索引
result_df = df_with_index.withColumn("numbers", get_array_udf("index")).drop("index")
result_df.show()

方法三:直接构造新DataFrame

如果原DataFrame数据量不大,可以直接收集行数据后合并数组再创建新DataFrame:

from pyspark.sql import SparkSession

spark = SparkSession.builder.appName("AddArrayColumn").getOrCreate()
df = spark.createDataFrame([(1, "foo"), (2, "bar")], ["id", "name"])
arrays = [[1, 2, 3], [4, 5, 6]]

# 收集原数据行并与数组合并
original_rows = df.collect()
new_rows = [(row.id, row.name, arr) for row, arr in zip(original_rows, arrays)]

# 创建目标DataFrame
result_df = spark.createDataFrame(new_rows, ["id", "name", "numbers"])
result_df.show()

注意事项

  • 确保数组arrays的长度与原DataFrame的行数完全一致,否则会出现数据匹配错位或丢失
  • 使用索引关联时,需指定明确的排序规则(如示例中的orderBy("id")),避免因行顺序不确定导致匹配错误

内容的提问来源于stack exchange,提问作者Youshikyou

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.03 18:35:23