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

В RethinkDB запросы map-reduce работают с последовательностями и состоят из двух или трех частей:
- Необязательная операция группировки, которая разделяет элементы последовательности на несколько групп.
- Операция map, которая фильтрует и/или преобразует элементы в последовательности (или в каждой группе) в новую последовательность (или сгруппированные последовательности).
- Операция reduce, которая агрегирует значения, созданные операцией map, в одно значение (или одно значение для каждой группы).
Некоторые другие реализации map-reduce, такие как Hadoop, используют этап отображения для выполнения группировки тоже; реализация RethinkDB явно их разделяет. Иногда это называется «группировка-отображение-сведение» или GMR. RethinkDB эффективно распределяет запросы GMR по таблицам и фрагментам. Вы пишете запросы GMR с помощью команд group, map и reduce, хотя, как мы увидим в наших примерах, многие команды ReQL компилируются в запросы GMR за кулисами — многие распространённые случаи map-reduce можно выполнить в одной или двух строках ReQL.
Простой пример
Предположим, что вы ведёте блог и хотите получить количество записей. Запрос map-reduce для выполнения этой операции будет состоять из следующих шагов:
- Этап map, который преобразует каждую запись в число
1(поскольку мы считаем каждую запись один раз). - Этап reduce, который суммирует количество записей.
Для этого примера нам не понадобится этап group.
Для нашего блога у нас есть таблица 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 функций.
Пример с группировкой
Предположим, в примере с блогом в предыдущем примере вы хотели бы получить количество записей по категориям. Запрос map-reduce для выполнения этой операции будет состоять из следующих шагов:
- Этап 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 не вызывается для элементов его входного потока слева направо. Она вызывается для элементов потока в любом порядке или для вывода предыдущих вызовов функции.
Вот пример неправильного написания предыдущего сгруппированного запроса map-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.
Будьте осторожны! Убедитесь, что ваша функция сведения не предполагает, что этап сведения выполняется слева направо!
Подробнее
Для получения дополнительной информации о map-reduce в целом, прочитайте статью в Википедии. Для получения дополнительной информации о реализации RethinkDB просмотрите документацию API.
© RethinkDB contributors
Licensed under the Creative Commons Attribution-ShareAlike 3.0 Unported License.
https://rethinkdb.com/docs/map-reduce/