Команда 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).
© RethinkDB contributors
Licensed under the Creative Commons Attribution-ShareAlike 3.0 Unported License.
https://rethinkdb.com/api/javascript/changes/