Skip to content

Commit d5ee249

Browse files
committed
self check
1 parent 7d104eb commit d5ee249

File tree

3 files changed

+5
-6
lines changed

3 files changed

+5
-6
lines changed

core/src/main/scala/org/apache/spark/Dependency.scala

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -65,7 +65,7 @@ abstract class NarrowDependency[T](_rdd: RDD[T]) extends Dependency[T] {
6565
* @param keyOrdering key ordering for RDD's shuffles
6666
* @param aggregator map/reduce-side aggregator for RDD's shuffle
6767
* @param mapSideCombine whether to perform partial aggregation (also known as map-side combine)
68-
* @param shuffleWriterProcessor the processor to control the write behavior in ShuffleMapTask.
68+
* @param shuffleWriterProcessor the processor to control the write behavior in ShuffleMapTask
6969
*/
7070
@DeveloperApi
7171
class ShuffleDependency[K: ClassTag, V: ClassTag, C: ClassTag](

core/src/main/scala/org/apache/spark/shuffle/ShuffleWriterProcessor.scala

Lines changed: 2 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -22,10 +22,9 @@ import org.apache.spark.internal.Logging
2222
import org.apache.spark.rdd.RDD
2323
import org.apache.spark.scheduler.MapStatus
2424

25-
2625
/**
2726
* The interface for customizing shuffle write process. The driver create a ShuffleWriteProcessor
28-
* and put it into [[ShuffleDependency]], and executors use it for write processing.
27+
* and put it into [[ShuffleDependency]], and executors use it in each ShuffleMapTask.
2928
*/
3029
private[spark] trait ShuffleWriteProcessor extends Serializable with Logging {
3130

@@ -75,7 +74,7 @@ private[spark] trait ShuffleWriteProcessor extends Serializable with Logging {
7574

7675

7776
/**
78-
* Default shuffle write processor use the shuffle write metrics reporter in context.
77+
* Default shuffle write processor which use the shuffle write metrics reporter in context.
7978
*/
8079
private[spark] class DefaultShuffleWriteProcessor extends ShuffleWriteProcessor {
8180
override def createMetricsReporter(

sql/core/src/main/scala/org/apache/spark/sql/execution/exchange/ShuffleExchangeExec.scala

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -350,8 +350,8 @@ object ShuffleExchangeExec {
350350
}
351351

352352
/**
353-
* Create a customized [[ShuffleWriteProcessor]] for SQL which wrapping the default metrics
354-
* reporter with [[SQLShuffleWriteMetricsReporter]].
353+
* Create a customized [[ShuffleWriteProcessor]] for SQL which wrap the default metrics reporter
354+
* with [[SQLShuffleWriteMetricsReporter]] as new reporter for [[ShuffleWriteProcessor]].
355355
*/
356356
def createShuffleWriteProcessor(metrics: Map[String, SQLMetric]): ShuffleWriteProcessor = {
357357
(reporter: ShuffleWriteMetricsReporter) => {

0 commit comments

Comments
 (0)