Потоки изменений в RethinkDB
Changefeeds лежат в основе реального времени функциональности RethinkDB.
Они позволяют клиентам получать изменения в таблице, в отдельном документе или даже в результатах конкретного запроса по мере их возникновения. Почти любой запрос ReQL можно преобразовать в changefeed.
Базовое использование
Подпишитесь на канал, вызвав changes для таблицы:
r.table('users').changes.run(conn).each{|change| p(change)}
Команда changes возвращает курсор (как и команды table или filter). Вы можете итерировать его содержимое с помощью ReQL. В отличие от других курсоров, выход changes бесконечен: курсор будет блокироваться до тех пор, пока не появятся новые элементы. Каждый раз, когда вы вносите изменения в таблицу или документ, за которыми следит changes канал, в курсор будет возвращаться новый объект. Например, если вы вставляете пользователя {id: 1, name: Slava, age: 31} в таблицу users, RethinkDB опубликует этот документ в changefeeds, подписанных на users:
{
:old_val => nil,
: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{ |row|
row['room_id'].eq(ROOM_ID)
}.changes().run(conn)
Вы также можете использовать более сложные выражения. Допустим, у вас есть таблица scores, содержащая последние результаты игры для каждого пользователя вашей игры. Вы можете создать канал всех игр, в которых пользователь превзошёл свой предыдущий результат, и получить только новое значение:
r.table('scores').changes().filter{ |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 nil (в примере с 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) | { name: “Fred”, admin: true } |
| update({name: “George”}) | { name: “George”, admin: true } |
| update({admin: false}) | { name: “George”, admin: false } |
| update({name: “Jay”}) | { name: “Jay”, admin: false } |
new_val | { name: “Jay”, admin: 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/ruby/