【问题标题】:spark scala: Performance degrade with simple UDF over large number of columnsspark scala:在大量列上使用简单的 UDF 会降低性能
【发布时间】:2022-08-19 04:00:37
【问题描述】:

我有一个包含 1 亿行和 10,000 列的数据框。这些列有两种类型,标准 (C_i) 和动态 (X_i)。这个dataframe是经过一些处理得到的,性能很快。现在只剩下两个步骤:

目标:

  1. 需要使用相同的 C_i 列子集对每个 X_i 执行特定操作。
  2. 将每个 X-i 列转换为 FloatType

    困难:

    1. 随着列数的增加,性能会严重下降。
    2. 一段时间后,似乎只有 1 个执行程序可以工作(%CPU 使用率 < 200%),即使是在具有 100 行和 1,000 列的样本数据上也是如此。如果我将它推到 1,500 列,它就会崩溃。

      最小代码:

      import spark.implicits._
      import org.apache.spark.sql.types.FloatType
      
      // sample_udf
      val foo = (s_val: String, t_val: String) => {
          t_val + s_val.takeRight(1)
      }
      val foos_udf = udf(foo)
      spark.udf.register(\"foos_udf\", foo)
      
      val columns = Seq(\"C1\", \"C2\", \"X1\", \"X2\", \"X3\", \"X4\")
      val data = Seq((\"abc\", \"212\", \"1\", \"2\", \"3\", \"4\"),(\"def\", \"436\", \"2\", \"2\", \"1\", \"8\"),(\"abc\", \"510\", \"1\", \"2\", \"5\", \"8\"))
      
      val rdd = spark.sparkContext.parallelize(data)
      var df = spark.createDataFrame(rdd).toDF(columns:_*)
      df.show()
      
      for (cols <- df.columns.drop(2)) {
          df = df.withColumn(cols, foos_udf(col(\"C2\"),col(cols)))
      }
      df.show()
      
      for (cols <- df.columns.drop(2)) {
          df = df.withColumn(cols,col(cols).cast(FloatType))
      }
      df.show()
      

      1,500 列数据错误:

      Exception in thread \"main\" java.lang.StackOverflowError
          at org.apache.spark.sql.catalyst.plans.logical.LogicalPlan.isStreaming(LogicalPlan.scala:37)
          at org.apache.spark.sql.catalyst.plans.logical.LogicalPlan$$anonfun$isStreaming$1.apply(LogicalPlan.scala:37)
          at org.apache.spark.sql.catalyst.plans.logical.LogicalPlan$$anonfun$isStreaming$1.apply(LogicalPlan.scala:37)
          at scala.collection.LinearSeqOptimized$class.exists(LinearSeqOptimized.scala:93)
          at scala.collection.immutable.List.exists(List.scala:84)
      ...
      

      想法:

      1. 也许var 可以替换,但数据大小接近RAM 的40%。
      2. 也许for 循环为dtype 转换可能会导致性能下降,但我不知道如何,以及有哪些替代方案。通过在互联网上搜索,我看到有人建议基于 foldLeft 的方法,但显然仍然在内部转换为 for 循环。

        对此的任何投入将不胜感激。

    标签: scala apache-spark apache-spark-standalone


    【解决方案1】:

    不确定这是否会修复您的 10000~ 列的性能,但我能够使用以下代码在 1500 本地运行它。

    我提到了第 1 点和第 2 点,它们可能对性能产生了一些影响。请注意,据我了解,foldLeft 应该是一个没有内部 for 循环的纯递归函数,因此在这种情况下它可能会对性能产生影响。

    此外,两个 for 循环可以简化为一个 for 循环,我将其重构为 foldLeft

    如果我们用 spark 函数替换 udf,我们也可能会获得性能提升。

      import spark.implicits._
      import org.apache.spark.sql.types.FloatType
      import org.apache.spark.sql.functions._
    
      // sample_udf
      val foo = (s_val: String, t_val: String) => {
        t_val + s_val.takeRight(1)
      }
      val foos_udf = udf(foo)
      spark.udf.register("foos_udf", foo)
    
      val numberOfColumns = 1500
      val numberOfRows = 100
    
    
      val colNames = (1 to numberOfColumns).map(s => s"X$s")
      val colValues = (1 to numberOfColumns).map(_.toString)
    
      val columns = Seq("C1", "C2") ++ colNames
      val schema = StructType(columns.map(field => StructField(field, StringType)))
    
      val rowFields = Seq("abc", "212") ++ colValues
      val listOfRows = (1 to numberOfRows).map(_ => Row(rowFields: _*))
      val listOfRdds = spark.sparkContext.parallelize(listOfRows)
      val df = spark.createDataFrame(listOfRdds, schema)
    
      df.show()
    
      val newDf = df.columns.drop(2).foldLeft(df)((df, colName) => {
        df.withColumn(colName, foos_udf(col("C2"), col(colName)) cast FloatType)
      })
    
      newDf.show()
    
    

    希望这可以帮助!

    *** 编辑

    找到了一种绕过循环的更好的解决方案。只需使用 SelectExpr 制作一个表达式,这样 sparks 会一次性投射所有列,而无需任何递归。从我之前的例子:

    而不是向左折叠,只需用这些行替换它。我刚刚在本地计算机上用 10k 列 100 行对其进行了测试,持续了几秒钟

      val selectExpression = Seq("C1", "C2") ++ colNames.map(s => s"cast($s as float)")
      val newDf = df.selectExpr(selectExpression:_*)
    

    【讨论】:

    • CPU 使用率反映了非常低的缓存命中率。逐行方法可能是解决方案。
    • 我在编辑中添加的df.selectExpr 方法是否仍然存在问题?
    • 感谢您的帮助。 cast 问题不是主要问题,因为计算本身崩溃了,我可以从 csv 文件中读取数据,我可以将数据转储到该文件中(在最坏的情况下)。但是,能够做到这一点将有助于跳过中间步骤继续计算。计算失败是令人惊讶的,一旦我弄清楚如何修复它(很可能通过转移到逐行计算),我将使用这个非常有用的技巧来完成任务。
    【解决方案2】:

    更快的解决方案是在行本身上调用 UDF,而不是在每一列上调用。由于 Spark 将数据存储为行,因此早期的方法表现出糟糕的性能。

    def my_udf(names: Array[String]) = udf[String,Row]((r: Row) => {
        val row = Array.ofDim[String](names.length)
        for (i <- 0 until row.length) {
                row(i) = r.getAs(i)
        }
        ...
    }
    ...
    val df2 = df1.withColumn(results_col,my_udf(df1.columns)(struct("*"))).select(col(results_col))
    

    可以按照 Riccardo 的建议进行类型转换

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2015-10-11
      • 1970-01-01
      • 1970-01-01
      • 2020-06-25
      • 2016-02-25
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多