Spec-Zone.ru › RethinkDB javascript

Команда ReQL: changes

Синтаксис команды

stream.changes([options]) → stream
singleSelection.changes([options]) → stream

Описание

Преобразуйте запрос в поток изменений (changefeed), бесконечный поток объектов, представляющих изменения в результатах запроса по мере их возникновения. Поток изменений может возвращать изменения в таблице или в отдельном документе (потоке изменений «точки»). Перед командой changes можно использовать такие команды, как filter или map, для преобразования или фильтрации вывода, и многие команды, работающие с последовательностями, могут быть объединены после changes.

Существует шесть необязательных аргументов для changes.

  • squash: Управляет тем, как уведомления об изменениях группируются в пакеты. Допустимые значения — true, false и числовое значение:
    • true: Когда несколько изменений одного и того же документа возникают до отправки пакета уведомлений, изменения «сжимаются» в одно изменение. Клиент получает уведомление, которое полностью синхронизирует его с сервером.
    • false: Все изменения отправляются клиенту в неизменённом виде. Это значение по умолчанию.
    • n: Числовое значение (с плавающей точкой). Аналогично true, но сервер подождет n секунд для ответа, чтобы как можно больше изменений объединить, уменьшив сетевой трафик. Первый пакет всегда возвращается немедленно.
  • changefeedQueueSize: количество изменений, которые сервер будет буферизовать между чтением клиентом, прежде чем начать отбрасывать изменения и генерировать ошибку (по умолчанию: 100 000).
  • includeInitial: если true, поток изменений начнётся с текущего содержимого таблицы или выбора, за которым ведётся наблюдение. Эти начальные результаты будут иметь поля new_val, но не поля old_val. Начальные результаты могут быть перемешаны с фактическими изменениями, при условии, что для изменённого документа уже предоставлен начальный результат. Если для документа был отправлен начальный результат, и в этом документе произошли изменения, которые перенесли его в неотправленную часть набора результатов (например, поток изменений отслеживает 100 лучших авторов, первые 50 уже отправлены, и автор 48 стал автором 52), будет отправлено уведомление «uninitial», содержащее поле old_val, но не поле new_val.
  • includeStates: если true, поток изменений будет содержать специальные документы состояния, состоящие из поля state и строки, указывающей на изменение состояния потока. Эти документы могут появляться в любом месте потока между уведомлениями, описанными ниже. Если includeStates имеет значение false (значение по умолчанию), документы состояния не будут отправляться.
  • includeOffsets: если true, поток изменений на orderBy.limit потоке изменений будет содержать поля old_offset и new_offset в документах состояния, которые включают old_val и new_val. Это позволяет приложениям поддерживать упорядоченный список набора результатов потока. Если old_offset установлено и не равно null, элемент в позиции old_offset удаляется; если new_offset установлено и не равно null, то new_val вставляется в позицию new_offset. Установка includeOffsets в true для потока изменений, который его не поддерживает, приведёт к ошибке.
  • includeTypes: если true, каждый результат в потоке изменений будет содержать поле type со строкой, указывающей тип изменения, представленного результатом: add, remove, change, initial, uninitial, state. По умолчанию — false.

В настоящее время существуют два состояния:

  • {state: 'initializing'} указывает, что следующие документы представляют начальные значения в потоке, а не изменения. Это будет первый документ потока, возвращающий начальные значения.
  • {state: 'ready'} указывает, что следующие документы представляют изменения. Это будет первый документ потока, который не возвращает начальные значения; в противном случае он указывает, что начальные значения уже все отправлены.

Начиная с RethinkDB 2.2, документы состояния будут только отправляться, если параметр includeStates имеет значение true, даже в потоках изменений «точки». Начальные значения будут отправляться только если includeInitial имеет значение true. Если includeStates имеет значение true и includeInitial ложно, первый документ в потоке будет {state: 'ready'}.

Если таблица станет недоступной, поток изменений будет отключен, и драйвер выбросит исключение времени выполнения.

Уведомления о изменениях имеют вид объекта с двумя полями:

{
    "old_val": <document before change>,
    "new_val": <document after change>
}

Когда includeTypes имеет значение true, будет три поля:

{
    "old_val": <document before change>,
    "new_val": <document after change>,
    "type": <result type>
}

Когда документ удаляется, new_val будет null; когда документ вставляется, old_val будет null.

Определённые команды преобразования документов могут быть объединены перед потоками изменений. Дополнительную информацию см. в обсуждении потоков изменений в документации по «Языку запросов».

Примечание: Потоки изменений игнорируют флаг read_mode в run, и всегда ведут себя так, как будто он установлен в single (т. е. возвращаемые ими значения находятся в памяти на первичной реплике, но не обязательно были записаны на диск). Подробнее см. Гарантии согласованности.

Сервер будет буферизовать до changefeedQueueSize элементов (по умолчанию 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, function(err, cursor) {
  cursor.each(console.log);
});

По мере выполнения этих запросов во втором клиенте, первый клиент будет получать и выводить следующие объекты:

> r.table('games').insert({id: 1}).run(conn, callback);
{old_val: null, new_val: {id: 1}}

> r.table('games').get(1).update({player1: 'Bob'}).run(conn, callback);
{old_val: {id: 1}, new_val: {id: 1, player1: 'Bob'}}

> r.table('games').get(1).replace({id: 1, player1: 'Bob', player2: 'Alice'}).run(conn, callback);
{old_val: {id: 1, player1: 'Bob'},
 new_val: {id: 1, player1: 'Bob', player2: 'Alice'}}

> r.table('games').get(1).delete().run(conn, callback)
{old_val: {id: 1, player1: 'Bob', player2: 'Alice'}, new_val: null}

> r.tableDrop('games').run(conn, callback);
ReqlRuntimeError: Changefeed aborted (table unavailable)

Пример: Вернуть все изменения, которые увеличивают счёт игрока.

r.table('test').changes().filter(
  r.row('new_val')('score').gt(r.row('old_val')('score'))
).run(conn, callback)

Пример: Вернуть все изменения в счёте конкретного игрока, которые увеличивают его значение свыше 10.

r.table('test').get(1).filter(r.row('score').gt(10)).changes().run(conn, callback)

Пример: Вернуть все вставки в таблицу.

r.table('test').changes().filter(r.row('old_val').eq(null)).run(conn, callback)

Пример: Вернуть все изменения в игре 1, с уведомлениями о состоянии и начальными значениями.

r.table('games').get(1).changes({includeInitial: true, includeStates: true}).run(conn, callback);
// 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').orderBy(
    { index: r.desc('score') }
).limit(10).changes().run(conn, callback);

Пример: Поддерживать состояние массива на основе потока изменений.

r.table('data').changes(
    {includeInitial: true, includeOffsets: true}
).run(conn, function (err, change) {
    // delete item at old_offset before inserting at new_offset
    if (change.old_offset != null) {
        myArray.splice(change.old_offset, 1);
    }
    if (change.new_offset != null) {
        myArray.splice(change.new_offset, 0, change.new_val);
    }
});

(Это упрощённая реализация; для более сложной обработки см. функцию applyChange в исходном коде Horizon client/src/ast.js).

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

  • table

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

Spec-Zone.ru

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