Spec-Zone.ru › RethinkDB ruby

Команда 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, но принципы применимы ко всем языкам.)

Связанные команды

  • table

© RethinkDB contributors
Licensed under the Creative Commons Attribution-ShareAlike 3.0 Unported License.
https://rethinkdb.com/api/ruby/changes/

Spec-Zone.ru

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