Spec-Zone.ru › RethinkDB javascript

Map-reduce в RethinkDB

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

  • Простой пример
  • Пример с группировкой
  • Более сложный пример
  • Как выполняются запросы GMR
  • Подробнее

Map-reduce Illustration

В 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. Вычисляется количество документов на фрагменте 1. Запрос возвращает значение 4 для фрагмента.
  2. Вычисляется количество документов на фрагменте 2. Запрос возвращает значение 6 для фрагмента.
  3. Выполняется заключительный этап сведения для объединения значений двух фрагментов. Вместо вычисления 4 + 6, запрос выполняет 4 + 1.

Будьте осторожны! Убедитесь, что ваша функция сведения не предполагает, что этап сведения выполняется слева направо!

Подробнее

Для получения дополнительной информации о map-reduce в целом, прочитайте статью в Википедии. Для получения дополнительной информации о реализации RethinkDB просмотрите документацию API.

  • group
  • map
  • reduce
  • ungroup
  • concat_map

© RethinkDB contributors
Licensed under the Creative Commons Attribution-ShareAlike 3.0 Unported License.
https://rethinkdb.com/docs/map-reduce/

Spec-Zone.ru

Настройки Оффлайн Что нового Помощь О нас
Spec-Zone .ru
спецификации, руководства, описания, API