Spec-Zone.ru › RethinkDB python

Команда 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 установлено и не равно None, элемент по индексу old_offset удаляется; если new_offset установлено и не равно None, то 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 будет None; когда документ вставляется, old_val будет None.

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

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

Сервер будет буферизовать до 100 000 элементов. Если лимит буфера достигнут, более ранние изменения будут отброшены, и клиент получит объект вида {"error": "Changefeed cache over array size limit, skipped X elements."}, где X — количество пропущенных элементов.

Команды, работающие с потоками (например, filter или map), обычно могут быть объединены после changes. Однако, поскольку поток, генерируемый changes, не имеет конца, команды, которым необходимо обработать весь поток перед возвратом (например, reduce или count), не могут.

Пример: Подписаться на изменения в таблице.

Начать мониторинг потока изменений в одном клиенте:

for change in r.table('games').changes().run(conn):
  print change

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

> r.table('games').insert({'id': 1}).run(conn)
{'old_val': None, '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': None}

> r.table_drop('games').run(conn)
ReqlRuntimeError: Changefeed aborted (table unavailable)

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

r.table('test').changes().filter(
  r.row['new_val']['score'] > r.row['old_val']['score']
).run(conn)

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

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

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

r.table('test').changes().filter(r.row['old_val'].eq(None)).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)

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

for change in r.table('data').changes(include_initial=True, include_offsets=True).run(conn):
    # delete item at old_offset before inserting at new_offset
    if change.old_offset != None:
        my_array.pop(change.old_offset)
    if change.new_offset != None:
        my_array.insert(change.new_offset, change.new_val);

(Это упрощённая реализация, и в продакшене вы должны использовать асинхронную модель событий, определённую с помощью set_loop_type. Для более сложного примера см. функцию applyChange в исходном коде Horizon’s client/src/ast.js; она написана на JavaScript, но принципы применимы ко всем языкам.)

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

  • table

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

Spec-Zone.ru

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