如何将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
相关产品推荐
相关产品推荐

