Команда ReQL: changes
Синтаксис команды
stream.changes() → stream singleSelection.changes() → stream
Описание
Преобразуйте запрос в поток изменений (changefeed), бесконечный поток объектов, представляющих изменения в результатах запроса по мере их возникновения. Поток изменений может возвращать изменения в таблице или отдельном документе (потоке изменений «точки»). Перед командой changes можно использовать такие команды, как filter или map, чтобы преобразовать или отфильтровать вывод, а многие команды, работающие со последовательностями, можно объединять после changes.
Вы можете указать один из шести необязательных аргументов через optArg.
-
squash: Управляет тем, как группируются уведомления об изменениях. Допустимые значения —true,falseи числовое значение:-
true: Если перед отправкой пакета уведомлений произошли несколько изменений в одном и том же документе, изменения «объединяются» в одно изменение. Клиент получает уведомление, которое полностью синхронизирует его с сервером. -
false: Все изменения будут отправлены клиенту без изменений. Это значение по умолчанию. -
n: Числовое значение (с плавающей точкой). Похоже наtrue, но сервер подождетnсекунды, чтобы ответить, чтобы объединить как можно больше изменений, уменьшая сетевой трафик. Первый пакет всегда будет возвращен немедленно.
-
-
changefeed_queue_size: количество изменений, которые сервер будет буферизовать между чтениями клиента, прежде чем начать отбрасывать изменения и генерировать ошибку (по умолчанию: 100 000). -
include_initial: Еслиtrue, поток изменений начнётся с текущего содержимого таблицы или выбранного набора данных, за которым ведётся наблюдение. Эти начальные результаты будут иметь поляnew_val, но не поляold_val. Начальные результаты могут быть перемешаны с фактическими изменениями, при условии, что для изменённого документа уже был предоставлен начальный результат. Если начальный результат для документа был отправлен, и в этом документе произошли изменения, которые переместили бы его в неотправленную часть набора результатов (например, поток изменений отслеживает 100 лучших участников, первые 50 отправлены, а участник 48 стал участником 52), будет отправлено уведомление «uninitial», содержащее полеold_val, но не полеnew_val. -
include_states: Еслиtrue, поток изменений будет включать специальные статусные документы, состоящие из поляstateи строки, указывающей на изменение состояния потока. Эти документы могут появляться в любом месте потока между уведомлениями о документах, описанных ниже. Еслиinclude_statesравноfalse(значение по умолчанию), статусные документы не будут отправляться. -
include_offsets: Еслиtrue, поток изменений на потоке измененийorderBy.limitбудет включать поляold_offsetиnew_offsetв статусных документах, которые включаютold_valиnew_val. Это позволяет приложениям поддерживать упорядоченные списки набора результатов потока. Еслиold_offsetустановлено и не равноnull, элемент вold_offsetудаляется; еслиnew_offsetустановлено и не равноnull, тоnew_valвставляется вnew_offset. Установкаinclude_offsetsвtrueна потоке изменений, который его не поддерживает, вызовет ошибку. -
include_types: Еслиtrue, каждый результат на потоке изменений будет содержать полеtypeсо строкой, указывающей тип изменения, которое представляет результат:add,remove,change,initial,uninitial,state. Значение по умолчанию —false.
В настоящее время существует два состояния:
-
{state: 'initializing'}указывает, что следующие документы представляют начальные значения в потоке, а не изменения. Это будет первый документ потока, возвращающего начальные значения. -
{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 будет null; когда документ вставляется, old_val будет null.
Некоторые команды преобразования документов могут быть объединены до потоков изменений. Дополнительную информацию можно найти в документации по «Языку запросов» в разделе обсуждения потоков изменений.
Примечание: Потоки изменений игнорируют флаг
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), не могут.
Пример: Подпишитесь на изменения в таблице.
Начните отслеживать поток изменений в одном клиенте:
Result<Object> changes = r.table("games").changes().run(conn);
for (Object change : changes) {
System.out.println(change);
}
По мере выполнения этих запросов во втором клиенте первый клиент будет получать и печатать следующие объекты:
r.table("games").insert(r.hashMap("id", 1)).run(conn);
{"old_val": null, "new_val": {"id": 1}}
r.table("games").get(1).update(r.hashMap("player1", "Bob")).run(conn);
{"old_val": {"id": 1}, "new_val": {"id": 1, "player1": "Bob"}}
r.table("games").get(1).replace(
r.hashMap("id", 1).with("player1", "Bob").with("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": null}
r.tableDrop("games").run(conn);
ReqlRuntimeError: Changefeed aborted (table unavailable)
Пример: Вернуть все изменения, которые увеличивают очки игрока.
r.table("test").changes().filter(
row -> row.g("new_val").g("score").gt(row.g("old_val").g("score"))
).run(conn);
Пример: Вернуть все изменения в очках конкретного игрока, которые увеличивают их более 10.
r.table("test").get(1).filter(row -> row.g("score").gt(10)).changes().run(conn);
Пример: Вернуть все вставки в таблицу.
r.table("test").changes().filter(
row -> row.g("old_val").eq(null)
).run(conn);
Пример: Вернуть все изменения в игре 1 с уведомлениями о состоянии и начальными значениями.
r.table("games").get(1).changes()
.optArg("include_initial", true).optArg("include_states", true).run(conn);
Возвращаемый результат в потоке изменений:
{"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().optArg("index", r.desc("score"))
.limit(10).changes().run(conn);
Пример: Поддерживать состояние списка на основе потока изменений.
Result<Object> changes = r.table("data").changes()
.optArg("include_initial", true)
.optArg("include_offsets", true)
.run((conn);
for (Object change : changes) {
// Delete item at old_offset before inserting at new_offset
if (change.old_offset != null) {
myList.remove(change.old_offset);
}
if (change.new_offset != null) {
myList.add(change.new_offset, change_new.val);
}
};
(Это упрощенная реализация. Для более сложного примера обратитесь к функции applyChange в исходном коде клиента Horizon client/src/ast.js; она написана на JavaScript, но принципы применимы ко всем языкам).
© RethinkDB contributors
Licensed under the Creative Commons Attribution-ShareAlike 3.0 Unported License.
https://rethinkdb.com/api/java/changes/