Spark SQL orderBy跨分区全局范围排序分区问题咨询
需求说明
需要对Spark DataFrame实现满足跨分区全局范围有序的排序:不仅每个分区内部数据有序,任意两个分区的值域完全不重叠,即一个分区的所有元素值要么全部小于等于、要么全部大于等于另一个分区的所有元素值。该结果用于后续配合Window.partitionBy("partitionID")使用窗口函数,当前运行环境为Scala + Spark 1.6。
测试过程与问题复现
初始测试数据集生成
初始代码生成5个分区的测试DataFrame,分区内、分区间均无排序规则,结果符合预期:
val df = sc.parallelize(List(10,8,5,9,1,6,4,7,3,2),5) .toDF("val") .withColumn("partitionID",spark_partition_id) df.show
输出结果:
+---+-----------+ |val|partitionID| +---+-----------+ | 10| 0| | 8| 0| | 5| 1| | 9| 1| | 1| 2| | 6| 2| | 4| 3| | 7| 3| | 3| 4| | 2| 4| +---+-----------+
直接调用orderBy的异常结果
尝试直接使用orderBy实现全局排序,代码如下:
val df2 = df.orderBy("val").withColumn("partitionID2",spark_partition_id) df2.show
实际输出中val列虽然展示层面全局有序,但分区未按值域切分,分区间存在值重叠,不符合需求:
+---+-----------+------------+ |val|partitionID|partitionID2| +---+-----------+------------+ | 1| 2| 2| | 2| 4| 4| | 3| 4| 4| | 4| 3| 3| | 5| 1| 1| | 6| 2| 2| | 7| 3| 3| | 8| 0| 0| | 9| 1| 1| | 10| 0| 0| +---+-----------+------------+
预期效果为排序后连续元素归属同一分区,分区间值域无重叠,参考输出如下:
+---+-----------+------------+ |val|partitionID|partitionID2| +---+-----------+------------+ | 1| 2| 2| | 2| 4| 2| | 3| 4| 4| | 4| 3| 4| | 5| 1| 1| | 6| 2| 1| | 7| 3| 3| | 8| 0| 3| | 9| 1| 0| | 10| 0| 0| +---+-----------+------------+
核心逻辑误区
orderBy算子仅保证最终返回的查询结果行顺序全局有序,不保证物理分区按排序键的值域做严格无重叠切分。Spark SQL默认排序为平衡性能与数据倾斜,采用采样估算方式生成分区边界,属于近似范围分区,天然可能出现跨分区值域重叠的问题。- 排序操作后直接获取
spark_partition_id无法得到连续值对应的分区标识,shuffle读阶段仅按拉取顺序输出行,不会自动将同值域数据归并到同一物理分区。
Spark 1.6环境实现方案
Spark 1.6版本未提供高版本的repartitionByRange精确范围分区算子,可通过以下两种方案实现需求:
方案1:手动按排序位置切分分区(100%保证值域无重叠)
先对排序键生成全局连续行号,按目标分区数将连续行号划分到对应分区,再按分区ID重分区并做分区内排序,完全满足分区间值域无重叠要求,示例代码:
import org.apache.spark.sql.expressions.Window import org.apache.spark.sql.functions._ // 给排序后的每行生成从0开始的全局连续排名 val rankedDf = df.orderBy("val") .withColumn("rn", row_number().over(Window.orderBy("val")) - 1) // 按目标分区数计算每行所属的分区ID val totalCount = df.count() val targetPartitionNum = 5 val withTargetPartitionDf = rankedDf .withColumn("partitionID2", (col("rn") * targetPartitionNum / totalCount).cast("int")) // 按目标分区ID重分区,分区内按val排序 val finalDf = withTargetPartitionDf .repartition(targetPartitionNum, col("partitionID2")) .sortWithinPartitions("val")
方案2:使用RDD层sortByKey实现
RDD层的sortByKey算子范围分区采样精度高于Spark SQL层的orderBy,指定分区数后可基本满足值域不重叠要求(极端数据分布下仍可能出现偏移,要求绝对严谨优先选方案1),示例代码:
val sortedRdd = df.select("val").rdd .map(row => (row.getInt(0), row)) .sortByKey(numPartitions = 5) .map(_._2) // 转换回DataFrame即可使用 val sortedDf = sqlContext.createDataFrame(sortedRdd, df.schema)
内容的提问来源于stack exchange,提问作者Thomas
相关产品推荐
相关产品推荐

