【问题标题】:Spark Accumulator throws a class cast exception when trying to count the number of records in the datasetSpark Accumulator 在尝试计算数据集中的记录数时抛出类转换异常
【发布时间】:2021-02-19 06:31:13
【问题描述】:

我正在尝试计算我的数据集中的记录数。我正在使用累加器尝试以下逻辑。

    val accum = sc.longAccumulator("My_Accum")
    val fRDD = tempDS.rdd.persist(StorageLevel.MEMORY_AND_DISK).foreach(x=>{
      accum.add(1)
      x
    })

    val recordCount = accum.value

    println("record Count is : "+recordCount)

我在 accum.add(1) 线上遇到了一个类转换异常 java.lang.ClassCastException: org.apache.spark.sql.catalyst.expressions.GenericRowWithSchema 无法转换为 packagename.CaseClassName。

使用相同的逻辑,我可以在上一步中获得累加器值

谁能帮我解决这个问题。除了 count() 和累加器之外,还有其他方法可以计算吗?

【问题讨论】:

  • 看看从 1 改 1L 是否有帮助
  • 这段代码对我来说很好用,spark 2.4 和 scala 2.11
  • 它在一种情况下对我有用,但是当我尝试与它相同的逻辑时,我收到了这个错误。我认为当我尝试执行 dataset.foreach 或 dataset.map 时,它无法识别架构并将其视为 GenericRow 而不是我的案例类

标签: scala apache-spark dataset rdd


【解决方案1】:

我认为问题不在于累加器,似乎 Spark 无法将您的 DataFrame,即 Dataset[Row] 转换为 Dataset[packagename.CaseClassName](您没有在代码中显示)。

此外,这是一种非常规的计算行数的方法,我不建议这样做。最快的方法是在你的Dataset 上使用.count,在大多数情况下使用RDD 会更慢

【讨论】:

  • 我使用的是 Dataset[MyClass] 而不是数据框。计数也是一项昂贵的操作,所以我想使用累加器
猜你喜欢
  • 2020-11-16
  • 2023-04-09
  • 1970-01-01
  • 1970-01-01
  • 2017-08-28
  • 2015-06-23
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多