如何规范化以嵌套数组为主的Spark DataFrame?
规范化嵌套数组类型的Spark DataFrame
针对你提供的嵌套数组DataFrame,我们可以分几步来做规范化处理,涵盖清洗元素空格、字符串转数组、展开嵌套数组、扁平化深层嵌套这些常见场景:
1. 清洗数组内元素的多余空格
你的foo和baz数组里的元素都带有多余的前后空格(比如"1 "),可以用Spark的transform函数遍历数组每个元素,调用trim去除空格;同时把bar这种空格分隔的字符串转成规范数组:
import org.apache.spark.sql.functions.{trim, transform, split} // 清洗数组元素空格,同时把bar字段转成无空元素的数组 val cleanedDF = f .withColumn("foo_clean", transform($"foo", item => trim(item))) .withColumn("baz_clean", transform($"baz", item => trim(item))) .withColumn("bar_array", split(trim($"bar"), "\\s+")) // 先trim整个bar字符串,再按任意数量空格拆分
执行后cleanedDF的输出会把数组元素的空格去掉,bar也从杂乱的字符串变成了规整的数组:
+------+-------------------+-------------+-------------------+------------+------------+---------+ |id |foo |bar |baz |foo_clean |baz_clean |bar_array| +------+-------------------+-------------+-------------------+------------+------------+---------+ |thinga|[1 ] |1 2 3 |[2 ] |[1] |[2] |[1,2,3] | |thinga|[1 2 3 4 ] | 0 0 0 |[2 3 4 5 ] |[1 2 3 4] |[2 3 4 5] |[0,0,0] | |thingb|[1 2 ] |1 2 3 4 5 |[1 2 ] |[1 2] |[1 2] |[1,2,3,4,5]| |thingb|[0 , 0 , 0 ] |1 2 3 4 5 |[1 2 3 ] |[0, 0, 0] |[1 2 3] |[1,2,3,4,5]| +------+-------------------+-------------+-------------------+------------+------------+---------+
2. 展开嵌套数组为多行
如果需要把数组的每个元素拆成单独的行(方便做聚合或明细分析),用explode函数即可:
import org.apache.spark.sql.functions.explode // 展开foo_clean数组为多行 val explodedDF = cleanedDF .select($"id", explode($"foo_clean").alias("foo_item"), $"bar_array", $"baz_clean") explodedDF.show(false)
输出会把每个数组元素拆成独立行:
+------+---------+---------+------------+ |id |foo_item |bar_array|baz_clean | +------+---------+---------+------------+ |thinga|1 |[1,2,3] |[2] | |thinga|1 2 3 4 |[0,0,0] |[2 3 4 5] | |thingb|1 2 |[1,2,3,4,5]|[1 2] | |thingb|0 |[1,2,3,4,5]|[1 2 3] | |thingb|0 |[1,2,3,4,5]|[1 2 3] | |thingb|0 |[1,2,3,4,5]|[1 2 3] | +------+---------+---------+------------+
如果需要同时展开多个数组并保持元素的位置对应(比如foo和baz的元素一一配对),可以用arrays_zip先把数组打包,再explode:
import org.apache.spark.sql.functions.{arrays_zip, explode_outer} val zippedExplodedDF = cleanedDF .select($"id", $"bar_array", arrays_zip($"foo_clean", $"baz_clean").alias("foo_baz_pair")) .select($"id", $"bar_array", explode_outer($"foo_baz_pair").alias("pair")) .select($"id", $"pair.foo_clean".alias("foo_item"), $"pair.baz_clean".alias("baz_item"), $"bar_array") zippedExplodedDF.show(false)
3. 扁平化深层嵌套的数组
如果你的数组元素本身是空格分隔的字符串(比如foo_clean里的"1 2 3 4"),需要进一步拆成单个元素的数组,可以用split+flatten组合:
import org.apache.spark.sql.functions.{split, flatten} val flattenedDF = cleanedDF .withColumn("foo_flattened", flatten(transform($"foo_clean", item => split(item, "\\s+")))) .withColumn("baz_flattened", flatten(transform($"baz_clean", item => split(item, "\\s+")))) flattenedDF.show(false)
这会把数组里的空格分隔字符串完全拆成单个元素,得到彻底扁平化的数组:
+------+-------------------+-------------+-------------------+------------+------------+---------+---------------+---------------+ |id |foo |bar |baz |foo_clean |baz_clean |bar_array|foo_flattened |baz_flattened | +------+-------------------+-------------+-------------------+------------+------------+---------+---------------+---------------+ |thinga|[1 ] |1 2 3 |[2 ] |[1] |[2] |[1,2,3] |[1] |[2] | |thinga|[1 2 3 4 ] | 0 0 0 |[2 3 4 5 ] |[1 2 3 4] |[2 3 4 5] |[0,0,0] |[1,2,3,4] |[2,3,4,5] | |thingb|[1 2 ] |1 2 3 4 5 |[1 2 ] |[1 2] |[1 2] |[1,2,3,4,5]|[1,2] |[1,2] | |thingb|[0 , 0 , 0 ] |1 2 3 4 5 |[1 2 3 ] |[0, 0, 0] |[1 2 3] |[1,2,3,4,5]|[0,0,0] |[1,2,3] | +------+-------------------+-------------+-------------------+------------+------------+---------+---------------+---------------+
内容的提问来源于stack exchange,提问作者Georg Heiler
相关产品推荐
相关产品推荐

