Spark中重复执行相同Repartition操作仍产生Shuffle的问题咨询
我之前也碰到过一模一样的情况,其实这属于Spark的预期行为,背后主要有这几个原因:
1. 文本读写后的分区数不一致
你第一次用partitionBy(new HashPartitioner(10))把RDD分成了10个分区,但用saveAsTextFile保存时,Spark不会强制生成10个文件——如果某个分区没有数据(比如空分区),Spark就不会创建对应的空文件。当你用sc.textFile读取S3上的文件时,Spark会根据文件数量和默认的分区策略(比如minPartitions参数)来确定读取后的分区数。如果这个分区数不等于10,后续再执行partitionBy(new HashPartitioner(10))时,Spark就会触发Shuffle来重新调整分区数量到10,这就会产生少量的Shuffle数据。
2. 字符串哈希的隐性不一致(概率较低)
虽然你的id是String类型,两次分区都用了HashPartitioner,但要注意:HashPartitioner依赖的是Java的String.hashCode()。如果在JSON序列化/反序列化过程中,某些特殊字符(比如Unicode字符、不可见字符)的处理出现细微差异(比如转义方式变化),会导致同一个id的哈希值在两次处理中不一样,这些数据就会被分配到不同的分区,从而触发Shuffle。不过这种情况比较少见,更多是前面的分区数问题。
3. 空值或异常数据的影响
如果数据中存在id为null的记录,第一次分区时null会被分配到固定的分区(通常是分区0)。但如果在JSON反序列化时,null被转换成了空字符串"",那么它的哈希值会和原来的null完全不同,第二次分区时就会被分配到其他分区,这部分数据就会产生Shuffle。
验证与解决方法
验证步骤
- 先检查读取后的RDD分区数:
val rddACloneRaw = sc.textFile(s3OutputPath) .map(x => objectMapper.readValue(x, classOf[APojo])) .map(x => (x.getId(), x)) println(rddACloneRaw.getNumPartitions) // 看是否等于10 - 采样每个分区的
id.hashCode() % 10值,确认是否和分区索引一致,判断是否有数据哈希不匹配的情况。
解决方法
方法1:保留分区信息的序列化方式
放弃文本格式,用saveAsObjectFile保存RDD,这样读取时会直接恢复原有的分区器,后续不需要再执行partitionBy:
// 第一次保存 val rddA = sc.textFile(s3InputPath) .map(x => objectMapper.readValue(x, classOf[APojo])) .map(x => (x.getId(), x)) .partitionBy(new HashPartitioner(10)) rddA.saveAsObjectFile(s3OutputPath) // 读取后直接使用,无Shuffle val rddAClone = sc.objectFile[(String, APojo)](s3OutputPath) .map(x => objectMapper.writeValueAsString(x._2))
方法2:强制指定分区数并手动设置分区器
如果必须用文本格式,可以先确保分区数正确,再手动设置分区器(前提是确认读取后的分区数据符合哈希规则):
val rddACloneRaw = sc.textFile(s3OutputPath) .map(x => objectMapper.readValue(x, classOf[APojo])) .map(x => (x.getId(), x)) .coalesce(10, shuffle = false) // 强制调整到10个分区,不触发Shuffle // 手动设置分区器,让后续操作感知分区规则 rddACloneRaw.partitioner = Some(new HashPartitioner(10)) val rddAClone = rddACloneRaw.map(x => objectMapper.writeValueAsString(x._2))
方法3:自定义分区验证逻辑
如果想彻底避免Shuffle,可以在读取后先验证每个分区的哈希范围,确认符合预期后再跳过partitionBy操作。
内容的提问来源于stack exchange,提问作者Bakmy

