Команда ReQL: changes
Синтаксис команды
stream.changes([options) → stream singleSelection.changes([options]) → stream
Описание
Преобразуйте запрос в изменение потока, бесконечный поток объектов, представляющих изменения в результатах запроса по мере их возникновения. Изменение потока может возвращать изменения в таблице или отдельном документе (потоке «точки»). Такие команды, как filter или map могут быть использованы перед командой changes, чтобы преобразовать или отфильтровать вывод, и многие команды, работающие с последовательностями, могут быть объединены после changes.
Существует шесть необязательных аргументов для changes.
-
squash: Управляет тем, как группируются уведомления об изменениях. Допустимые значения —true,falseи числовое значение:-
true: Когда несколько изменений в одном документе происходят перед отправкой пакета уведомлений, изменения «сжимаются» в одно изменение. Клиент получает уведомление, которое полностью синхронизирует его с сервером. -
false: Все изменения отправляются клиенту в первоначальном виде. Это значение по умолчанию. -
n: Числовое значение (с плавающей точкой). Аналогичноtrue, но сервер подождётnсекунд, чтобы ответить, чтобы как можно больше изменений объединились, уменьшая сетевой трафик. Первый пакет всегда возвращается немедленно.
-
-
changefeed_queue_size: количество изменений, которые сервер будет буферизовать между чтениями клиента, прежде чем начать отбрасывать изменения и генерировать ошибку (по умолчанию: 100 000). -
include_initial: еслиtrue, поток изменений начнётся с текущего содержимого таблицы или выбранного набора, за которым ведётся наблюдение. Эти начальные результаты будут иметь поляnew_val, но не поляold_val. Начальные результаты могут быть перемешаны с фактическими изменениями, при условии, что начальный результат для изменённого документа уже предоставлен. Если начальный результат для документа был отправлен, и в этом документе было произведено изменение, которое поместило бы его в часть набора результатов, которая ещё не была отправлена (например, поток изменений отслеживает 100 лучших участников, первые 50 уже отправлены, а 48-й участник стал 52-м), будет отправлено уведомление «неинициализированный», содержащее полеold_val, но без поляnew_val. -
include_states: еслиtrue, поток изменений будет включать специальные статусные документы, состоящие из поляstateи строки, указывающей на изменение состояния потока. Эти документы могут появиться в любой момент потока между документами уведомлений, описанными ниже. Еслиinclude_statesимеет значениеfalse(значение по умолчанию), статусные документы не отправляются. -
include_offsets: еслиtrue, поток изменений дляorder_by.limitпотока изменений будет включать поляold_offsetиnew_offsetв документах состояния, которые включаютold_valиnew_val. Это позволяет приложениям поддерживать упорядоченные списки набора результатов потока. Еслиold_offsetзадано и не равноnil, элемент по индексуold_offsetудаляется; еслиnew_offsetзадано и не равноnil, тоnew_valвставляется по индексуnew_offset. Установкаinclude_offsetsв значениеtrueдля потока изменений, который его не поддерживает, приведёт к ошибке. -
include_types: еслиtrue, каждый результат в потоке изменений будет включать полеtypeсо строкой, указывающей тип изменения, который представляет результат:add,remove,change,initial,uninitial,state. По умолчаниюfalse.
В настоящее время существуют два состояния:
-
{:state => "initializing"}указывает, что следующие документы представляют начальные значения в потоке, а не изменения. Это будет первый документ в потоке, возвращающем начальные значения. -
{:state => "ready"}указывает, что следующие документы представляют изменения. Это будет первый документ в потоке, который *не* возвращает начальные значения; в противном случае он укажет, что все начальные значения были отправлены.
Начиная с RethinkDB 2.2, документы состояния будут *только* отправляться, если параметр
include_statesимеет значениеtrue, даже в потоках изменений по точке. Начальные значения будут отправлены только еслиinclude_initialравноtrue. Еслиinclude_statesравноtrue, аinclude_initialложно, первый документ в потоке будет{:state => 'ready'}.
Если таблица становится недоступной, поток изменений отключается, и драйвер генерирует исключение во время выполнения.
Уведомления о изменениях имеют вид объекта с двумя полями:
{
:old_val => <document before change>,
:new_val => <document after change>
}
Когда include_types равно true, будет три поля:
{
:old_val => <document before change>,
:new_val => <document after change>,
:type => <result type>
}
Когда документ удаляется, new_val будет равно nil; когда документ вставляется, old_val будет равно nil.
Определённые команды преобразования документов могут быть объединены перед потоками изменений. Для получения дополнительной информации см. обсуждение потоков изменений в документации «Язык запросов».
Примечание: Потоки изменений игнорируют флаг
read_modeв значенииrun, и всегда ведут себя так, как будто он установлен вsingle(т. е. значения, которые они возвращают, находятся в памяти на первичной реплике, но не обязательно записаны на диск). Для получения более подробной информации см. Гарантии согласованности.
Сервер будет буферизовать до changefeed_queue_size элементов (по умолчанию 100 000). Если предел буфера достигнут, ранние изменения будут отброшены, и клиент получит объект вида {:error => "Changefeed cache over array size limit, skipped X elements."}, где X — количество пропущенных элементов.
Команды, которые работают с потоками (например, filter или map) обычно можно объединять после changes. Однако, поскольку поток, созданный changes, не имеет конца, команды, которым необходимо обработать весь поток до возврата (например, reduce или count), нельзя.
Пример: Подпишитесь на изменения в таблице.
Начните отслеживать поток изменений в одном клиенте:
r.table('games').changes().run(conn).each{|change| p(change)}
Когда эти запросы выполняются во втором клиенте, первый клиент получит и распечатает следующие объекты:
> r.table('games').insert({:id => 1}).run(conn)
{:old_val => nil, :new_val => {:id => 1}}
> r.table('games').get(1).update({:player1 => 'Bob'}).run(conn)
{:old_val => {:id => 1}, :new_val => {:id => 1, :player1 => 'Bob'}}
> r.table('games').get(1).replace({:id => 1, :player1 => 'Bob', :player2 => 'Alice'}).run(conn)
{:old_val => {:id => 1, :player1 => 'Bob'},
:new_val => {:id => 1, :player1 => 'Bob', :player2 => 'Alice'}}
> r.table('games').get(1).delete().run(conn)
{:old_val => {:id => 1, :player1 => 'Bob', :player2 => 'Alice'}, :new_val => nil}
> r.table_drop('games').run(conn)
ReqlRuntimeError: Changefeed aborted (table unavailable)
Пример: Верните все изменения, которые увеличивают очки игрока.
r.table('test').changes().filter{ |row|
row['new_val']['score'] > row['old_val']['score']
}.run(conn)
Пример: Верните все изменения очков конкретного игрока, которые увеличивают их более чем на 10.
r.table('test').get(1).filter { |row|
row['score'] > 10
}.run(conn)
Пример: Верните все вставки в таблицу.
r.table('test').changes().filter{ |row|
row['old_val'].eq(nil)
}.run(conn)
Пример: Верните все изменения игры 1 со статусом уведомлений и начальными значениями.
r.table('games').get(1).changes({:include_initial => true, :include_states => true}).run(conn)
# result returned on changefeed
{:state => "initializing"}
{:new_val => {:id => 1, :score => 12, :arena => "Hobbiton Field"}}
{:state => "ready"}
{
:old_val => {:id => 1, :score => 12, :arena => "Hobbiton Field"},
:new_val => {:id => 1, :score => 14, :arena => "Hobbiton Field"}
}
{
:old_val => {:id => 1, :score => 14, :arena => "Hobbiton Field"},
:new_val => {:id => 1, :score => 17, :arena => "Hobbiton Field", :winner => "Frodo"}
}
Пример: Верните все изменения в 10 лучших играх. Предполагается наличие вторичного индекса score на таблице games.
r.table('games').order_by(
{:index => r.desc('score')}
).limit(10).changes().run(conn)
Пример: Поддерживайте состояние массива на основе потока изменений.
EventMachine.run {
r.table('data').changes(
{:include_initial => true, :include_offsets => true}
).em_run(conn) { |change|
# delete item at old_offset before inserting at new_offset
my_array.delete_at(change.old_offset) if change.old_offset != nil
my.array.insert(change.new_offset, change.new_val) if change.new_offset != nil
}
}
(Это упрощённая реализация. Для получения более подробной информации о поддержке EventMachine в RethinkDB, обратитесь к документации em_run. Для более сложного примера см. функцию applyChange в исходном коде клиента Horizon client/src/ast.js, которая написана на JavaScript, но принципы применимы ко всем языкам.)
© RethinkDB contributors
Licensed under the Creative Commons Attribution-ShareAlike 3.0 Unported License.
https://rethinkdb.com/api/ruby/changes/