Spec-Zone.ru › RethinkDB ruby

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

Changefeeds лежат в основе реального времени функциональности RethinkDB.

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

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

Data Modeling Illustration

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

Подпишитесь на канал, вызвав 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 после большинства команд, которые преобразуют или выбирают данные:

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

Вы также можете использовать цепочку 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/

Spec-Zone.ru

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