Учебник веб-разработки
Разделы учебника
На этой странице

VIII. Системы хранения

MapReduce: группировка, частичные итоги и границы распределения

Оглавление · Аналитика · FP

Задача: разложить агрегацию на map, группировку и reduce, затем проверить, можно ли объединять частичные результаты. Нужны функции и ассоциативность.

От записей к итогам

Map выдаёт промежуточные пары ключ/значение. Shuffle собирает значения одного ключа, reduce вычисляет итог группы. Распределённая реализация дополнительно организует передачу данных, планирование и восстановление задач. Работа Dean и Ghemawat.

Исходник схемы
flowchart LR
  Input["События"] --> Map["Map: topic → 1"]
  Map --> Shuffle["Shuffle: все значения одного topic вместе"]
  Shuffle --> Reduce["Reduce: сумма по topic"]
  Reduce --> Output["Итоги тем"]

Метод массива reduce — локальная свёртка. Он не запускает кластер, не перемещает данные между узлами и не восстанавливает потерянную работу.

Исполняемая локальная модель

Сохраните блок во временный lesson.mts и запустите Node 24.

export function mapReduce<A, K, V, R>(
  input: readonly A[],
  map: (item: A) => readonly (readonly [K, V])[],
  reduce: (key: K, values: readonly V[]) => R,
): Map<K, R> {
  const groups = new Map<K, V[]>();
  for (const item of input) {
    for (const [key, value] of map(item)) {
      const group = groups.get(key);
      if (group) group.push(value);
      else groups.set(key, [value]);
    }
  }
  const result = new Map<K, R>();
  for (const [key, values] of groups) result.set(key, reduce(key, values));
  return result;
}

export type LessonEvent = Readonly<{ id: number; topic: string }>;
export function countTopics(events: readonly LessonEvent[]): Map<string, number> {
  return mapReduce<LessonEvent, string, number, number>(
    events,
    (event) => [[event.topic, 1]],
    (_topic, values) => values.reduce((sum, value) => sum + value, 0),
  );
}
export function countPartitions(
  partitions: readonly (readonly LessonEvent[])[],
): Map<string, number> {
  const partials = partitions.flatMap((part) => [...countTopics(part)]);
  return mapReduce<[string, number], string, number, number>(
    partials,
    (pair) => [pair],
    (_topic, values) => values.reduce((sum, value) => sum + value, 0),
  );
}

const events: readonly LessonEvent[] = [
  { id: 1, topic: "react" },
  { id: 2, topic: "react" },
  { id: 3, topic: "data" },
  { id: 4, topic: "react" },
];
console.log([...countTopics(events)].sort(([a], [b]) => (a < b ? -1 : a > b ? 1 : 0)));
console.log(
  [...countPartitions([events.slice(0, 2), events.slice(2)])].sort(([a], [b]) =>
    a < b ? -1 : a > b ? 1 : 0,
  ),
);
// Оба результата: [["data", 1], ["react", 3]]

Ключи generic Map сравнивает по правилам JavaScript: два разных объекта не станут одной группой по равенству полей. В примере ключи — строки. Требуются корректные проверенные записи, чистые callbacks и счётчики в пределах точных целых number. Readonly ограничивает API при компиляции, не замораживает callback-вход в runtime.

Локальные мутации принадлежат одному вызову; исходные события не меняются. Map хранит группы целиком в памяти, порядок итерации следует первому появлению ключа. Для сравнения результатов мы явно сортируем строки; такой порядок не обещается любым распределённым runtime.

Закон объединения

Для количества combine(a, b) = a + b, нейтральный элемент — 0. В области точных целых сложение ассоциативно и коммутативно: разбиение и порядок частей не меняют итог. Повтор одной части меняет сумму: сложение не идемпотентно. Свойства объединения и дедупликация — разные вопросы.

Среднее средних не работает при разных размерах групп: для [0, 10] и [100] оно равно 52.5, а среднее всех трёх — 110/3. Объединяйте (sum, count) покомпонентно, делите в конце. Для float порядок сложения может менять округление; проверьте требования к воспроизводимости, прежде чем объявлять закон точным.

Чего модель не проверяет

Исходник схемы
flowchart TD
  Parts["Части входа"] --> Tasks["Попытки задач"]
  Tasks --> Transfer["Передача и группировка"]
  Transfer --> Commit["Согласованная фиксация итогов"]
  Tasks --> Retry["Повтор после сбоя"]
  Retry --> Tasks

Это обязанности распределённой системы, отсутствующие в коде выше. Нельзя просто прибавить результат каждой попытки: повтор потерянной/неподтверждённой задачи способен удвоить часть итога. Внешний эффект внутри mapper требует своей политики повторов; чистое преобразование проще переисполнять.

Партиционирование по ключу направляет одинаковые ключи одному владельцу группы. Популярный ключ может перегрузить его — data skew. Combiner уменьшает пересылку только при допустимом объединении; функция reduce не обязана подходить для этой роли. Hadoop и Spark — дальнейшее знакомство с системами исполнения, не зависимости этого учебного примера.

Практика

Сравните прямой итог со всеми разрезами маленького массива, пустыми частями и перестановками. Повторите запись: id в этом коде не дедуплицируется, она считается ещё одним событием. Сопоставьте это с count(*) SQL из аналитики.

Покажите контрпример для среднего средних, повторённой части и float-ассоциативности. Затем нарисуйте, как результат попытки задачи отличать от результата самой задачи. Локальная проверка не подтверждает отказоустойчивость кластера или время выполнения SQL.