Spec-Zone.ru › RethinkDB ruby

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