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/