MapReduce: map, shuffle и reduce на примере Word Count
Разбираем MapReduce без магии: map, shuffle и reduce, Word Count на Python, повтор задач, combiner и границы пакетной модели.
MapReduce — это одновременно программная модель и система выполнения пакетных вычислений над большими наборами данных. Разработчик описывает две предметные функции: map превращает входную запись в промежуточные пары «ключ — значение», а reduce объединяет все значения одного ключа. Разбиение данных, запуск задач, пересылку промежуточных результатов и повтор работы после сбоя берёт на себя среда выполнения.
Главная идея статьи Джеффри Дина и Санджая Гемавата 2004 года не в том, что подсчёт слов требует двух функций. Она в том, что узкий контракт между пользовательским кодом и runtime позволяет одинаково исполнять одну задачу на сотнях и тысячах машин.
Три этапа без магии
Представим два документа:
- A: «raft map raft»;
- B: «map bloom».
Map независимо обрабатывает каждую запись и выдаёт пары:
- A → (raft, 1), (map, 1), (raft, 1);
- B → (map, 1), (bloom, 1).
Затем происходит shuffle. Система группирует промежуточные значения по ключу и доставляет каждую группу нужному reducer:
- bloom → [1];
- map → [1, 1];
- raft → [1, 1].
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 делал вместо прикладного кода
Описанная авторами реализация выполняла следующую работу:
- Разбивала вход на M фрагментов и назначала map-задачи свободным workers.
- Буферизовала промежуточные пары и разделяла их на R областей локального диска.
- Передавала reducer адреса этих областей; reducer читал их удалённо и группировал записи по ключу.
- Записывала результат каждого reducer в отдельный выходной файл.
- Отслеживала workers и повторно назначала незавершённые задачи после сбоя.
Локальность данных была частью дизайна, а не мелкой оптимизацией. Планировщик старался запустить map рядом с репликой входного блока, чтобы не пересылать весь вход по сети. Для медленных «хвостовых» задач система могла запустить резервную копию и принять результат первой завершившейся попытки.
Почему важна детерминированность
Повтор задачи безопасен не автоматически. В статье эквивалентность последовательному выполнению формулируется для случая, когда пользовательские map и reduce детерминированы относительно входа. Если функция читает текущее время, генерирует случайное значение или выполняет внешний побочный эффект, две попытки могут дать разные результаты.
Практическое правило: пользовательская функция должна вычислять значение из входной записи, а фиксацию результата должен контролировать runtime. Запись во внешнюю платёжную систему, отправка письма или неидемпотентный HTTP-запрос внутри mapper превращают штатный retry в риск дублирования.
Combiner: полезная, но не универсальная оптимизация
Word Count создаёт много пар вида (слово, 1). Чтобы не передавать их все по сети, map-worker может заранее сложить локальные единицы. В статье это делает необязательная combiner-функция.
Такое сокращение корректно, когда операция допускает частичное объединение в произвольной группировке. Сумма подходит: (a + b) + c = a + (b + c), а порядок слагаемых не меняет результат. Среднее значение в виде одного числа не подходит: среднее средних групп разного размера искажает итог. Для среднего нужно передавать пару (сумма, количество), а окончательное деление выполнять после объединения всех пар.
Где модель уместна
MapReduce хорошо выражает независимую обработку большого числа записей с последующей агрегацией:
- подсчёт частот и статистики по логам;
- построение инвертированного индекса;
- группировку событий по пользователю или объекту;
- распределённую сортировку и подготовку пакетных витрин.
Полезный диагностический вопрос: можно ли разбить вход на независимые части, представить промежуточный результат парами с устойчивым ключом и объединить значения каждого ключа без скрытого общего состояния? Если да, задача близка к модели.
Где одного MapReduce недостаточно
Оригинальная система создавалась для конечных пакетных заданий. У неё нет ответа на все современные режимы обработки:
- интерактивный запрос с миллисекундной задержкой плохо сочетается с запуском полного batch-job;
- для итеративного алгоритма цепочка jobs может многократно материализовывать промежуточные данные;
- «горячий» ключ способен перегрузить один reducer, даже если остальные простаивают;
- бесконечный поток нельзя честно представить как окончательно собранный набор без определения окон, времени события и правил обработки опоздавших данных.
Последний класс задач отдельно разбирает работа о модели Dataflow 2015 года: для неограниченных и приходящих не по порядку данных приходится явно выбирать компромисс между корректностью, задержкой и стоимостью. Это развитие пространства задач, а не опровержение MapReduce.
Ограничения исходной работы
Статья 2004 года описывает реализацию и производственные нагрузки Google того времени. Это не независимый сравнительный benchmark всех систем распределённой обработки. Конкретные размеры блоков, сеть, диски и архитектура master/workers историчны; переносить их как настройки современного кластера нельзя.
Но сохраняется инженерный урок: небольшой функциональный интерфейс полезен только вместе с точно определёнными shuffle, retry, commit и failure semantics. Реализация map и reduce — малая часть промышленной системы.
Чек-лист перед реализацией похожего pipeline
- Зафиксируйте ключ промежуточной группировки и оцените перекос его распределения.
- Сделайте вычисления детерминированными либо явно спроектируйте идемпотентные побочные эффекты.
- Проверьте, можно ли применять локальный combiner без изменения результата.
- Измеряйте отдельно чтение, shuffle, вычисление и запись: сеть часто скрывает стоимость за простым API.
- Не обещайте exactly-once только потому, что runtime повторяет задачи; гарантия зависит от способа фиксации результата.
MapReduce стоит изучать не как название старого продукта, а как образец границы абстракции: предметная логика остаётся короткой, потому что система исполнения берёт на себя большую часть распределённой сложности.
Связанные уроки
- Частотный анализ: словарь как модель данных задачи — Словарь хранит пары “ключ — значение”. Для школьника это удобный способ считать частоты, хранить результаты участников, связывать имя с баллом или быстро проверять накопленные…
- CSV и внешние данные: формат, валидация и поток обработки — Файл позволяет программе работать с данными, которые не вводятся вручную каждый раз. В школьных задачах часто встречаются текстовые файлы, таблицы CSV и наборы строк с числами.
- Словари в Python: пары ключ-значение — Словари в Python: пары ключ-значение — это урок 34 школьного трека Python. Он нужен, чтобы хранить данные, к которым удобно обращаться по имени ключа.
Источники
Формат и права
- Формат
- Авторский разбор
Атрибуция
Самостоятельный редакционный разбор ЯдроКода по первичным публикациям Jeffrey Dean, Sanjay Ghemawat и Tyler Akidau с соавторами. Факты и идеи изложены редакцией самостоятельно; фрагменты исходного текста, код и иллюстрации не воспроизводятся.
Код, данные и иллюстрации
Схемы и иллюстрации первоисточников не копируются. Для публикации используются только собственные текстовые модели и примеры.