Как createCombiner, mergeValue, mergeCombiner работает в CombineByKey в Spark (используя Scala)

Я пытаюсь понять, как работает каждый шаг в combineByKeys.

Может кто-нибудь, пожалуйста, помогите мне понять то же самое для ниже RDD?

val rdd = sc.parallelize(List(
  ("A", 3), ("A", 9), ("A", 12), ("A", 0), ("A", 5),("B", 4), 
  ("B", 10), ("B", 11), ("B", 20), ("B", 25),("C", 32), ("C", 91),
   ("C", 122), ("C", 3), ("C", 55)), 2)

rdd.combineByKey(
    (x:Int) => (x, 1),
    (acc:(Int, Int), x) => (acc._1 + x, acc._2 + 1),
    (acc1:(Int, Int), acc2:(Int, Int)) => (acc1._1 + acc2._1, acc1._2 + acc2._2))

Ответы

Ответ 1

Во-первых, позвольте сломать процесс:

Во-первых, createCombiner создает начальное значение (сумматор) для ключевого первого столкновения с разделом, если он не найден --> (firstValueEncountered, 1). Таким образом, это просто инициализация кортежа с первым значением и счетчиком ключей 1.

Затем mergeValue запускается только в том случае, если для найденного ключа в этом разделе уже создан комбайнер (кортеж в нашем случае) --> (existingTuple._1 + subSequentValue, existingTuple._2 + 1). Это добавляет существующее значение кортежа (в первом слоте) с вновь обнаруженным значением и берет существующий счетчик кортежей (во втором слоте) и увеличивает его. Итак, мы

Затем mergeCombiner берет комбинаторы (кортежи), созданные на каждом разделе, и объединяет их вместе --> (tupleFromPartition._1 + tupleFromPartition2._1, tupleFromPartition1._2 + tupleFromPartition2._2). Это просто добавление значений из каждого кортежа вместе, а счетчики - в один кортеж.

Затем разделим подмножество ваших данных на разделы и увидим его в действии:

("A", 3), ("A", 9), ("A", 12),("B", 4), ("B", 10), ("B", 11)

Раздел 1

A=3 --> createCombiner(3) ==> accum[A] = (3, 1)
A=9 --> mergeValue(accum[A], 9) ==> accum[A] = (3 + 9, 1 + 1)
B=11 --> createCombiner(11) ==> accum[B] = (11, 1)

Раздел 2

A=12 --> createCombiner(12) ==> accum[A] = (12, 1)
B=4 --> createCombiner(4) ==> accum[B] = (4, 1)
B=10 --> mergeValue(accum[B], 10) ==> accum[B] = (4 + 10, 1 + 1)

Объединение разделов вместе

A ==> mergeCombiner((12, 2), (12, 1)) ==> (12 + 12, 2 + 1)
B ==> mergeCombiner((11, 1), (14, 2)) ==> (11 + 14, 1 + 2)

Итак, вы должны вернуть массив примерно так:

Array((A, (24, 3)), (B, (25, 3)))