【问题标题】:Scala & Spark: Cast multiple columns at onceScala & Spark:一次投射多列
【发布时间】:2017-06-19 05:51:04
【问题描述】:

由于VectorAssembler 崩溃,如果传递的列具有除NumericTypeBooleanType 之外的任何其他类型,并且我正在处理很多TimestampType 列,我想知道:

有没有简单的方法,一次投射多列

基于this answer,我已经有了一种投单列的便捷方式:

def castColumnTo(df: DataFrame, 
    columnName: String, 
    targetType: DataType ) : DataFrame = {
      df.withColumn( columnName, df(columnName).cast(targetType) )
}

我曾考虑递归调用castColumnTo,但我强烈怀疑这是(高性能)方法。

【问题讨论】:

  • 是什么阻止您遍历列并调用此函数(无需递归)?
  • 为什么需要递归?你的意思是迭代?请记住,Spark 是懒惰的,因此没有明显的理由说明它的性能不够

标签: scala apache-spark


【解决方案1】:

基于 cmets(谢谢!)我想出了以下代码(未实现错误处理):

def castAllTypedColumnsTo(df: DataFrame, 
   sourceType: DataType, targetType: DataType) : DataFrame = {

      val columnsToBeCasted = df.schema
         .filter(s => s.dataType == sourceType)

      //if(columnsToBeCasted.length > 0) {
      //   println(s"Found ${columnsToBeCasted.length} columns " +
      //      s"(${columnsToBeCasted.map(s => s.name).mkString(",")})" +
      //      s" - casting to ${targetType.typeName.capitalize}Type")
      //}

      columnsToBeCasted.foldLeft(df){(foldedDf, col) => 
         castColumnTo(foldedDf, col.name, LongType)}
}

感谢鼓舞人心的 cmets。 foldLeft(解释为 herehere)保存 for 循环以迭代 var 数据帧。

【讨论】:

  • 您可以使用foldLeft 代替for 循环并避免var,因为我认为这是改进的样式。 val dfReturn = columnsToBeCasted.foldLeft(df){(accdf, col) => castColumnTo(accdf, col.name, LongType)}
  • @TheArchetypalPaul 出于好奇..为什么你不发布答案,即使你清楚地知道答案:)
  • @rogue-one 不能被打扰 :) 我真的不需要代表,只是在调整 Boen 的答案,所以他应该得到荣誉。而且我在工作,所以无论如何我都处于驾车模式
  • 感谢 cmets !我会更新我的答案。对于任何对折叠感兴趣的人:oldfashionedsoftware.com/2009/07/30/…
  • @Boern,fold 比这更酷:cs.nott.ac.uk/~pszgmh/fold.pdf
【解决方案2】:

在 scala 中使用惯用方法投射所有列

def castAllTypedColumnsTo(df: DataFrame, sourceType: DataType, targetType: DataType) = {
df.schema.filter(_.dataType == sourceType).foldLeft(df) {
    case (acc, col) => acc.withColumn(col.name, df(col.name).cast(targetType))
 }
}

【讨论】:

  • @TheArchetypalPaul 删除了不必要的分配 stmt.. 并进行了建议的更改
  • 是的,现在很清楚了。但是,我确实赞成将 OP 分成两行,因为它提供了添加他使用的日志记录的机会。你可以在一行 scala 中做很多事情,但有时当你回到它时,要花很长时间才能弄清楚它做了什么。但这是一个品味问题,也许只是我。我将删除我的第一条评论
  • 两个答案都有帮助,但是对于这两个答案,如果我更改 1000 列,将创建多少个 DataFrame 实例? foldLeft 会改善垃圾收集吗?
  • 我问是因为 withColumn 函数在每次调用时都会创建一个新的数据框,尽管 foldLeft 维护一个实例
【解决方案3】:
FastDf = (spark.read.csv("Something.csv", header = False, mode="DRPOPFORMED"))
FastDf.OldTypes = [feald.dataType for feald in FastDf.schema.fields]
FastDf.NewTypes = [StringType(), FloatType(), FloatType(), IntegerType()]
FastDf.OldColnames = FastDf.columns
FastDf.NewColnames = ['S_tring', 'F_loat', 'F_loat2', 'I_nteger']
FastDfSchema = FastDf.select(*
                             (FastDf[colnumber]
                              .cast(FastDf.NewTypes[colnumber])
                              .alias(FastDf.NewColnames[colnumber]) 
                                  for colnumber in range(len(FastDf.NewTypes)
                                                )
                             )
                            )

我知道它在 pyspark 中,但逻辑可能很方便。

【讨论】:

    【解决方案4】:

    我正在为 python 翻译 scala 程序。我为您的问题找到了明智的答案。 该列名为 V1 - V28、时间、金额、类别。 (我不是 Scala 专业人士) 解决方案如下所示。

    // cast all the column to Double type.
    val df = raw.select(((1 to 28).map(i => "V" + i) ++ Array("Time", "Amount", "Class")).map(s => col(s).cast("Double")): _*)
    

    链接:https://github.com/intel-analytics/analytics-zoo/blob/master/apps/fraudDetection/Fraud%20Detction.ipynb

    【讨论】:

    • 这可以扩展为提供模式:val df = raw.select(schema.map{case (s, coltype) => col(s).cast(coltype)}.toList: _*) 其中模式是地图集合,例如val schema = Map("evar39" -> "DOUBLE", "evar46" -> "STRING")
    猜你喜欢
    • 2021-12-15
    • 2019-02-21
    • 2018-10-13
    • 1970-01-01
    • 2017-12-15
    • 1970-01-01
    • 1970-01-01
    • 2021-07-29
    • 2021-11-20
    相关资源
    最近更新 更多