Scala中如何将Map转换为以键为列名的Spark DataFrame?
Scala 实现 Map 转 DataFrame(键为列名、值为数据)
需要将Scala中的Map转换为Spark DataFrame,要求Map的键作为DataFrame的列名,Map的值作为对应列的行数据。Python和PySpark中仅需一行代码即可实现,但Scala中尝试多种方式均失败,现提供正确实现方法。
Python/PySpark 参考实现
# Python 示例代码 Example_Map_aka_Dictionary = {"Key 1":["Value 1"], "Key 2":[111111.1111], "Key 3":[["Value_n"]]} # 方法1:Pandas import pandas as pd pd.DataFrame(Example_Map_aka_Dictionary) # 方法2:PySpark pandas API import pyspark.pandas as ps ps.DataFrame(Example_Map_aka_Dictionary) # 方法3:Spark原生API spark.createDataFrame( data=[[Value[0] for Value in Example_Map_aka_Dictionary.values()]], schema=list(Example_Map_aka_Dictionary.keys()) )
注意:方法3中使用双层括号,代表单条行数据的集合。
Scala 中的尝试与报错
Map 构建代码
val Example_Data = Seq(("Big_Data_Value_1", 123.45, List("Big_Data_Value_n"))) val Example_Columns = Seq("Big_Data_Column_1", "Big_Data_Column_2", "Big_Data_Column_n") val Example_df = Example_Data.toDF(Example_Columns: _*) Example_df.show() /* +-----------------+-----------------+------------------+ |Big_Data_Column_1|Big_Data_Column_2| Big_Data_Column_n| +-----------------+-----------------+------------------+ | Big_Data_Value_1| 123.45|[Big_Data_Value_n]| +-----------------+-----------------+------------------+ */ val Example_Map_after_Complicated_Operations = Example_df.columns .map(Column_Title => Column_Title -> "Example String after If Statements") .toMap
失败尝试与报错原因
尝试直接将Map值转成List后调用toDF:
val Values_format_8 = Example_Map_after_Complicated_Operations.values.toList val Keys=Example_Map_after_Complicated_Operations.keys.toList Values_format_8.toDF(Keys: _*)
报错信息:
IllegalArgumentException: requirement failed: The number of columns doesn't match. Old column names (1): value New column names (3): Big_Data_Column_1, Big_Data_Column_2, Big_Data_Column_n
报错原因:直接将List[String]转成DataFrame时,Spark会默认生成一列名为value的DataFrame,与传入的多列名数量不匹配;而网上常见的m.toSeq.toDF("name", "score")仅能生成两列DF,不符合需求。
正确实现方法
方法1:单条行数据场景(对应Python方法3)
要重现Python中的双层括号结构,需将Map的值包装成单元素Seq,且元素为元组/Row,让Spark识别为一行多列的数据:
import spark.implicits._ import org.apache.spark.sql.Row import org.apache.spark.sql.types.{StringType, StructField, StructType} // 保持列顺序(Scala默认Map无序,需按原列名顺序提取) val columns = Example_df.columns.toSeq // 方式1:使用Tuple(适合列数固定的场景) val dataTuple = Seq(columns.map(col => Example_Map_after_Complicated_Operations(col)).toTuple) val resultDF1 = dataTuple.toDF(columns: _*) // 方式2:使用Row(适合动态列数场景) val schema = StructType(columns.map(col => StructField(col, StringType))) val dataRow = Seq(Row.fromSeq(columns.map(col => Example_Map_after_Complicated_Operations(col)))) val resultDF2 = spark.createDataFrame(dataRow, schema) resultDF2.show()
方法2:Map值为Column类型的场景
如果Map的值是org.apache.spark.sql.Column类型(比如动态生成列表达式),可通过select构建DF:
import org.apache.spark.sql.functions.lit // 构建值为Column类型的Map val columnMap = Example_df.columns .map(colName => colName -> lit(s"Computed value for $colName")) .toMap // 生成一行数据并选择目标列 val resultDF = spark.range(1).select(columnMap.values.toSeq: _*) resultDF.show()
内容的提问来源于stack exchange,提问作者Sam
相关产品推荐
相关产品推荐

