Spec-Zone.ru › RethinkDB java

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

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

  • table

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

Spec-Zone.ru

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