Команда 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, но принципы применимы ко всем языкам.)
© RethinkDB contributors
Licensed under the Creative Commons Attribution-ShareAlike 3.0 Unported License.
https://rethinkdb.com/api/python/changes/