Spark并行执行两个Action时,如何确保源数据仅被读取一次?
Spark并行触发count()操作的竞态条件与缓存机制解析
你的代码示例:
Dataset data = spark.read().createView("test").cache(); CompletableFuture.allOf( supplyAsync(() -> data.count()), supplyAsync(() -> data.count()) );
问题解答:
不会出现竞态条件
Spark的Dataset是不可变对象,所有操作都只会生成新的Dataset实例,不会修改原对象的状态。并行调用count()只是基于同一个不可变的Dataset执行计算,不存在多个线程对共享可变数据的修改操作,自然不会产生竞态条件。不会重复读取源数据的原因
你代码里调用了.cache()方法,这个方法会标记Dataset需要被缓存到Spark的存储系统(默认是内存优先,内存不足时写入磁盘):
- 第一个
count()触发时,Spark会执行完整的读取源数据、计算流程,同时把计算结果写入缓存; - 第二个
count()执行时,Spark会先检查缓存中是否存在该Dataset的缓存数据,确认存在后直接从缓存中读取数据计算,不会再去访问源数据。
需要注意的是,.cache()是懒加载的——只有当第一次触发终端操作(比如count())时,才会实际执行缓存写入操作,后续的终端操作都会复用缓存内容。
内容的提问来源于stack exchange,提问作者sharin gan
相关产品推荐
相关产品推荐

