groupByKey function for RDDs returns pairs of type (<key_type>, ResultIterable). While I was trying to understand the reason why the ResultIterable class exist, I found the following docstring in the class definition:
A special result iterable. This is used because the standard iterator can not be pickled
Question 1: Why is this iterator "pickleable" and the standard one is not?
I even tried to locally re-write the groupByKey function, but commenting the final .mapValues(ResultIterable), basically copy-pasting only the necessary code that would make the groupByKey function working:
import pyspark
from pyspark.shuffle import Aggregator, ExternalMerger, \
get_used_memory, ExternalSorter, ExternalGroupBy
import os
from pyspark.resultiterable import ResultIterable
sc = pyspark.SparkContext('local[*]')
def portable_hash(x):
if 'PYTHONHASHSEED' not in os.environ:
raise RuntimeError("Randomness of hash of string should be disabled via PYTHONHASHSEED")
if x is None:
return 0
if isinstance(x, tuple):
h = 0x345678
for i in x:
h ^= portable_hash(i)
h *= 1000003
h &= sys.maxsize
h ^= len(x)
if h == -1:
h = -2
return int(h)
return hash(x)
def groupByKey(rdd, numPartitions=None, partitionFunc=portable_hash):
def createCombiner(x):
return [x]
def mergeValue(xs, x):
xs.append(x)
return xs
def mergeCombiners(a, b):
a.extend(b)
return a
memory = rdd._memory_limit()
serializer = rdd._jrdd_deserializer
agg = Aggregator(createCombiner, mergeValue, mergeCombiners)
def combine(iterator):
merger = ExternalMerger(agg, memory * 0.9, serializer)
merger.mergeValues(iterator)
return merger.items()
locally_combined = rdd.mapPartitions(combine, preservesPartitioning=True)
shuffled = locally_combined.partitionBy(numPartitions, partitionFunc)
def groupByKey(it):
merger = ExternalGroupBy(agg, memory, serializer)
merger.mergeCombiners(it)
return merger.items()
return shuffled.mapPartitions(groupByKey, True)#.mapValues(ResultIterable)
test_rdd = sc.parallelize([("a", 1), ("b", 1), ("a", 1)])
sorted(groupByKey(test_rdd).collect())
And this produces the expected result, that is
[('a', [1, 1]), ('b', [1])]