【发布时间】: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