Потоки изменений в RethinkDB
Потоки изменений лежат в основе реального времени функциональности RethinkDB.
Они позволяют клиентам получать изменения в таблице, отдельном документе или даже результатах конкретного запроса по мере их происхождения. Практически любой запрос ReQL может быть преобразован в поток изменений.
Базовое использование
Подпишитесь на поток, вызвав changes для таблицы:
Cursor changeCursor = r.table("users").changes().run(conn);
for (Object change : changeCursor) {
System.out.println(change);
}
Команда changes возвращает курсор (как и команды table или filter). Вы можете итерироваться по его содержимому с помощью ReQL. В отличие от других курсоров, результат changes бесконечен: курсор будет ожидать, пока станут доступны новые элементы. Каждый раз, когда вы вносите изменения в таблицу или документ, за которыми следит поток changes, новый объект будет возвращен курсору. Например, если вы вставляете пользователя {"id": 1, "name": "Slava", "age": 31} в таблицу users, RethinkDB опубликует этот документ в потоки изменений, подписанные на users.
{
"old_val": null,
"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 присутствуют.
Точечные (изменение одного документа) потоки изменений
«Точечный» поток изменений возвращает изменения в одном документе внутри таблицы, а не в таблице в целом.
r.table("users").get(100).changes().run(conn);
Формат вывода точечного потока изменений идентичен потоку изменений таблицы.
Потоки изменений с фильтрацией и агрегирующими запросами
Как и любая команда ReQL, changes интегрируется с остальной частью языка запросов. Вы можете вызвать changes после большинства команд, которые преобразуют или выбирают данные:
Вы также можете цепочкой подключить changes перед любой командой, которая работает с последовательностью документов, при условии, что эта команда не потребляет всю последовательность. (Например, count и orderBy не могут следовать за командой changes).
Предположим, у вас есть чат-приложение с несколькими клиентами, которые публикуют сообщения в различные чаты. Вы можете создать потоки, которые подписываются на сообщения, отправленные в конкретную комнату:
r.table("messages").filter(
row -> row.g("room_id").eq(ROOM_ID)
).changes().run(conn);
Вы также можете использовать более сложные выражения. Допустим, у вас есть таблица scores, которая содержит последние результаты игры для каждого пользователя вашей игры. Вы можете создать поток всех игр, в которых пользователь улучшил свой предыдущий результат, и получить только новое значение:
r.table("scores").changes().filter(
change -> change.g("new_val").g("score").gt(change.g("old_val").g("score"))
).g("new_val").run(conn);
Существуют некоторые ограничения и оговорки относительно цепочек с потоками изменений.
-
min,maxиorderByдолжны использоваться с индексами. -
orderByтребуетlimit; ни одна из команд не работает сама по себе. -
orderByдолжен использоваться с вторичным индексом или первичным индексом; он не может использоваться с неиндексированным полем. - Вы не можете использовать потоки изменений после concatMap или других преобразований, результаты которых нельзя передать фрагментам.
- Вы не можете применить
filterпослеorderBy.limitв потоке изменений. - Преобразования применяются до того, как рассчитываются изменения.
Включение изменений состояния
Необязательный аргумент include_states для changes позволяет получать дополнительные документы «состояния» в потоках изменений. Это может позволить вашему приложению различать начальные значения, возвращаемые в начале потока, и последующие изменения. Прочитайте документацию API changes для полного объяснения и примера.
Включение начальных значений
Указав true в необязательном аргументе include_initial, поток изменений начнётся с текущего содержимого контролируемой таблицы или выборки. Первоначальные результаты будут иметь new_val поля, но не old_val поля, что легко позволяет их отличить от событий изменения.
Если для документа был отправлен первоначальный результат, и происходит изменение, которое переместит его в неотправленную часть набора результатов (например, поток изменений отслеживает 100 лучших авторов, первые 50 были отправлены, а автор 48 стал автором 52), будет отправлено уведомление «uninitial», с полем old_val, но без поля new_val. Это отличается от события удаления, для которого new_val будет иметь значение null. (В примере с 100 лучшими авторами это может указывать на то, что автор был удален или выпал из 100 лучших).
Если вы укажете true для include_states и include_initial, поток изменений начнётся с документа состояния {"state": "initializing"}, за которым последуют начальные значения. Документ состояния {"state": "ready"} будет отправлен, когда все начальные значения будут отправлены.
Включение типов результатов
Необязательный аргумент include_types добавляет третье поле, type, к каждому отправленному результату. Значения строк для type в значительной степени самоочевидны:
-
add: добавление нового значения в набор результатов. -
remove: удаление старого значения из набора результатов. -
change: изменение существующего значения в наборе результатов. -
initial: уведомление о начальном значении. -
uninitial: уведомление об uninitial значении. -
state: документ состояния отinclude_states.
Включение поля type может упростить код, обрабатывающий различные случаи результатов потока изменений.
Обработка задержки
В зависимости от скорости, с которой ваше приложение вносит изменения в отслеживаемые данные и насколько быстро оно обрабатывает уведомления о изменениях, возможно, что произойдёт более одного изменения между вызовами команды 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, может быть числом с плавающей точкой. Обратите внимание, что запрошенного интервала нет гарантии, но он выполняется как лучшее усилие.
Примечание: Потоки изменений игнорируют флаг read_mode для run, и всегда ведут себя так, как если бы он был установлен на single (то есть значения, которые они возвращают, находятся в памяти на первичной реплике, но, возможно, ещё не записаны на диск). Для получения более подробной информации см. Гарантии согласованности.
Масштабирование
Потоки изменений хорошо масштабируются, хотя они создают дополнительные сообщения внутри кластера пропорционально количеству серверов с открытыми подключениями к потоку на каждом записи. Это можно смягчить, запустив сервер прокси RethinkDB (опция запуска rethinkdb proxy); см. Запуск прокси-узла для получения подробностей.
Поскольку потоки изменений являются однонаправленными и не возвращают подтверждения от клиентов, они не могут гарантировать доставку. Если вам нужна реальная обновляемость с гарантиями доставки, рассмотрите модель, которая распределяет обновления клиентам через брокер сообщений, такой как RabbitMQ.
Подробнее
- Ссылка на API команды changes
- Введение в ReQL
- Типы данных ReQL
© RethinkDB contributors
Licensed under the Creative Commons Attribution-ShareAlike 3.0 Unported License.
https://rethinkdb.com/docs/changefeeds/java/