Spark Scala按Schema转换DataFrame数据类型问题求助
解决Spark 2.0.2中程序化应用自定义Schema到无类型DataFrame的问题
嘿,我来帮你搞定这个问题!
首先得说说你原来代码的问题:你用schema.map(x => df.withColumn(x._1, col(x._1).cast(x._2)))的时候,忽略了Spark DataFrame的不可变性——每次调用withColumn都会生成一个全新的DataFrame实例,但你只是把这些新实例丢在了map操作里,根本没把它们串联起来或者赋值给变量,所以原来的df完全没被修改,这就是为什么你查df.schema的时候类型还是没变的原因。
然后关于你提到的foldLeft,其实它不是Spark DataFrame或RDD的方法,而是Scala标准库中集合的原生方法,所以在Spark 2.0.2 + Scala 2.11的环境下完全可以用它来实现你的需求,这也是处理这种逐字段迭代转换的最佳实践。
正确的实现代码
先确保你导入了必要的函数:
import org.apache.spark.sql.functions.col
然后用foldLeft遍历你的schema列表,逐步转换每个字段的类型:
// 你的自定义schema列表 val schemaList = List(("ID","double"),("LOG_IN_DATE","date"),("USER","string")) // 用foldLeft迭代转换,生成最终的类型正确的DataFrame val convertedDF = schemaList.foldLeft(df) { (currentDF, fieldMeta) => val (fieldName, targetType) = fieldMeta currentDF.withColumn(fieldName, col(fieldName).cast(targetType)) }
代码解释
foldLeft的初始值是你的原始无类型DataFramedf- 每次迭代时,我们取出当前的DataFrame
currentDF,以及schema列表中的一个字段元组(fieldName, targetType) - 调用
withColumn生成新的DataFrame,将指定字段转换为目标类型,这个新DF会作为下一次迭代的输入 - 最终得到的
convertedDF就是已经应用了所有类型转换的DataFrame,你可以通过convertedDF.schema来验证字段类型是否正确
小提示
如果你的原始DataFrame字段顺序和schema列表的顺序不一致也没关系,withColumn会根据字段名精准匹配,不会影响转换结果。
内容的提问来源于stack exchange,提问作者Sidi
相关产品推荐
相关产品推荐

