Scala中如何从Spark推荐结果DataFrame生成两列新DataFrame?
解决Spark推荐结果DataFrame格式转换问题
嘿,看你正在用Spark做协同过滤推荐,已经拿到用户31511的Top10推荐结果啦,现在要把嵌套的数组结构拆成扁平的productionId和rating两列,这个不难,我给你一步步说:
首先,先回顾下你当前results的结构和输出:
当前results的输出和Schema
scala> results.show +------+--------------------+ |userId| recommendations| +------+--------------------+ | 31511|[[328, 0.7845393]...| +------+--------------------+ scala> results.printSchema root |-- userId: integer (nullable = false) |-- recommendations: array (nullable = true) | |-- element: struct (containsNull = true) | | |-- productionId: integer (nullable = true) | | |-- rating: float (nullable = true)
你的目标是转换成这样的扁平结构:
+------------+------+ |productionId|rating| +------------+------+ | 1 | 1 | +------------+------+ | 2 | 4.5 | +------------+------+ | 3 | 5 | +------------+------+ | 4 | 2.5 | +------------+------+
实现步骤
要把数组列拆成多行,我们需要用Spark的explode函数,它能把数组里的每个元素拆成单独的行。具体代码如下:
// 先导入explode函数 import org.apache.spark.sql.functions.explode // 转换DataFrame val flattenedDF = results .select(explode($"recommendations").alias("rec")) // 展开数组列,给每个struct元素起别名rec .select($"rec.productionId", $"rec.rating") // 从struct里提取需要的字段 // 查看结果 flattenedDF.show()
这样就能得到你想要的扁平结构啦!如果还需要保留userId字段,也可以在select的时候加上$"userId",比如:
val flattenedDFWithUserId = results .select($"userId", explode($"recommendations").alias("rec")) .select($"userId", $"rec.productionId", $"rec.rating")
内容的提问来源于stack exchange,提问作者Daniel
相关产品推荐
相关产品推荐

