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

如何将PySpark DataFrame中多列竖线分隔值拆分为对应行?

解决PySpark中拆分多列分隔值并对应行的问题

原始数据与需求

给定如下PySpark DataFrame:

import pyspark

arrayData = [
        ('1','temperature','21', 'Celsius'),
        ('2','humidity|temperature','88|21.8', 'Percent|Celsius'),
        ('3','temperature','21.2', 'Celsius'),
        ('4','temperature','19.9', 'Celsius'),
        ('5','humidity','85.5', 'Percent')]
df = spark.createDataFrame(data=arrayData, schema = ['id','types','values','labels'])
df.show()

输出结果:

+----+--------------------+---------+---------------+
|  id|               types|   values|         labels|
+----+--------------------+---------+---------------+
|   1|         temperature|       21|        Celsius|
|   2|humidity|temperature|  88|21.8|Percent|Celsius|
|   3|         temperature|     21.2|        Celsius|
|   4|         temperature|     19.9|        Celsius|
|   5|            humidity|     85.5|        Percent|
+----+--------------------+---------+---------------+

需要将types、values、labels列按|拆分,让同一位置的元素对应生成单独行,目标输出:

+----+--------------------+---------+---------------+
|  id|                type|    value|          label|
+----+--------------------+---------+---------------+
|   1|         temperature|       21|        Celsius|
|   2|         temperature|     21.8|        Celsius|
|   2|            humidity|       88|        Percent|
|   3|         temperature|     21.2|        Celsius|
|   4|         temperature|     19.9|        Celsius|
|   5|            humidity|     85.5|        Percent|
+----+--------------------+---------+---------------+

错误尝试分析

直接多次使用explode会生成笛卡尔积,导致大量多余行:

from pyspark.sql.functions import explode,col,split
df_2 = df.withColumn("id",col("id"))\
         .withColumn("type",explode(split("types", "\\|")))\
         .withColumn("value",explode(split("values", "\\|")))\
         .withColumn("label",explode(split("labels", "\\|")))
df_2.show()

输出会出现所有可能的组合,不符合需求。

正确实现方法

使用arrays_zip将拆分后的三个数组按位置配对成结构体数组,再explode展开,最后提取字段:

from pyspark.sql.functions import split, explode, arrays_zip, col

# 拆分各列为数组,打包成位置对应的结构体数组
df_result = df.withColumn("type_value_label", arrays_zip(
    split(col("types"), "\\|"),
    split(col("values"), "\\|"),
    split(col("labels"), "\\|")
))

# 爆炸结构体数组,提取对应字段
df_result = df_result.withColumn("tv_l", explode(col("type_value_label")))\
           .select(
               col("id"),
               col("tv_l.0").alias("type"),
               col("tv_l.1").alias("value"),
               col("tv_l.2").alias("label")
           )

df_result.show()

原理说明:

  • split将各列按|拆分为数组;
  • arrays_zip把三个数组中相同索引的元素组合成结构体,保证位置对应;
  • explode将结构体数组拆分为多行;
  • 最后通过select提取结构体中的元素并重命名为目标列。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.13 14:15:38