You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Spark分区与Executor绑定问题及重复重分区优化咨询

问题解答

基础问题:分区与Executor的分配逻辑

当多次对同一DataFrame使用相同分区键+相同分区数的HashPartitioner执行repartition时:

  • 分区内容会保持一致:因为key_hash % partitions_number的计算结果固定,相同键的数据始终会被分到同一个分区ID中。
  • 分区ID不会固定绑定到某个Executor:Spark的调度器会根据集群实时资源状态(如Executor负载、内存可用性)动态分配分区到不同Executor处理。不过,如果数据已经存储在某个节点的内存/磁盘上,Spark会优先重用本地数据,避免不必要的网络传输——也就是说,分区内容不会在Executor间无意义地移动。

实际场景:提前repartition的价值分析

你的代码逻辑

spark.sql("select key, topic, row_number() over (partition by key,topic order by key) as rn from x").createOrReplaceTempView("y")
// MapPartitions操作,生成DataFrame z
spark.sql("select key, collect_list(row_number) from z group by key")

执行计划中重复出现的Exchange hashpartitioning(rowkey#817, 10),是因为MapPartitions操作破坏了Spark的分区感知能力,导致后续group by需要再次触发重分区。

提前repartition的价值

提前按key执行repartition是有实际价值的,核心原因如下:

  1. 避免真实的数据移动:即使Spark触发了第二次重分区,由于两次重分区的键和分区数完全一致,Spark的Exchange算子会检测到数据已经符合目标分区规则,仅调整分区元数据,不会执行实际的网络shuffle。
  2. 优化中间计算的本地性:提前按key分区后,后续的窗口函数、MapPartitions操作都可以在本地分区内执行,无需先按(key, topic)分区再重新洗牌,减少了中间阶段的数据移动开销。
  3. 降低重计算成本:如果你的MapPartitions操作涉及HBase读取、Avro反序列化等重计算,提前分区可以让这些计算在正确的分区上完成,避免后续重分区带来的重复计算。

总结:提前执行相同规则的repartition,能有效减少Executor间的数据传输,尤其在中间操作涉及大量计算时,收益更为明显。


中文翻译后的执行计划

+- 窗口函数 [row_number() over(窗口定义(rowkey#817, column_name#815 ASC NULLS FIRST, 指定窗口范围(RowFrame, 无界前置$(), 当前行$())) AS segment_index#839], 分区键[rowkey#817], 排序规则[column_name#815 ASC NULLS FIRST]
                  +- *(5) 排序 [rowkey#817 ASC NULLS FIRST, column_name#815 ASC NULLS FIRST], 不全局排序, 0
                     +- 交换操作 hash分区(rowkey#817, 10), ENSURE_REQUIREMENTS, [id=#647]
                        +- *(4) 投影 [rowkey#817, topic#808, field1#809, field2#810, field3#811, field4#812, field5#813, column_name#815, struct(trid, cast(value#828.rt_trid as string), dataDt, value#828.rt_data_dt, dataDtUtc, value#828.rt_start_date_utc, eventIdx, cast(value#828.rt_event_idx as int), lastEvent, value#828.rt_last_event, recordType, value#828.rt_record_type, latitude, value#828.geo.rt_latitude, longitude, value#828.geo.rt_longitude, altitude, value#828.geo.rt_altitude, country, value#828.address.rt_country, county, value#828.address.rt_county, state, value#828.address.rt_state, ... 58个更多字段) AS json_struct#881]
                           +- *(4) 投影 [cast(rowkey#807 as string) AS rowkey#817, topic#808, field1#809, field2#810, field3#811, field4#812, field5#813, column_name#815, avro反序列化表达式(value#816, hbase, (hbase.columns.cq,schema), (hbase.namespace,dev), (hbase.columns.name.value,EnrichedJourney), (hbase.columns.namespace.value,com.vodafone.automotive.wasp.telematics.datamodel.journey), (hbase.columns.cf,0), (hbase.table,SCHEMA_REPOSITORY), true) AS value#828]
                              +- *(4) 从对象序列化 [if (断言非空(input[0, org.apache.spark.sql.Row, true]).isNullAt) null else 验证外部类型(get外部行字段(断言非空(input[0, org.apache.spark.sql.Row, true]), 0, rowkey), BinaryType) AS rowkey#807, if (断言非空(input[0, org.apache.spark.sql.Row, true]).isNullAt) null else 静态调用(class org.apache.spark.unsafe.types.UTF8String, StringType, fromString, 验证外部类型(get外部行字段(断言非空(input[0, org.apache.spark.sql.Row, true]), 1, topic), StringType), true, false, true) AS topic#808, if (断言非空(input[0, org.apache.spark.sql.Row, true]).isNullAt) null else 验证外部类型(get外部行字段(断言非空(input[0, org.apache.spark.sql.Row, true]), 2, field1), BinaryType) AS field1#809, if (断言非空(input[0, org.apache.spark.sql.Row, true]).isNullAt) null else 验证外部类型(get外部行字段(断言非空(input[0, org.apache.spark.sql.Row, true]), 3, field2), BinaryType) AS field2#810, if (断言非空(input[0, org.apache.spark.sql.Row, true]).isNullAt) null else 验证外部类型(get外部行字段(断言非空(input[0, org.apache.spark.sql.Row, true]), 4, field3), BinaryType) AS field3#811, if (断言非空(input[0, org.apache.spark.sql.Row, true]).isNullAt) null else 验证外部类型(get外部行字段(断言非空(input[0, org.apache.spark.sql.Row, true]), 5, field4), BinaryType) AS field4#812, if (断言非空(input[0, org.apache.spark.sql.Row, true]).isNullAt) null else 验证外部类型(get外部行字段(断言非空(input[0, org.apache.spark.sql.Row, true]), 6, field5), BinaryType) AS field5#813, if (断言非空(input[0, org.apache.spark.sql.Row, true]).isNullAt) null else 静态调用(class org.apache.spark.unsafe.types.UTF8String, StringType, fromString, 验证外部类型(get外部行字段(断言非空(input[0, org.apache.spark.sql.Row, true]), 8, column_name), StringType), true, false, true) AS column_name#815, if (断言非空(input[0, org.apache.spark.sql.Row, true]).isNullAt) null else 验证外部类型(get外部行字段(断言非空(input[0, org.apache.spark.sql.Row, true]), 9, value), BinaryType) AS value#816]
                                 +- MapPartitions ######.hbase.HBaseReader$$$Lambda$3992/510160290@44b639d3, obj#806: org.apache.spark.sql.Row
                                    +- 反序列化为对象 创建外部行(rowkey#787, topic#738.toString, field1#749, field2#751, field3#753, field4#755, field5#757, StructField(rowkey,BinaryType,true), StructField(topic,StringType,true), StructField(field1,BinaryType,true), StructField(field2,BinaryType,true), StructField(field3,BinaryType,true), StructField(field4,BinaryType,true), StructField(field5,BinaryType,true)), obj#805: org.apache.spark.sql.Row
                                       +- *(3) 投影 [cast(rowkey#737 as binary) AS rowkey#787, topic#738, field1#749, field2#751, field3#753, field4#755, field5#757]
                                          +- *(3) 过滤 castProcessBooleanUDF(consolidationResult#759)
                                             +- 排序聚合 键=[rowkey#737, topic#738, column_family#739], 函数=[max(if ((column_name#740 <=> field1)) value#741 else null), max(if ((column_name#740 <=> field2)) value#741 else null), max(if ((column_name#740 <=> field3)) value#741 else null), max(if ((column_name#740 <=> field4)) value#741 else null), max(if ((column_name#740 <=> field5)) value#741 else null), max(if ((column_name#740 <=> consolidationResult)) value#741 else null)]
                                                +- 排序聚合 键=[rowkey#737, topic#738, column_family#739], 函数=[partial_max(if ((column_name#740 <=> field1)) value#741 else null), partial_max(if ((column_name#740 <=> field2)) value#741 else null), partial_max(if ((column_name#740 <=> field3)) value#741 else null), partial_max(if ((column_name#740 <=> field4)) value#741 else null), partial_max(if ((column_name#740 <=> field5)) value#741 else null), partial_max(if ((column_name#740 <=> consolidationResult)) value#741 else null)]
                                                   +- *(2) 排序 [rowkey#737 ASC NULLS FIRST, topic#738 ASC NULLS FIRST, column_family#739 ASC NULLS FIRST], 不全局排序, 0
                                                      +- 交换操作 hash分区(rowkey#737, 10), REPARTITION_BY_COL, [id=#630]

内容的提问来源于stack exchange,提问作者D. belvedere

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.29 20:17:01