Spark自定义AccumulatorV2使用ArrayBuffer报下标越界问题求解
异常根本原因
- 核心问题出在你对
reset()方法的错误实现,结合Spark累加器的运行流程触发了越界:
Spark为每个Task创建累加器副本时,会先调用copy()生成新实例,再自动调用reset()将副本重置为零状态。你当前的reset()直接调用了ArrayBuffer.clear(),这个操作会清空所有元素并把ArrayBuffer的长度置为0。后续执行add()方法时,你尝试访问arrayBuffer(i),但此时集合长度为0,自然触发下标越界异常。 - 为什么替换为Array后运行正常?
因为Array是定长集合,没有clear()方法,你修改为Array后对应的reset()方法只能将每个位置的元素重置为0,不会改变数组的长度,所以不会出现下标访问越界的问题。
修复方案
修改reset()方法的实现,重置时保持ArrayBuffer的长度为5,仅重置元素值为0即可:
override def reset(): Unit = { // 方式1:直接覆盖每个位置的值为初始的0,不改变集合长度 for (i <- arrayBuffer.indices) { arrayBuffer(i) = 0L } // 方式2:清空后重新填充5个0,也可以保持长度符合预期 // arrayBuffer.clear() // arrayBuffer ++= Seq.fill(5)(0L) }
内容的提问来源于stack exchange,提问作者xj y
相关产品推荐
相关产品推荐

