为何Spark Accumulator输出类型需线程安全?附自定义累加器疑问
问题
根据AccumulatorV2的定义,其输出类型需支持原子性读取(如Int、Long)或具备线程安全性(如同步集合),因为该类型会被其他线程读取。
假设我定义了一个名为CheckSumAccumulator的类,它继承自AccumulatorV2,输出类型为CheckSum。CheckSumAccumulator包含一个私有字段checkSum;CheckSum类包含私有字段count,以及对应的公共getter和setter方法。
代码示例:
public class CheckSumAccumulator extends AccumulatorV2<String, CheckSum> { private CheckSum checkSum; ... } public class CheckSum extends Serializable { private long count; public long getCount() { return count; } public void setCount(long count) { this.count = count; } }
请问这种情况可能会引发什么问题?每个Executor中的Accumulator实例是否运行在单线程中?
回答
一、可能引发的问题
- 线程安全风险:
CheckSum里的count字段既不是原子类型(比如AtomicLong),读写操作也没加同步控制。Executor端执行任务时,同一个CheckSumAccumulator实例可能被多个任务线程同时更新或读取,会导致count出现竞态条件——比如多线程同时做递增操作,最终结果会比预期值小,或者读取到中间的不一致状态。 - 结果不一致:Driver端获取Accumulator最终值时,因为
CheckSum非线程安全,可能读到未完全更新的count,导致统计结果不准确。
二、Executor中的Accumulator实例是否单线程运行?
不是。Executor中的Accumulator实例会被多个任务线程共享使用。Spark Executor默认以多线程方式执行任务(线程池大小通常对应CPU核心数),同一个Accumulator实例会被这些线程同时操作,所以必须保证其输出类型的线程安全性。
内容的提问来源于stack exchange,提问作者user1532146
相关产品推荐
相关产品推荐

