【问题标题】:Spark columnar performanceSpark柱状性能
【发布时间】:2017-09-13 08:10:06
【问题描述】:

我是 Spark 的相对初学者。我有一个宽数据框(1000 列),我想根据相应列是否缺少值来添加列

所以

+----+ |一个 | +----+ | 1 | +----+ |空| +----+ | 3 | +----+

变成

+----+-------+ |一个 |管理信息系统 | +----+-------+ | 1 | 0 | +----+-------+ |空| 1 | +----+-------+ | 3 | 1 | +----+-------+

这是自定义 ml 转换器的一部分,但算法应该是清晰的。

override def transform(dataset: org.apache.spark.sql.Dataset[_]): org.apache.spark.sql.DataFrame = {
  var ds = dataset
  dataset.columns.foreach(c => {
    if (dataset.filter(col(c).isNull).count() > 0) {
      ds = ds.withColumn(c + "_MIS", when(col(c).isNull, 1).otherwise(0))
    }
  })


  ds.toDF()
}

循环遍历列,如果 > 0 个空值创建一个新列。

传入的数据集被缓存(使用 .cache 方法),相关配置设置为默认值。 这目前在一台笔记本电脑上运行,即使只有最少的行,1000 列的运行时间约为 40 分钟。 我认为问题出在数据库中,所以我尝试使用镶木地板文件,但结果相同。查看作业 UI,它似乎正在执行文件扫描以进行计数。

有没有办法改进这个算法以获得更好的性能,或者以某种方式调整缓存?增加 spark.sql.inMemoryColumnarStorage.batchSize 只是让我出现 OOM 错误。

【问题讨论】:

    标签: scala apache-spark spark-dataframe


    【解决方案1】:

    删除条件:

    if (dataset.filter(col(c).isNull).count() > 0) 
    

    只留下内部表达式。正如其所写,Spark 需要 #columns 数据扫描。

    如果您希望修剪列计算统计信息一次,如Count number of non-NaN entries in each column of Spark dataframe with Pyspark 中所述,并使用单个drop 调用。

    【讨论】:

    • 谢谢。尝试了这种方法,不到 2 分钟,这对我来说是可以接受的。我将在下面发布我的代码。
    【解决方案2】:

    这是解决问题的代码。

    override def transform(dataset: Dataset[_]): DataFrame = {
      var ds = dataset
      val rowCount = dataset.count()
      val exprs = dataset.columns.map(count(_))
      val colCounts = dataset.agg(exprs.head, exprs.tail: _*).toDF(dataset.columns: _*).first()
      dataset.columns.foreach(c => {
        if (colCounts.getAs[Long](c) > 0 && colCounts.getAs[Long](c) < rowCount   ) {
          ds = ds.withColumn(c + "_MIS", when(col(c).isNull, 1).otherwise(0))
        }
      })
      ds.toDF()
    }
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2016-06-04
      • 1970-01-01
      • 1970-01-01
      • 2020-04-16
      • 2018-02-18
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多