Spark SQL中limit查询多次执行结果不一致问题咨询
Spark 无排序Limit结果不稳定问题原理说明
复现代码
var data = Seq[(String, Int)]() for (i <- 1 until 10000) { val str = f"value: ${i}" data = data :+ (str, i) } val df = spark.sparkContext.parallelize(data).toDF() df.createOrReplaceTempView("v_logs") val a = spark.sql( """ SELECT * FROM v_logs limit 20 """ ) a.show() // 标记1 a.show() // 标记2 a.show() // 标记3 a.select(col("_2")).show() // 标记4 a.select(col("_2")).show() // 标记5 a.select(col("_2")).show() // 标记6
上述Scala编写的Spark代码中,预期标记1、2、3处的三次a.show()返回结果完全一致,标记4、5、6处的三次查询返回结果也完全一致,但实际运行结果不符合预期。添加order by _2子句即可得到稳定结果,该现象由Spark的核心运行机制导致,具体原理如下:
- DataFrame为惰性执行,无缓存时每次action都会重新计算
代码中定义的val a仅保存了查询的逻辑执行计划,并没有将计算结果持久化存储。每一次调用show()这类action算子,都会独立向集群提交一个Spark作业,完整走完从数据源读取到计算输出的全流程,不会复用上一次show()的计算结果,除非手动调用cache()/persist()方法将中间结果缓存到内存/磁盘。 - 无全局排序时,分布式查询不承诺结果的顺序和固定批次
不带ORDER BY的LIMIT语句会触发Spark的快速limit优化:执行阶段首先在每个数据分区本地拉取最多20条记录,再将各分区的返回结果汇总到Driver节点,直接截取前20条输出,整个过程不会做全局排序。
很多开发者会误以为本地Seq通过parallelize转成RDD后会保持原集合的顺序,实际上parallelize会按照默认并行度将集合切分为多个分区,跨分区的读取顺序本身就没有强一致性保证。加上分布式环境下各分区任务的调度顺序受集群资源状态、节点负载、任务本地化策略影响,哪个分区的结果先返回没有固定规则;即使在本地单机模式运行,执行线程调度、JVM GC停顿等随机事件也会改变分区结果的返回顺序。最终Driver端汇总得到的数据集顺序不固定,截取的20条结果自然每次运行都可能存在差异。 - 添加ORDER BY后结果稳定的原因
增加order by _2子句后,查询会触发全局排序Shuffle:所有数据会按照_2字段的排序规则重新分发分区、完成分区内排序,最终limit取数是从全局有序的结果集中截取固定位置的记录。示例代码中_2字段是1到9999的唯一整数值,排序逻辑完全确定,因此每次返回的结果都会保持一致。 - 标记4、5、6处结果不稳定的逻辑完全相同:
select(col("_2"))只是在原有逻辑计划上增加了列裁剪节点,每次调用show()依然会独立触发作业,没有确定性排序规则的前提下,结果自然无法保持稳定。
内容的提问来源于stack exchange,提问作者SungHo Kim
相关产品推荐
相关产品推荐

