【问题标题】:How to avoid using withColumn iteratively in Spark Scala?如何避免在 Spark Scala 中迭代地使用 withColumn?
【发布时间】:2019-08-13 02:46:03
【问题描述】:

目前,我们有一个迭代使用 withColumn 的代码。 When-Otherwise 条件检查,我们在此基础上进行算术计算。

示例代码:

df.withColumn("col4", when(col("col1")>10, col("col2").+col("col3")).otherwise(col("col2")))

对于另外 40 列,其他算术计算以迭代方式发生。 总记录数 - 2M。

重复的代码:

  df.withColumn(colName1, when((col(amountToSubtract).<("0")) && (col(colName1).===("0")) && (col(amountToSubtract).<=(col(colName2))) && (col(colName2).!==("0")), col(colName2)).otherwise(col(colName1))).
      withColumn(amountToSubtract, when(col(amountToSubtract).!==("0"), col(amountToSubtract).-(col(colName2))).otherwise(col(amountToSubtract))).
      withColumn(colName1, when((col(amountToSubtract).>("0")) && (col(colName1).===("0")), col(colName2).+(col(amountToSubtract))).otherwise(col(colName1))).
      withColumn(amountToSubtract, when((col(amountToSubtract).>("0")) && (col(colName1).===("0")), "0").otherwise(col(amountToSubtract)))

随后进行 7 或 8 组其他计算。

此时作业挂起的时间更长。有时它会引发 GC 开销错误。 意识到 .withColumn 的迭代使用不利于性能,我找不到实现上述条件检查的替代方法。请协助。

【问题讨论】:

  • 您可以使用 foldleft 并添加具有相同逻辑的列。您获取列列表并使用数据框作为零元素应用 foldleft。在 foldleft 中应用 withColumn 方法。
  • @firas 感谢您的意见!我提到的重复代码 sn-p 是一个块。并且相同的 4 块 withColumns 又被使用了 8 次。因此,我发现使用 FoldLeft 很难实现它。如果可以,请为我的示例代码分享一个示例 foldleft - 我可能会进一步扩展它。
  • 我发布了答案。如果有帮助,你能把它标记为好。如果没有告诉我,我会尝试更新。
  • @firas Havent 尝试了您的解决方案。因为 UDF 对我来说更熟悉,所以我已经尝试过这里提供的其他解决方案。下周将尝试您的建议。将让您保持不变。

标签: scala dataframe apache-spark apache-spark-sql


【解决方案1】:

您可以尝试将当前使用“何时/否则”实现的所有逻辑封装到一个 UDF 中,该 UDF 将生成新列值时需要考虑的所有列值的数组作为输入,并作为输出返回所有生成的列值的数组。 Udfs 有时有自己的性能问题,但可能值得一试。这是我正在考虑的技术的简单说明:

object SO extends App {

  val sparkSession = SparkSession.builder().appName("simple").master("local[*]").getOrCreate()
  sparkSession.sparkContext.setLogLevel("ERROR")

  import sparkSession.implicits._

  case class Record(col1: Int, col2: Int, amtToSubtract: Int)

  val recs = Seq(
    Record(1, 2, 3),
    Record(11, 2, 3)
  ).toDS()


  val colGenerator : Seq[Int] => Seq[Int] =
    (arr: Seq[Int]) =>  {
      val (in_c1, in_c2, in_amt_sub) = (arr(0), arr(1), arr(2))

      val newColName1_a = if (in_amt_sub < 0 && in_c1 == 0 && in_amt_sub < in_c2 &&  in_c2 != 0) {
        in_c2
      }
      else {
        in_c1
      }
      val newAmtSub_a = if (in_amt_sub != 0) {
        in_amt_sub - in_c2
      } else {
        in_amt_sub
      }

      val newColName1_b  = if (  newAmtSub_a  > 0 &&  newColName1_a  == 0 ) {
        in_c2  + newAmtSub_a
      } else {
        newColName1_a
      }

      val newAmtSub_b = if (newAmtSub_a  > 0  && newColName1_b   == 0) {
        0
      } else {
        newAmtSub_a
      }

      Seq(newColName1_b,  newAmtSub_b)
    }

  val colGeneratorUdf = udf(colGenerator)

  // Here the first column in the generated array is 'col4', the UDF could equivalently generate as many
  // other values as you want from the input array of column values.
  //
  val afterUdf = recs.withColumn("colsInStruct",   colGeneratorUdf (array($"col1", $"col2", $"amtToSubtract")))
  afterUdf.show()
  // RESULT
  //+----+----+-------------+------------+
  //|col1|col2|amtToSubtract|colsInStruct|
  //+----+----+-------------+------------+
  //|   1|   2|            3|      [1, 1]|
  //|  11|   2|            3|     [11, 1]|
  //+----+----+-------------+------------+


}

【讨论】:

  • 感谢您的帮助。这里的挑战是,1. 全部 40 列的何时-否则逻辑不同。 2. 我在这里看到的主要问题是 - 迭代 withColumns 的使用。因此,将 UDF 集成到这种已经降低性能的迭代 withColumn 代码只会使情况变得更糟,这是我所理解的。如果我的理解是错误的,请纠正。尽管如此,已经开始尝试 UDF 选项。将让您保持不变。
  • 知道了。我认为您可以做的是移动 UDF 中所有 40 列的逻辑。那有意义吗?如果不是...请随时使用更多逻辑来更新您的示例以获取其他列,如果不清楚,我可以向您展示如何将其移至 UDF。
  • 嗨,克里斯,根据要求 - 用多次重复的代码更新了问题。
  • 我更新了答案.. 不确定我的商业逻辑 100% 正确.. 但希望能传达总体思路;^)
  • @DasarathyDR - 进展如何?我很好奇您是否尝试过并且性能有所提高。希望它有效!
【解决方案2】:

我不确定您尝试应用的逻辑。但是这里是 foldleft 的想法:

  val colList = List("col1", "col2", "col3")
  val df: DataFrame = ???
  colList.foldLeft(df){case(df, colName1) => df
     .withColumn(colName1, when((col(amountToSubtract).<("0")) && (col(colName1).=== ("0")) && (col(amountToSubtract).<=(col(colName2))) && (col(colName2).!==("0")), col(colName2)).otherwise(col(colName1))).
     .withColumn(amountToSubtract, when(col(amountToSubtract).!==("0"), col(amountToSubtract).-(col(colName2))).otherwise(col(amountToSubtract))).
     .withColumn(colName1, when((col(amountToSubtract).>("0")) && (col(colName1).===("0")), col(colName2).+(col(amountToSubtract))).otherwise(col(colName1))).
     .withColumn(amountToSubtract, when((col(amountToSubtract).>("0")) && (col(colName1).===("0")), "0").otherwise(col(amountToSubtract)))

}

【讨论】:

    猜你喜欢
    • 2019-03-20
    • 1970-01-01
    • 1970-01-01
    • 2018-12-26
    • 1970-01-01
    • 2021-07-13
    • 2020-10-04
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多