Skip to content

Commit 51c41da

Browse files
CrazyJvmrxin
authored andcommitted
misleading task number of groupByKey
"By default, this uses only 8 parallel tasks to do the grouping." is a big misleading. Please refer to #389 detail is as following code : def defaultPartitioner(rdd: RDD[_], others: RDD[_]*): Partitioner = { val bySize = (Seq(rdd) ++ others).sortBy(_.partitions.size).reverse for (r <- bySize if r.partitioner.isDefined) { return r.partitioner.get } if (rdd.context.conf.contains("spark.default.parallelism")) { new HashPartitioner(rdd.context.defaultParallelism) } else { new HashPartitioner(bySize.head.partitions.size) } } Author: Chen Chao <crazyjvm@gmail.com> Closes #403 from CrazyJvm/patch-4 and squashes the following commits: 42f6c9e [Chen Chao] fix format 829a995 [Chen Chao] fix format 1568336 [Chen Chao] misleading task number of groupByKey (cherry picked from commit 9c40b9e) Signed-off-by: Reynold Xin <rxin@apache.org>
1 parent f0abf5f commit 51c41da

File tree

1 file changed

+2
-2
lines changed

1 file changed

+2
-2
lines changed

docs/scala-programming-guide.md

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -189,8 +189,8 @@ The following tables list the transformations and actions currently supported (s
189189
<tr>
190190
<td> <b>groupByKey</b>([<i>numTasks</i>]) </td>
191191
<td> When called on a dataset of (K, V) pairs, returns a dataset of (K, Seq[V]) pairs. <br />
192-
<b>Note:</b> By default, this uses only 8 parallel tasks to do the grouping. You can pass an optional <code>numTasks</code> argument to set a different number of tasks.
193-
</td>
192+
<b>Note:</b> By default, if the RDD already has a partitioner, the task number is decided by the partition number of the partitioner, or else relies on the value of <code>spark.default.parallelism</code> if the property is set , otherwise depends on the partition number of the RDD. You can pass an optional <code>numTasks</code> argument to set a different number of tasks.
193+
</td>
194194
</tr>
195195
<tr>
196196
<td> <b>reduceByKey</b>(<i>func</i>, [<i>numTasks</i>]) </td>

0 commit comments

Comments
 (0)