如何用PySpark将表中数组列拆分映射为指定名称的列
使用PySpark将数组列拆分为指定名称的独立列
需求说明
原始数据表(最后一列为数组类型):
| job_id | timestamp | item_values |
|---|---|---|
| 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] |
对应的item名称列表:["item1", "item2", "item3"]
需要将数组列item_values中的元素分别映射为对应名称的独立列,得到目标表:
| job_id | timestamp | item1 | item2 | item3 |
|---|---|---|---|---|
| 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 |
实现方法
方法一:遍历生成新列
通过数组索引逐个提取元素,创建对应名称的列,逻辑直观易懂:
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
相关产品推荐
相关产品推荐

