Мап-редус в RethinkDB
Мап-редус — это способ обобщить и выполнить агрегационные функции на больших наборах данных, потенциально хранящихся на многих серверах, эффективным образом. Он работает путём параллельной обработки данных на каждом сервере, а затем объединения этих результатов в один набор. Впервые он был разработан компанией Google и позднее был реализован в системах баз данных, таких как Apache Hadoop и MongoDB.

В RethinkDB запросы мап-редус работают с последовательностями и состоят из двух или трёх частей:
- Необязательная операция группировки, которая разбивает элементы последовательности на несколько групп.
- Операция мапирования, которая фильтрует и/или преобразует элементы в последовательности (или каждой группе) в новую последовательность (или сгруппированные последовательности).
- Операция редукции, которая агрегирует значения, полученные в результате мапирования, в одно значение (или одно значение для каждой группы).
Некоторые другие реализации мап-редус, например, Hadoop, также используют этап мапирования для выполнения группировки; реализация RethinkDB явно отделяет их. Это иногда называется «группировка-мапирование-редукция» или GMR. RethinkDB эффективно распределяет запросы GMR по таблицам и фрагментам. Вы пишете запросы GMR с помощью команд group, map и reduce, хотя, как мы увидим в наших примерах, многие команды ReQL компилируются в запросы GMR в фоновом режиме — многие распространённые случаи мап-редуса можно выполнить в одну или две строки ReQL.
Простой пример
Предположим, что вы ведёте блог и хотели бы получить количество записей. Запрос мап-редус для выполнения этой операции будет состоять из следующих шагов:
- Этап мапирования, который преобразует каждую запись в число
1(поскольку мы считаем каждую запись один раз). - Этап редукции, который суммирует количество записей.
Для этого примера нам не понадобится этап группировки.
Для нашего блога у нас есть таблица posts, которая содержит записи блога. Вот пример документа из таблицы. (Мы будем использовать Python для этого примера, но другие драйверы ReQL очень похожи.)
{
"id": "7644aaf2-9928-4231-aa68-4e65e31bf219"
"title": "The line must be drawn here"
"content": "This far, no further! ..."
"category": "Fiction"
}
Сначала мы отобразим каждую запись на число 1:
r.table('posts').map(lambda post: 1)
И суммируем записи с помощью reduce:
r.table('posts').map(lambda post: 1).reduce(lambda a, b: a + b).run(conn)
Во многих случаях, когда может потребоваться запрос GMR, ReQL предоставляет даже более простые функции агрегирования. Этот пример намного проще написать с помощью count:
r.table('posts').count().run(conn)
RethinkDB имеет сокращения для пяти общих операций агрегирования: count, sum, avg, min, и max. На практике вы часто сможете использовать их с group вместо написания собственных map и reduce функций.
Пример с группировкой
Предположим, что в блоге в предыдущем примере вы хотите получить количество записей по каждой категории. Запрос мап-редус для выполнения этой операции будет состоять из следующих шагов:
- Этап группировки, который группирует записи на основе их категории.
- Этап мапирования из предыдущего примера.
- Этап редукции, который суммирует количество записей для каждой группы.
Сначала мы сгруппируем записи:
r.table('posts').group(lambda post: post['category'])
Затем, как и прежде, мы отобразим каждую запись на число 1. Команды после команды group будут применяться к каждому сгруппированному набору.
r.table('posts').group(lambda post: post['category']).map(
lambda post: 1)
И снова, мы суммируем записи с помощью reduce, что в этот раз выводит итоги для каждой группы:
r.table('posts').group(lambda post: post['category']).map(
lambda post: 1).reduce(lambda a, b: a + b).run(conn)
И, конечно же, мы можем использовать count для сокращения этого кода. Мы можем ещё больше сократить его: ReQL позволит вам указать group с именем поля вместо лямбда-функции. Таким образом, упрощённая функция выглядит так:
r.table('posts').group('category').count().run(conn)
Более сложный пример
Этот пример основан на примере из MongoDB. Представьте таблицу заказов, где каждый документ в таблице структурирован следующим образом:
{
"customer_id": "cs11072",
"date": r.time(2014, 27, 2, 12, 13, 09, '-07:00'),
"id": 103,
"items": [
{
"price": 91,
"quantity": 1,
"item_id": "sku10491"
} ,
{
"price": 9,
"quantity": 3,
"item_id": "sku14667"
} ,
{
"price": 37 ,
"quantity": 3,
"item_id": "sku16857"
}
],
"total": 229
}
Сначала вернём общую стоимость по каждому клиенту. Поскольку это предварительно вычисляется по каждому заказу в поле total, это легко сделать с помощью одной из функций агрегирования RethinkDB.
r.table('orders').group('customer_id').sum('total').run(conn)
Теперь для чего-то более сложного: вычисления общего и среднего количества проданных товаров по каждому товару. Для этого мы будем использовать функцию concat_map, которая объединяет мапирование и конкатенацию вместе. В этом случае мы хотим получить последовательность всех проданных товаров во всех заказах с их идентификаторами товаров и количествами. Мы также добавим поле «count», установленное в 1; мы будем использовать его так же, как и мапирование каждой записи в примере с блогом.
r.table('orders').concat_map(lambda order:
order['items'].map(lambda item:
{'item_id': item['item_id'], 'quantity': item['quantity'], 'count': 1}
))
Внутренняя функция map используется только для итерации по товарам в каждом заказе. На этом этапе наш запрос вернёт список объектов, каждый из которых содержит три поля: item_id, quantity и count.
Теперь мы сгруппируем по полю item_id и используем пользовательскую функцию reduce для суммирования количеств и счётчиков.
r.table('orders').concat_map(lambda order:
order['items'].map(lambda item:
{'item_id': item['item_id'], 'quantity': item['quantity'], 'count': 1}
)).group('item_id').reduce(lambda left, right: {
'item_id': left['item_id'],
'quantity': left['quantity'] + right['quantity'],
'count': left['count'] + right['count']
})
Наконец, мы используем ungroup для преобразования этих сгруппированных данных в массив объектов с ключами group и reduction. Поле group будет идентификатором товара для каждой группы; поле reduction будет содержать все товары из функции concat_map, которые относятся к каждой группе. Затем мы ещё раз используем map для итерации по этому массиву, вычисляя среднее значение на этом проходе.
r.table('orders').concat_map(lambda order:
order['items'].map(lambda item:
{'item_id': item['item_id'], 'quantity': item['quantity'], 'count': 1}
)).group('item_id').reduce(lambda left, right: {
'item_id': left['item_id'],
'quantity': left['quantity'] + right['quantity'],
'count': left['count'] + right['count']
}).ungroup().map(lambda group: {
'item_id': group['group'],
'quantity': group['reduction']['quantity'],
'avg': group['reduction']['quantity'] / group['reduction']['count']
}).run(conn)
Вывод будет иметь такой формат:
[
{
"avg": 3.3333333333333,
"quantity": 20,
"item_id": "sku10023"
},
{
"avg": 2.2142857142857,
"quantity": 31,
"item_id": "sku10042"
},
...
]
(Обратите внимание, что JavaScript или другой язык, где операторы + и / не переопределены для работы с ReQL, потребуют использования div и add.)
Как выполняются запросы GMR
Запросы GMR RethinkDB распределяются и параллелизуются по фрагментам и ядрам процессора, когда это возможно. Хотя это позволяет им выполнять запросы эффективно, важно помнить, что функция reduce не вызывается для элементов её входного потока слева направо. Она вызывается для элементов потока в любом порядке или для результата предыдущих вызовов функции.
Вот пример неправильного способа написания предыдущего группированного запроса мап-редукции, просто инкрементируя первое значение, переданное функции редукции:
# Incorrect!
r.table('posts').group(lambda post: post['category']).map(
lambda post: 1).reduce(lambda a, b: a + 1).run(conn)
Предположим, что у нас есть десять документов в одной категории в фрагментированной таблице. Четыре документа находятся в фрагменте 1, шесть — во фрагменте 2. При выполнении неправильного запроса это будет его путь:
- Вычисляется количество документов в фрагменте 1. Запрос возвращает значение
4для фрагмента. - Вычисляется количество документов во фрагменте 2. Запрос возвращает значение
6для фрагмента. - Выполняется окончательный шаг редукции для объединения значений двух фрагментов. Вместо вычисления
4 + 6, запрос выполняет4 + 1.
Будьте внимательны! Убедитесь, что ваша функция редукции не предполагает, что шаг редукции выполняется слева направо!
Дополнительная информация
Для получения дополнительной информации о мап-редусе в целом прочитайте статью в Википедии. Для получения дополнительной информации о реализации RethinkDB просмотрите документацию API.
© RethinkDB contributors
Licensed under the Creative Commons Attribution-ShareAlike 3.0 Unported License.
https://rethinkdb.com/docs/map-reduce/