PySpark中如何为数组元素迭代命名?
在PySpark中将数组元素转为带自定义命名键的Map
首先明确:PySpark的原生ArrayType是基于数字索引的有序集合,不支持自定义命名键。你想要的xyz1: "xxx"这种键值对结构,对应PySpark中的MapType(映射类型)。以下分场景实现需求:
先修正Split的正则转义问题
你的原代码中split(old_df["long_string"], "**")会报错,因为**是正则表达式的特殊量词,必须转义为\\*\\*才能按字面拆分字符串。
方法一:使用高阶函数(PySpark 3.1+ 推荐)
利用transform、arrays_zip和map_from_entries等高阶函数,无需展开数组即可完成转换,效率更高:
from pyspark.sql import functions as F # 1. 拆分字符串得到数组(修正正则转义) df = old_df.withColumn("array_col", F.split(F.col("long_string"), "\\*\\*")) # 2. 将数组转为带自定义键的Map df = df.withColumn( "named_map_col", F.map_from_entries( # 把数组元素和从1开始的序号打包,再转成(key, value)的结构体数组 F.transform( F.arrays_zip(F.col("array_col"), F.sequence(F.lit(1), F.size(F.col("array_col")))), lambda item: F.struct( F.concat(F.lit("xyz"), item["sequence"]).alias("key"), item["array_col"].alias("value") ) ) ) )
执行后,named_map_col的结构即为需求中的键值对形式:
map( xyz1 -> "this is", xyz2 -> "a very long", xyz3 -> "string with a lot", xyz4 -> "of words in it" )
方法二:兼容低版本PySpark(低于3.1)
如果你的PySpark版本不支持高阶函数,可以通过posexplode展开数组,再分组聚合生成Map:
from pyspark.sql import functions as F # 1. 拆分字符串得到数组(修正正则转义) df = old_df.withColumn("array_col", F.split(F.col("long_string"), "\\*\\*")) # 2. 展开数组,获取元素位置(pos从0开始,+1得到1起始的序号) exploded_df = df.select( "*", F.posexplode(F.col("array_col")).alias("pos", "value") ).withColumn("key", F.concat(F.lit("xyz"), F.col("pos") + 1)) # 3. 分组聚合回Map(groupBy需包含原表中所有需要保留的列) result_df = exploded_df.groupBy("long_string") # 替换为你的表主键或全部保留列 .agg( F.map_from_entries(F.collect_list(F.struct("key", "value"))).alias("named_map_col") )
内容的提问来源于stack exchange,提问作者fstr
相关产品推荐
相关产品推荐

