Changefeeds в RethinkDB
Changefeeds — основа реального времени в RethinkDB.
Они позволяют клиентам получать изменения в таблице, отдельном документе или даже результатах определённого запроса по мере их возникновения. Почти любой запрос ReQL можно преобразовать в changefeed.
Основные принципы работы
Подпишитесь на поток, вызвав changes для таблицы:
r.table('users').changes().run(conn, function(err, cursor) {
cursor.each(console.log);
})
Команда changes возвращает курсор (как команды table или filter). Вы можете перебирать его содержимое с помощью ReQL. В отличие от других курсоров, вывод changes бесконечен: курсор будет ожидать появления новых элементов. Каждый раз, когда вы вносите изменения в таблицу или документ, за которыми следит поток changes, в курсор будет возвращено новое значение. Например, если вы вставляете пользователя {id: 1, name: Slava, age: 31} в таблицу users, RethinkDB разместит этот документ в changefeeds, подписанных на 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 присутствуют.
Точечные (изменения одного документа) changefeeds
Точечный changefeed возвращает изменения в одном документе в таблице, а не в таблице в целом.
r.table('users').get(100).changes().run(conn, callback);
Формат вывода точечного changefeed идентичен changefeed таблицы.
Changefeeds с фильтрацией и агрегационными запросами
Как и любая команда ReQL, changes интегрируется с остальной частью языка запросов. Вы можете вызвать changes после большинства команд, преобразующих или выбирающих данные:
Вы также можете цеплять changes перед любой командой, которая работает со последовательностью документов, если эта команда не потребляет всю последовательность. (Например, count и orderBy не могут следовать за командой changes.)
Предположим, у вас есть чат-приложение с несколькими клиентами, отправляющими сообщения в разные чаты. Вы можете создать потоки, которые будут подписаны на сообщения, отправленные в определённый чат:
r.table('messages').filter(
r.row('room_id').eq(ROOM_ID)
).changes().run(conn, callback)
Вы также можете использовать более сложные выражения. Допустим, у вас есть таблица scores, содержащая последние игровые результаты каждого пользователя. Вы можете создать поток всех игр, в которых пользователь улучшил свой предыдущий результат и получить только новое значение:
r.table('scores').changes().filter(
r.row('new_val')('score').gt(r.row('old_val')('score'))
)('new_val').run(conn, callback)
Есть некоторые ограничения и замечания по цепочке changefeeds.
-
min,maxиorderByдолжны использоваться с индексами. -
orderByтребуетlimit; ни одна команда не работает сама по себе. -
orderByнеобходимо использовать с вторичным индексом или первичным индексом; его нельзя использовать с непроиндексированным полем. - Вы не можете использовать changefeeds после concatMap или других преобразований, результаты которых не могут быть переданы на фрагменты.
- Вы не можете применить
filterпослеorderBy.limitв changefeed. - Преобразования применяются до расчёта изменений.
Включение изменений состояния
Необязательный аргумент includeStates для changes позволяет получать дополнительные документы «статуса» в потоках changefeed. Это позволяет вашему приложению различать начальные значения, возвращённые в начале потока, и последующие изменения. Обратитесь к документации API changes для получения полного объяснения и примера.
Включение начальных значений
Указывая true для необязательного аргумента includeInitial, поток changefeed начнется с текущего содержимого контролируемой таблицы или выборки. Результаты первоначальных значений будут иметь new_val поля, но не old_val поля, что упрощает их различие от событий изменений.
Если для документа было отправлено начальное значение, и внесены изменения, которые поместят его в неотправленную часть набора результатов (например, changefeed отслеживает 100 лучших участников, первые 50 были отправлены, а участник 48 стал участником 52), будет отправлено уведомление «uninitial» с полем old_val, но без поля new_val . Это отличается от события удаления, которое будет иметь new_val как null. (В примере с 100 лучшими участниками это может означать удаление участника или его выпадение из 100 лучших).
Если вы укажете true как для includeStates, так и для includeInitial, поток changefeed начнётся с документа состояния {state: 'initializing'}, а затем с начальных значений. Документ состояния {state: 'ready'} будет отправлен, когда все начальные значения будут отправлены.
Включение типов результатов
Необязательный аргумент includeTypes добавляет третье поле, type, к каждому отправленному результату. Значения строк для type в значительной степени понятны из контекста:
-
add: новое значение добавлено в набор результатов. -
remove: старое значение удалено из набора результатов. -
change: существующее значение изменено в наборе результатов. -
initial: уведомление о начальном значении. -
uninitial: уведомление об uninitial значении. -
state: документ статуса отincludeStates.
Включение поля type может упростить код, обрабатывающий различные случаи результатов changefeed.
Обработка задержки
В зависимости от скорости, с которой ваше приложение вносит изменения в контролируемые данные и скорости обработки уведомлений о изменениях, возможно, что между вызовами команды changes произойдёт больше одного изменения. Вы можете контролировать, что произойдёт в этом случае, с помощью необязательного аргумента squash.
По умолчанию, если между вызовами changes произойдёт больше одного изменения, ваше приложение получит один объект изменения, у которого new_val будет содержать все изменения данных. Предположим, что три обновления произошли с отслеживаемым документом между чтениями change:
| Изменение | Данные |
|---|---|
Начальное состояние (old_val) | { имя: “Фред”, admin: true } |
| update({имя: “Джордж”}) | { имя: “Джордж”, admin: true } |
| update({admin: false}) | { имя: “Джордж”, admin: false } |
| update({имя: “Джей”}) | { имя: “Джей”, admin: false } |
new_val | { имя: “Джей”, admin: false } |
Ваше приложение по умолчанию получит объект, таким, каким он был в базе данных после последнего изменения. Два предыдущих обновления будут «объединены» в третье.
Если вы хотели получить все изменения, включая промежуточные состояния, вы можете сделать это, передав squash: false. Сервер будет буферизовать до 100 000 изменений. (Это число можно изменить с помощью необязательного аргумента changefeedQueueSize.)
Третий вариант — указать, сколько секунд ожидать между объединениями. Передача squash: 5 команде changes говорит RethinkDB объединять изменения вместе каждые пять секунд. В зависимости от вашего случая использования, это может уменьшить нагрузку на сервер. Число, переданное в squash, может быть дробным. Обратите внимание, что запрошенный интервал не гарантируется, а является лишь лучшей оценкой.
Примечание: Changefeeds игнорируют флаг read_mode для run, и всегда ведут себя так, как будто он установлен в single (т.е. возвращаемые значения находятся в памяти на первичной реплике, но не обязательно записаны на диск). Для получения более подробной информации см. Гарантии согласованности.
Масштабирование
Changefeeds хорошо масштабируются, хотя они создают дополнительные внутрикластерные сообщения пропорционально числу серверов с открытыми соединениями changefeed при каждом записи. Это можно смягчить, запустив прокси-сервер 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/javascript/