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

如何用PySpark将表中数组列拆分映射为指定名称的列

使用PySpark将数组列拆分为指定名称的独立列

需求说明

原始数据表(最后一列为数组类型):

job_idtimestampitem_values
job12022-02-15T23:10:00.000+0000[0.2,3.4,13.2]
Three2022-02-15T23:20:00.000+0000[0.1,2.9,11.2]
job22022-02-15T23:30:00.000+0000[1.2,3.1,16.0]
job32022-02-15T23:40:00.000+0000[0.4,0.4,16.2]
job42022-02-15T23:50:00.000+0000[0.7,8.4,11.2]
job52022-02-15T24:00:00.000+0000[0.3,1.5,19.1]
job62022-02-15T24:10:00.000+0000[0.7,7.4,13.2]

对应的item名称列表:["item1", "item2", "item3"]

需要将数组列item_values中的元素分别映射为对应名称的独立列,得到目标表:

job_idtimestampitem1item2item3
job12022-02-15T23:10:00.000+00000.23.413.2
Three2022-02-15T23:20:00.000+00000.12.911.2
job22022-02-15T23:30:00.000+00001.23.116.0
job32022-02-15T23:40:00.000+00000.40.416.2
job42022-02-15T23:50:00.000+00000.78.411.2
job52022-02-15T24:00:00.000+00000.31.519.1
job62022-02-15T24:10:00.000+00000.77.413.2

实现方法

方法一:遍历生成新列

通过数组索引逐个提取元素,创建对应名称的列,逻辑直观易懂:

from pyspark.sql import SparkSession
from pyspark.sql.functions import col

# 初始化SparkSession
spark = SparkSession.builder.appName("ArrayToColumns").getOrCreate()

# 模拟原始数据集
data = [
    ("job1", "2022-02-15T23:10:00.000+0000", [0.2, 3.4, 13.2]),
    ("Three", "2022-02-15T23:20:00.000+0000", [0.1, 2.9, 11.2]),
    ("job2", "2022-02-15T23:30:00.000+0000", [1.2, 3.1, 16.0]),
    ("job3", "2022-02-15T23:40:00.000+0000", [0.4, 0.4, 16.2]),
    ("job4", "2022-02-15T23:50:00.000+0000", [0.7, 8.4, 11.2]),
    ("job5", "2022-02-15T24:00:00.000+0000", [0.3, 1.5, 19.1]),
    ("job6", "2022-02-15T24:10:00.000+0000", [0.7, 7.4, 13.2])
]

df = spark.createDataFrame(data, ["job_id", "timestamp", "item_values"])

# 定义item名称列表
item_names = ["item1", "item2", "item3"]

# 遍历列表,提取数组对应索引的元素作为新列
for idx, item_name in enumerate(item_names):
    df = df.withColumn(item_name, col("item_values")[idx])

# 移除原数组列(可选,根据实际需求决定)
df = df.drop("item_values")

# 查看结果
df.show()

方法二:使用selectExpr批量生成

通过字符串拼接生成SQL表达式,一次性完成列的提取,适合列数较多的场景:

# 基于上述已创建的df和item_names
select_expr = ["job_id", "timestamp"] + [f"item_values[{idx}] as {item_name}" for idx, item_name in enumerate(item_names)]
result_df = df.selectExpr(*select_expr)

# 查看结果
result_df.show()

注意事项

  • 确保数组列的元素长度与item名称列表的长度一致,避免出现索引越界或列缺失的问题。
  • 两种方法均可实现需求,可根据实际场景选择更合适的方式。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.27 08:32:53