Scala中如何将Map[String,DataFrame]合并为带JSON列的单个DataFrame
问题解答
关于DataFrame Map类型列的问题
Spark DataFrame完全支持Map类型的列,属于Spark SQL的原生支持数据类型,你可以通过org.apache.spark.sql.types.MapType自定义schema声明Map列,也可以通过map()、map_from_entries()等内置函数直接构造Map类型的列。
你当前的需求是将每个DataFrame的全量数据转为JSON数组存入单列,不需要用到Map列也可以实现;如果你需要存储结构化的键值对数据而非JSON字符串,也可以选择用Map类型的数组列。
函数实现
实现逻辑
- 遍历传入Map中的所有DataFrame
- 对每个DataFrame,先将每行的所有字段封装为结构体,再通过聚合操作将所有行的结构体收集为数组,最终将数组转为JSON字符串
- 将所有DataFrame处理得到的单条结果合并为最终的输出DataFrame
完整代码
import org.apache.spark.sql.DataFrame import org.apache.spark.sql.functions.{collect_list, struct, to_json, lit, col} import org.apache.spark.sql.types.{StructType, StructField, StringType} def sample(dfs : Map[String, DataFrame]): DataFrame = { // 定义输出DF的schema val outputSchema = StructType(Seq( StructField("column1", StringType, nullable = true) )) // 处理输入Map为空的边界情况 if (dfs.isEmpty) { return org.apache.spark.sql.SparkSession.getActiveSession.get .createDataFrame( org.apache.spark.sql.SparkSession.getActiveSession.get.sparkContext.emptyRDD[org.apache.spark.sql.Row], outputSchema ) } // 处理每个DF得到单条结果 val resultRows = dfs.values.map { df => df.agg( to_json(collect_list(struct("*"))).alias("column1") ).head() }.toSeq // 构造最终输出DF dfs.values.head.sparkSession.createDataFrame(resultRows, outputSchema) }
输出说明
上述代码的输出完全符合你给出的预期格式:输出DataFrame只有1个column1字段,每一行对应输入Map中一个DataFrame的全量数据JSON数组。
Map类型列实现(可选)
如果你不需要转JSON,而是希望存储结构化的Map类型数组,可以调整聚合逻辑,得到的column1类型为ArrayType(MapType(StringType, StringType)):
df.agg( collect_list( map_from_entries( array(df.columns.map(c => struct(lit(c), col(c).cast(StringType))): _*) ) ).alias("column1") )
内容的提问来源于stack exchange,提问作者minnu
相关产品推荐
相关产品推荐

