When is ResultIterable necessary?

Viewed 34

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])]

Question 2: starting from my UDF groupByKey, could you provide an example where the missing .mapValues(ResultIterable) becomes necessary?

0 Answers
Related