Spec-Zone.ru › RethinkDB python

Потоки изменений в RethinkDB

Изменения данных лежат в основе функций RethinkDB в реальном времени.

  • Базовое использование
  • Точечные (изменения одного документа) changefeeds
  • Changefeeds с фильтрацией и агрегационными запросами
  • Включение изменений состояния
  • Включение начальных значений
  • Включение типов результатов
  • Обработка задержек
  • Масштабирование
  • Подробнее

Они позволяют клиентам получать изменения в таблице, отдельном документе или даже результатах определенного запроса по мере их возникновения. Почти любой запрос ReQL можно преобразовать в changefeed.

Data Modeling Illustration

Базовое использование

Подпишитесь на поток, вызвав changes для таблицы:

feed = r.table('users').changes().run(conn)
for change in feed:
    print change

Команда changes возвращает курсор (как и команды table или filter). Вы можете итерироваться по его содержимому с помощью ReQL. В отличие от других курсоров, вывод changes является бесконечным: курсор будет блокироваться до тех пор, пока не появятся новые элементы. Каждый раз, когда вы вносите изменения в таблицу или документ, за которыми следит changes поток, в курсор будет возвращаться новый объект. Например, если вы вставляете пользователя {id: 1, name: Slava, age: 31} в таблицу users, RethinkDB опубликует этот документ в changefeeds, подписанные на users.

{
  'old_val': None,
  'new_val': { 'id': 1, 'name': 'Slava', 'age': 31 }
}

Здесь old_val - старая версия документа, а new_val - новая версия. При insert, old_val будет null; при delete, new_val будет null. При update, и old_val, и new_val присутствуют.

Точечные (изменения одного документа) changefeeds

Точечный changefeed возвращает изменения в одном документе в таблице, а не в таблице в целом.

r.table('users').get(100).changes().run(conn)

Формат вывода точечного changefeed идентичен формату вывода changefeed для таблицы.

Changefeeds с фильтрацией и агрегационными запросами

Как и любая команда ReQL, changes интегрируется с остальной частью языка запросов. Вы можете вызвать changes после большинства команд, которые преобразуют или выбирают данные:

  • filter
  • get_all
  • map
  • pluck
  • between
  • union
  • min
  • max
  • order_by.limit

Вы также можете использовать changes перед любой командой, которая работает с последовательностью документов, при условии, что эта команда не потребляет всю последовательность. (Например, count и orderBy не могут стоять после команды changes.)

Предположим, у вас есть приложение чата с несколькими клиентами, которые публикуют сообщения в разные чаты. Вы можете создать потоки, которые подписываются на сообщения, опубликованные в определенном чате:

r.table('messages').filter(r.row['room_id'] == ROOM_ID).changes().run(conn)

Вы также можете использовать более сложные выражения. Допустим, у вас есть таблица scores, которая содержит последние результаты игры для каждого пользователя вашей игры. Вы можете создать поток всех игр, в которых пользователь превзошел свой предыдущий результат, и получить только новое значение:

r.table('scores').changes().filter(
    lambda change: change['new_val']['score'] > change['old_val']['score']
)['new_val'].run(conn)

Ограничения и замечания по цепочке с changefeeds:

  • min, max и order_by должны использоваться с индексами.
  • order_by требует limit; ни одна команда не работает сама по себе.
  • order_by должен использоваться с вторичным индексом или первичным индексом; он не может использоваться с нес индексированным полем.
  • Вы не можете использовать changefeeds после concat_map или других преобразований, результаты которых не могут быть отправлены на фрагменты.
  • Вы не можете применить filter после order_by.limit в changefeed.
  • Преобразования применяются до вычисления изменений.

Включение изменений состояния

Необязательный аргумент include_states для changes позволяет получать дополнительные «статусные» документы в потоках changefeed. Это может позволить вашему приложению различать начальные значения, возвращаемые в начале потока, и последующие изменения. Обратитесь к документации API changes для получения полного объяснения и примера.

Включение начальных значений

Указав True для необязательного аргумента include_initial, поток changefeed начнётся с текущего содержимого таблицы или выбора, за которым он следит. Начальные результаты будут иметь поля new_val, но не поля old_val, поэтому их легко отличить от событий изменений.

Если начальный результат для документа был отправлен, и в этом документе произошли изменения, которые сместили бы его в неотправленную часть набора результатов (например, changefeed отслеживает 100 лучших авторов, первые 50 были отправлены, а автор 48 стал автором 52), будет отправлено уведомление «uninitial», с полем old_val, но без поля new_val. Это отличается от события удаления, в котором new_val имеет значение None. (В примере с 100 лучшими авторами это может указывать на удаление автора или его выбытие из 100 лучших.)

Если вы укажете True для include_states и include_initial, поток changefeed начнётся со статусного документа {'state': 'initializing'}, за которым следуют начальные значения. Статусный документ {'state': 'ready'} будет отправлен, когда все начальные значения будут отправлены.

Включение типов результатов

Необязательный аргумент include_types добавляет третье поле, type, к каждому отправленному результату. Строковые значения для type в основном самодокументированы:

  • add: новое значение, добавленное в набор результатов.
  • remove: старое значение, удалённое из набора результатов.
  • change: существующее значение изменено в наборе результатов.
  • initial: уведомление о начальном значении.
  • uninitial: уведомление об «uninitial» значении.
  • state: статусный документ от include_states.

Включение поля type может упростить код, обрабатывающий разные случаи результатов changefeed.

Обработка задержек

В зависимости от того, насколько быстро ваше приложение вносит изменения в отслеживаемые данные и как быстро оно обрабатывает уведомления об изменениях, возможно, что между вызовами команды changes произойдёт несколько изменений. Вы можете контролировать, что происходит в этом случае, с помощью необязательного аргумента squash.

По умолчанию, если между вызовами changes произойдёт более одного изменения, ваше приложение получит один объект изменения, чьё new_val будет включать все изменения в данных. Предположим, что три обновления произошли в отслеживаемом документе между чтениями change:

Изменение Данные
Начальное состояние (old_val ) { имя: «Фред», администратор: true }
update({имя: «Джордж»}) { имя: «Джордж», администратор: true }
update({администратор: false}) { имя: «Джордж», администратор: false }
update({имя: «Джей»}) { имя: «Джей», администратор: false }
new_val { имя: «Джей», администратор: false }

По умолчанию ваше приложение получит объект в том виде, в котором он существовал в базе данных после последнего изменения. Два предыдущих обновления будут «сжаты» в третье.

Если вы хотите получить все изменения, включая промежуточные состояния, вы можете сделать это, передав squash: False . Сервер будет буферизовать до 100 000 изменений. (Это число можно изменить с помощью необязательного аргумента changefeed_queue_size.)

Третий вариант — указать, сколько секунд ждать между сжатиями. Передача squash: 5 в команду changes сообщает RethinkDB о сжатии изменений каждые пять секунд. В зависимости от вашего варианта использования это может уменьшить нагрузку на сервер. Число, переданное в squash может быть дробным. Обратите внимание, что запрошенный интервал не гарантируется, а представляет собой оптимальное значение.

Примечание: Changefeeds игнорируют флаг read_mode для run, и всегда ведут себя так, как будто он установлен в single (т. е., значения, которые они возвращают, находятся в памяти на первичном репликаторе, но не обязательно записаны на диск). Для получения дополнительной информации прочитайте Гарантии согласованности.

Масштабирование

Changefeeds хорошо масштабируются, хотя они создают дополнительные внутрикластерные сообщения пропорционально числу серверов с открытыми подключениями к потоку при каждом записи. Это можно смягчить, запустив сервер прокси RethinkDB (опция запуска rethinkdb proxy); подробнее см. Запуск прокси-узла.

Поскольку changefeeds являются однонаправленными и не возвращают подтверждений от клиентов, они не могут гарантировать доставку. Если вам требуется обновление в реальном времени с гарантией доставки, рассмотрите модель, которая распределяет данные клиентам через посредника сообщений, такой как RabbitMQ.

Подробнее

  • Ссылка на API команды changes
  • Введение в ReQL
  • Типы данных ReQL

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

Spec-Zone.ru

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