MapReduce: map, shuffle и reduce на примере Word Count

Разбираем MapReduce без магии: map, shuffle и reduce, Word Count на Python, повтор задач, combiner и границы пакетной модели.

MapReduce — это одновременно программная модель и система выполнения пакетных вычислений над большими наборами данных. Разработчик описывает две предметные функции: map превращает входную запись в промежуточные пары «ключ — значение», а reduce объединяет все значения одного ключа. Разбиение данных, запуск задач, пересылку промежуточных результатов и повтор работы после сбоя берёт на себя среда выполнения.

Главная идея статьи Джеффри Дина и Санджая Гемавата 2004 года не в том, что подсчёт слов требует двух функций. Она в том, что узкий контракт между пользовательским кодом и runtime позволяет одинаково исполнять одну задачу на сотнях и тысячах машин.

Три этапа без магии

Представим два документа:

Map независимо обрабатывает каждую запись и выдаёт пары:

Затем происходит shuffle. Система группирует промежуточные значения по ключу и доставляет каждую группу нужному reducer:

Reduce сворачивает каждую группу и получает итог: bloom → 1, map → 2, raft → 2. Именно shuffle отличает распределённую модель от обычного вызова двух функций подряд: промежуточные записи нужно разделить, передать по сети и отсортировать или сгруппировать.

Минимальная модель на Python

Этот пример выполняется в одном процессе. Он не является заменой Hadoop, облачному runner или оригинальной системе Google, но делает контракт наблюдаемым:

from collections import defaultdict
from collections.abc import Iterable, Iterator


def map_words(document: str) -> Iterator[tuple[str, int]]:
    for word in document.lower().split():
        yield word, 1


def reduce_counts(word: str, counts: Iterable[int]) -> tuple[str, int]:
    return word, sum(counts)


def word_count(documents: Iterable[str]) -> dict[str, int]:
    shuffled: dict[str, list[int]] = defaultdict(list)

    for document in documents:
        for word, count in map_words(document):
            shuffled[word].append(count)

    return dict(
        reduce_counts(word, counts)
        for word, counts in sorted(shuffled.items())
    )


print(word_count(["raft map raft", "map bloom"]))
# {'bloom': 1, 'map': 2, 'raft': 2}

В реальном кластере список shuffled не собирается в памяти одного процесса. Map-задачи пишут разделы промежуточных данных, а reduce-задачи забирают относящиеся к ним разделы. В реализации из статьи пространство промежуточных ключей делилось на R частей, обычно через hash(key) mod R. Поэтому одинаковые ключи попадали к одному reducer.

Что runtime делал вместо прикладного кода

Описанная авторами реализация выполняла следующую работу:

  1. Разбивала вход на M фрагментов и назначала map-задачи свободным workers.
  2. Буферизовала промежуточные пары и разделяла их на R областей локального диска.
  3. Передавала reducer адреса этих областей; reducer читал их удалённо и группировал записи по ключу.
  4. Записывала результат каждого reducer в отдельный выходной файл.
  5. Отслеживала workers и повторно назначала незавершённые задачи после сбоя.

Локальность данных была частью дизайна, а не мелкой оптимизацией. Планировщик старался запустить map рядом с репликой входного блока, чтобы не пересылать весь вход по сети. Для медленных «хвостовых» задач система могла запустить резервную копию и принять результат первой завершившейся попытки.

Почему важна детерминированность

Повтор задачи безопасен не автоматически. В статье эквивалентность последовательному выполнению формулируется для случая, когда пользовательские map и reduce детерминированы относительно входа. Если функция читает текущее время, генерирует случайное значение или выполняет внешний побочный эффект, две попытки могут дать разные результаты.

Практическое правило: пользовательская функция должна вычислять значение из входной записи, а фиксацию результата должен контролировать runtime. Запись во внешнюю платёжную систему, отправка письма или неидемпотентный HTTP-запрос внутри mapper превращают штатный retry в риск дублирования.

Combiner: полезная, но не универсальная оптимизация

Word Count создаёт много пар вида (слово, 1). Чтобы не передавать их все по сети, map-worker может заранее сложить локальные единицы. В статье это делает необязательная combiner-функция.

Такое сокращение корректно, когда операция допускает частичное объединение в произвольной группировке. Сумма подходит: (a + b) + c = a + (b + c), а порядок слагаемых не меняет результат. Среднее значение в виде одного числа не подходит: среднее средних групп разного размера искажает итог. Для среднего нужно передавать пару (сумма, количество), а окончательное деление выполнять после объединения всех пар.

Где модель уместна

MapReduce хорошо выражает независимую обработку большого числа записей с последующей агрегацией:

Полезный диагностический вопрос: можно ли разбить вход на независимые части, представить промежуточный результат парами с устойчивым ключом и объединить значения каждого ключа без скрытого общего состояния? Если да, задача близка к модели.

Где одного MapReduce недостаточно

Оригинальная система создавалась для конечных пакетных заданий. У неё нет ответа на все современные режимы обработки:

Последний класс задач отдельно разбирает работа о модели Dataflow 2015 года: для неограниченных и приходящих не по порядку данных приходится явно выбирать компромисс между корректностью, задержкой и стоимостью. Это развитие пространства задач, а не опровержение MapReduce.

Ограничения исходной работы

Статья 2004 года описывает реализацию и производственные нагрузки Google того времени. Это не независимый сравнительный benchmark всех систем распределённой обработки. Конкретные размеры блоков, сеть, диски и архитектура master/workers историчны; переносить их как настройки современного кластера нельзя.

Но сохраняется инженерный урок: небольшой функциональный интерфейс полезен только вместе с точно определёнными shuffle, retry, commit и failure semantics. Реализация map и reduce — малая часть промышленной системы.

Чек-лист перед реализацией похожего pipeline

  1. Зафиксируйте ключ промежуточной группировки и оцените перекос его распределения.
  2. Сделайте вычисления детерминированными либо явно спроектируйте идемпотентные побочные эффекты.
  3. Проверьте, можно ли применять локальный combiner без изменения результата.
  4. Измеряйте отдельно чтение, shuffle, вычисление и запись: сеть часто скрывает стоимость за простым API.
  5. Не обещайте exactly-once только потому, что runtime повторяет задачи; гарантия зависит от способа фиксации результата.

MapReduce стоит изучать не как название старого продукта, а как образец границы абстракции: предметная логика остаётся короткой, потому что система исполнения берёт на себя большую часть распределённой сложности.

Источники

Формат и права

Формат
Авторский разбор

Атрибуция

Самостоятельный редакционный разбор ЯдроКода по первичным публикациям Jeffrey Dean, Sanjay Ghemawat и Tyler Akidau с соавторами. Факты и идеи изложены редакцией самостоятельно; фрагменты исходного текста, код и иллюстрации не воспроизводятся.

Код, данные и иллюстрации

Схемы и иллюстрации первоисточников не копируются. Для публикации используются только собственные текстовые модели и примеры.