Потоки изменений в RethinkDB
Изменения данных лежат в основе функций RethinkDB в реальном времени.
Они позволяют клиентам получать изменения в таблице, отдельном документе или даже результатах определенного запроса по мере их возникновения. Почти любой запрос ReQL можно преобразовать в changefeed.
Базовое использование
Подпишитесь на поток, вызвав 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 после большинства команд, которые преобразуют или выбирают данные:
Вы также можете использовать 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/