【问题标题】:Counting filtered items on dataframe SPARK计算数据框 SPARK 上的过滤项目
【发布时间】:2018-09-20 01:21:07
【问题描述】:

我有以下数据框:df

在某些时候,我需要根据时间戳(毫秒)过滤掉项目。 但是,保存过滤了多少记录对我来说很重要(如果太多,我想让工作失败) 天真地我能做到:

======Lots of calculations on df ======
val df_filtered = df.filter($"ts" >= startDay && $"ts"  <= endDay)
val filtered_count = df.count - df_filtered.count

但感觉完全是矫枉过正,因为 SPARK 将执行整个执行树 3 次(过滤和 2 个计数)。 Hadoop MapReduce 中的这项任务非常简单,因为我可以为过滤的每一行维护计数器。 有没有更有效的方法,只能找到累加器却无法连接过滤器。

建议的方法是在过滤器之前缓存 df,但是由于 DF 大小,我更喜欢这个选项作为最后的手段。

【问题讨论】:

  • 与计数相比,这种方法如何更快?有没有办法精确计算?它仍然没有处理计数本身之前的所有工作,这将在没有缓存的情况下执行所有计算的 3 倍
  • df.except(df_filtered).count 怎么样?
  • 累加器使用示例可参考here。但请注意,在转换中使用累加器并不是 100% 准确,因为在失败的情况下可以多次执行转换。但我想在你的情况下,这件事可以忽略。另一个缺点是当你发现错误太多时不能立即停止处理,因为在过滤器操作中无法获取累加器值。所以你需要在所有处理后失败。
  • @VladislavVarslavans 您能否给出简单的示例或链接,说明如何将累加器与数据框过滤器一起使用,因为我看不到它。如果我过滤这些值,我该如何为它们存储累加器

标签: scala apache-spark dataframe bigdata


【解决方案1】:

Spark 1.6.0 代码:

import org.apache.spark.sql.SQLContext
import org.apache.spark.{SparkConf, SparkContext}

object Main {

  val conf = new SparkConf().setAppName("myapp").setMaster("local[*]")
  val sc = new SparkContext(conf)
  val sqlContext = new SQLContext(sc)

  case class xxx(a: Int, b: Int)

  def main(args: Array[String]): Unit = {

    val df = sqlContext.createDataFrame(sc.parallelize(Seq(xxx(1, 1), xxx(2, 2), xxx(3,3))))

    val acc = sc.accumulator[Long](0)

    val filteredRdd = df.rdd.filter(r => {
      if (r.getAs[Int]("a") > 2) {
        true
      } else {
        acc.add(1)
        false
      }
    })

    val filteredRddDf = sqlContext.createDataFrame(filteredRdd, df.schema)

    filteredRddDf.show()

    println(acc.value)
  }
}

Spark 2.x.x 代码:

import org.apache.spark.sql.SparkSession

object Main {

  val ss = SparkSession.builder().master("local[*]").getOrCreate()
  val sc = ss.sparkContext

  case class xxx(a: Int, b: Int)

  def main(args: Array[String]): Unit = {

    val df = ss.createDataFrame(sc.parallelize(Seq(xxx(1, 1), xxx(2, 2), xxx(3,3))))

    val acc = sc.longAccumulator

    val filteredDf = df.filter(r => {
      if (r.getAs[Int]("a") > 2) {
        true
      } else {
        acc.add(1)
        false
      }
    }).toDF()


    filteredDf.show()

    println(acc.value)

  }
}

【讨论】:

  • 完成(见答案):)
  • 非常感谢!我不知道当我运行它时出现什么问题,我在累加器中得到 0。我将我的代码与您的代码相匹配,知道吗?
  • 好吧,我想我明白了,我只有在执行 show() 时才得到 accu 值。对生产代码而不是 show() 有什么建议吗?
  • 任何会触发计算所有元素的动作都可以。您可以使用foreachcount。但是您必须使用一个动作 - 因为它是 Spark 的基础,所以计算只能由动作触发。 show 可能涉及一些优化,因为它只打印一些元素。
猜你喜欢
  • 2014-11-23
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多