SparkUI Stages输入大小差异:repartitionByCassandraReplica与DirectJoin疑问
Spark与Cassandra两种Join方式的性能疑问解答
集群环境与测试背景
- 16节点集群,每节点同时部署Spark与Cassandra
- Cassandra副本因子:3
- Spark配置:
spark.sql.shuffle.partitions=96 - Spark-Cassandra-Connector版本:3.1.0
- 测试场景:基于分区键将DataFrame与84.05Gb的Cassandra表执行Join操作,对比两种Join方式的SparkUI Stages输入数据量与耗时:
- 方式1:
repartitionByCassandraReplica().JoinWithCassandraTable() - 方式2:始终开启DirectJoin的Join方式
- 方式1:
测试场景数据
场景1
- 配置:PartitionsPerHost=100(共1600个Spark分区)
- DirectJoin:Input 36.9Gb,耗时4.1分钟
- repartitionByCassandraReplica:Input 68.9Gb,耗时3.4分钟
场景2
- 配置:PartitionsPerHost=2(共32个Spark分区,与DirectJoin分区数一致)
- repartitionByCassandraReplica:Input 68.9Gb,耗时4.2分钟
疑问解答
1. 场景1中repartitionByCassandraReplica输入更大但耗时更少?是否仅因数据本地化?
数据本地化是核心,但不是唯一原因:
- 极致的本地化利用:该方式会将Spark分区与Cassandra副本节点严格绑定,每个Spark任务直接从本地Cassandra节点读取数据,完全消除跨节点数据拉取的网络延迟、带宽占用等开销;而DirectJoin的本地化策略并非严格绑定副本,部分任务仍需从远程节点拉取数据,额外开销被放大。
- 高并行度的资源利用:场景1中该方式的分区数(1600)远高于DirectJoin的32,更多并行任务能充分压榨集群的CPU、IO资源,抵消了数据量更大带来的处理开销。
- 无额外Shuffle开销:DirectJoin在部分场景下会触发Shuffle来对齐分区,而repartitionByCassandraReplica提前将DataFrame分区与Cassandra副本对齐,跳过了不必要的Shuffle步骤,节省了磁盘IO与网络传输时间。
2. 两种方式输入大小为何不同?
输入差异源于读取Cassandra数据的策略本质不同:
- DirectJoin是按需读取:它会先根据DataFrame中的分区键值,精准计算出Cassandra中需要读取的Token范围,只拉取与Join相关的分区数据,因此Input大小(36.9Gb)远小于原表的84.05Gb,相当于前置了一次分区过滤。
- repartitionByCassandraReplica是按副本全量读取:该方式按Cassandra的副本分布划分Spark分区,每个分区会读取对应节点上的全量副本数据,后续再在本地过滤掉与Join无关的部分,但SparkUI统计的是读取的原始输入大小,因此输入量更大(68.9Gb)。
3. 分区数相同时输入大小为何仍有差异?
场景2中两种方式分区数一致,但输入量仍不同,原因在于:
- repartitionByCassandraReplica的分区划分逻辑是绑定Cassandra的节点与Token范围,不管分区数多少,每个Spark分区对应的是节点上一段Token范围的全量副本数据——分区数仅影响并行处理的粒度,不改变读取的数据范围,因此输入量始终维持在68.9Gb。而DirectJoin始终基于Join键值精准过滤需要读取的数据,输入量始终是匹配的36.9Gb左右。
内容的提问来源于stack exchange,提问作者ktzan
相关产品推荐
相关产品推荐

