Асинхронные соединения
Некоторые драйверы RethinkDB поддерживают асинхронные подключения, интегрируясь с популярными асинхронными библиотеками. Это особенно полезно при работе с changefeeds и другими приложениями реального времени.
Благодаря своей событийно-ориентированной природе, JavaScript может легко выполнять запросы RethinkDB асинхронно. Официальные драйверы RethinkDB в настоящее время поддерживают интеграцию с EventMachine для Ruby и Tornado и Twisted для Python.
JavaScript
Для выполнения запросов RethinkDB асинхронно в JavaScript не требуются специальные процедуры или команды. Подробнее об использовании обратных вызовов и промисов с RethinkDB можно прочитать в документации к команде run.
Кроме того, курсоры и каналы RethinkDB реализуют интерфейс EventEmitter, совместимый с Node.js. Это позволяет вашему приложению настраивать обработчики событий для получения данных из запросов по мере их поступления.
Ruby с EventMachine
Драйвер RethinkDB для Ruby добавляет новую команду ReQL, em_run, предназначенную для работы с EventMachine. Кроме того, он предоставляет суперкласс RethinkDB::Handler, с методами, специфичными для событий (например, on_open, on_close), которые могут быть переопределены классом, определённым вашим приложением, и переданным в em_run.
Простой пример использования
Самый простой способ использования RethinkDB с EventMachine — передача блока в em_run. Если RethinkDB возвращает последовательность (включая поток), блок будет вызываться один раз для каждого элемента последовательности. В противном случае, блок будет вызываться только один раз со значениями возвращаемыми запросом.
Пример: Итерация по потоку
require 'eventmachine'
require 'rethinkdb'
include RethinkDB::Shortcuts
conn = r.connect(host: 'localhost', port: 28015)
EventMachine.run {
r.table('test').order_by(:index => 'id').em_run(conn) { |row|
# do something with returned row data
p row
}
}
Явное закрытие запроса
Команда em_run возвращает экземпляр QueryHandle. QueryHandle будет закрыт, когда все результаты будут получены или когда EventMachine перестанет работать. Вы можете явным образом закрыть его с помощью метода close.
EventMachine.run {
printed = 0
handle = r.table('test').order_by(:index => 'id').em_run(conn) { |row|
printed += 1
if printed > 3
handle.close
else
p row
end
}
}
Обработка ошибок
В указанном выше формате — с блоком, принимающим один аргумент — адаптер EventMachine RethinkDB будет выбрасывать ошибки обратно в ваше приложение для обработки в том же формате, как и при использовании RethinkDB без EventMachine. Если таблица test не существует в базе данных выше, вы получите стандартную ошибку ReqlRunTimeError:
RethinkDB::ReqlRunTimeError: Table `test.test` does not exist.
Backtrace:
r.table('test')
^^^^^^^^^^^^^^^
Вы также можете выбрать получение ошибок в блоке, приняв два аргумента.
EventMachine.run {
r.table('test').order_by(:index => 'id').em_run(conn) { |err, row|
if err
p [:err, err.to_s]
else
p [:row, row]
end
}
}
В этом формате блок получит nil в качестве первого аргумента, если ошибки нет. В случае ошибки, второй аргумент будет nil.
Использование RethinkDB::Handler
Для получения большего контроля, напишите класс, унаследованный от RethinkDB::Handler, переопределите методы обработки событий, а затем передайте экземпляр этого класса в em_run.
Пример: Итерация по потоку с использованием обработчика
require 'eventmachine'
require 'rethinkdb'
include RethinkDB::Shortcuts
conn = r.connect(host: 'localhost', port: 28015)
class Printer < RethinkDB::Handler
def on_open
p :open
end
def on_close
p :closed
end
def on_error(err)
p [:err, err.to_s]
end
def on_val(val)
p [:val, val]
end
end
EventMachine.run {
r.table('test').order_by(:index => 'id').em_run(conn, Printer)
}
# Sample output
:open
[:val, {"id"=>1}]
[:val, {"id"=>2}]
[:val, {"id"=>3}]
:closed
Различение типов данных
В дополнение к простому методу on_val, вы можете предоставить методы, специфичные для массивов, потоков и атомов.
class Printer < RethinkDB::Handler
def on_open
p :open
end
def on_close
p :closed
end
def on_error(err)
p [:err, err.to_s]
end
# Handle arrays
def on_array(array)
p [:array, array]
end
# Handle atoms
def on_atom(atom)
p [:atom, atom]
end
# Handle individual values received from streams
def on_stream_val(val)
p [:stream_val, val]
end
def on_val(val)
p [:val, val]
end
end
EventMachine.run {
r.table('test').order_by(:index => 'id').em_run(conn, Printer)
# print an array
r.expr([1, 2, 3]).em_run(conn, Printer)
# print a single row
r.table('test').get(1).em_run(conn, Printer)
}
# Sample output
:open
[:stream_val, {"id"=>0}]
[:stream_val, {"id"=>1}]
[:stream_val, {"id"=>2}]
:closed
:open
[:array, [1, 2, 3]]
:closed
:open
[:atom, {"id"=>0}]
:closed
Различные методы on_* обеспечивают взаимозаменяемость:
- массив будет обработан методом
on_array, если он определён; в противном случае он будет обработан методомon_atom. Если ни один из этих методов не определён, отдельные элементы массива будут обработаны методомon_stream_valили, если он не определён, методомon_val. - поток будет обработан методом
on_stream_valесли он определён; в противном случае он будет обработан методомon_val. - данные, которые не являются потоком, будут обработаны методом
on_atomесли он определён; в противном случае они будут обработаны методомon_val.
Таким образом, on_val выступает в качестве универсального обработчика для всех данных, которые не обрабатываются более специфичными методами.
Порядок вызова обратных вызовов в блоке EventMachine.run не гарантируется; в приведённом выше примере вывода [:array, [1, 2, 3]] мог быть напечатан первым.
Changefeeds
Changefeed обрабатывается как любой другой поток; когда вы передаёте блок в em_run, блок вызывается для каждого документа, полученного из канала. Если вы передаёте Handler, который определяет on_stream_val (или on_val), эти методы будут вызваны для каждого документа.
Кроме того, существуют методы, специфичные для changefeed:
-
on_initial_val: если changefeed возвращает начальные значения (include_initialуказан в качестве параметра к changes, эти значения будут переданы в этот метод. -
on_uninitial_val: changefeed, возвращающий начальные значения, также может возвращать значения «uninitial» для указания того, что документ, уже отправленный как начальное значение, был изменён (подробнее см. в документации кchanges); эти значения, если таковые имеются, будут переданы в этот метод. -
on_change: изменения будут переданы в этот метод. -
on_change_error: если канал включает документ, указывающий на ошибки, которые не приводят к прерыванию канала (например, уведомление о том, что сервер отбросил некоторые изменения), эти ошибки будут переданы в этот метод. -
on_state: канал может включать документы, указывающие на состояние потока; эти документы будут переданы в эту функцию, если она определена.
class FeedPrinter < RethinkDB::Handler
def on_open
p :open
end
def on_close
p :closed
end
def on_error(err)
p [:err, err.to_s]
end
def on_initial_val(val)
p [:initial, val]
end
def on_state(state)
p [:state, state]
end
def on_change(old, new)
p [:change, old, new]
end
end
# Subscribe to changes on the documents with the two lowest ids
EventMachine.run {
r.table('test').order_by(:index => 'id').limit(2).changes
.em_run(conn, FeedPrinter)
}
# Sample output
:open
[:state, "initializing"]
[:initial_val, {"id"=>1}]
[:initial_val, {"id"=>0}]
[:state, "ready"]
# Execute: r.table('test').insert({id: 0.5}).run(conn)
[:change, {"id"=>1}, {"id"=>0.5}]
# Execute: r.table_drop('test').run(conn)
[:err, "Changefeed aborted (table unavailable).\nBacktrace..."]
:closed
Использование одного обработчика с несколькими запросами
Вы можете зарегистрировать несколько запросов с одним экземпляром Handler. Если вы определите методы Handler с дополнительным аргументом (два аргумента вместо одного или один аргумент вместо нуля), этот аргумент получит соответствующий экземпляр QueryHandle.
class MultiQueryPrinter < RethinkDB::Handler
def on_open(qh)
p [:open, names[qh]]
end
def on_close(qh)
p [:close, names[qh]]
EventMachine.stop if @closed == 2
end
def on_val(val, qh)
p [:val, val, names[qh]]
end
end
EventMachine.run {
printer = Printer.new
h1 = r.expr(1).em_run(conn, printer)
h2 = r.expr(2).em_run(conn, printer)
names = { h1 => "h1", h2 => "h2" }
}
# Sample output
[:open, "h1"]
[:val, 1, "h1"]
[:close, "h1"]
[:open, "h2"]
[:val, 2, "h2"]
[:close, "h2"]
Остановка обработчика
Если вы вызовите метод stop на экземпляре Handler, он прекратит обработку изменений, а открытые потоки, использующие этот обработчик, будут закрыты. Зарегистрированные с этим экземпляром запросы не будут прерваны, если они в настоящее время обрабатываются (например, пакетное запись), но закроются вместо выполнения после того, как обработчик будет остановлен.
Пример: Выведите первые пять изменений в таблицу. После остановки обработчика запрос changefeed будет закрыт при следующем изменении в таблице, а не вернёт значение.
class FeedPrinter < RethinkDB::Handler
def initialize(max)
@counter = max
stop if @counter <= 0
end
def on_open
# Once the changefeed is open, insert 10 rows
r.table('test').insert([{}] * 10).run(conn, noreply: true)
end
def on_val(val)
# Every time we print a change, decrement @counter and stop if we hit 0
p val
@counter -= 1
stop if @counter <= 0
end
end
EventMachine.run {
r.table('test').changes.em_run(conn, Printer.new(5))
}
Python с Tornado или Twisted
Драйвер RethinkDB для Python интегрируется как с Tornado web framework, так и с Twisted networking engine. Используя команду set_loop_type, вы можете выбрать модель цикла событий 'tornado' или 'twisted', возвращая соответственно объекты Tornado Future или Twisted Deferred.
Tornado
Основные примеры использования
Перед connect, используйте команду set_loop_type("tornado") для настройки RethinkDB на использование асинхронных циклов событий, совместимых с Tornado.
from rethinkdb import RethinkDB
from tornado import ioloop, gen
from tornado.concurrent import Future, chain_future
import functools
r = RethinkDB()
r.set_loop_type("tornado")
connection = r.connect(host='localhost', port=28015)
После выполнения set_loop_type, r.connect вернёт объект Tornado Future, как и r.run.
Пример: Простой пример использования
@gen.coroutine
def single_row(connection_future):
# Wait for the connection to be ready
connection = yield connection_future
# Insert some data
yield r.table('test').insert([{"id": 0}, {"id": 1}, {"id": 2}]).run(connection)
# Print the first row in the table
row = yield r.table('test').get(0).run(connection)
print(row)
# Output
{u'id': 0}
Пример: Использование курсора
@gen.coroutine
def use_cursor(connection_future):
# Wait for the connection to be ready
connection = yield connection_future
# Insert some data
yield r.table('test').insert([{"id": 0}, {"id": 1}, {"id": 2}]).run(connection)
# Print every row in the table.
cursor = yield r.table('test').order_by(index="id").run(connection)
while (yield cursor.fetch_next()):
item = yield cursor.next()
print(item)
# Output
{u'id': 0}
{u'id': 1}
{u'id': 2}
Обратите внимание, что перебор курсора должен выполняться с помощью while и fetch_next, а не с помощью цикла for x in cursor.
Обработка ошибок
Если во время асинхронной операции возникает ошибка, оператор yield выбросит исключение, как обычно. Это может произойти немедленно (например, вы можете обратиться к таблице, которой не существует), но ваше приложение может получить большое количество данных перед ошибкой (например, ваша сеть может быть прервана после установления соединения).
Одна ошибка, заслуживающая особого внимания. Если у вас есть корутина, настроенная на неограниченное потребление changefeed, и соединение закрывается, корутина столкнётся с ошибкой ReqlRuntimeError.
Пример: Переброшенные ошибки
@gen.coroutine
def bad_table(connection):
yield r.table('non_existent').run(connection)
Traceback (most recent call last):
... elided ...
rethinkdb.errors.ReqlRuntimeError: Table `test.non_existent` does not exist. in:
r.table('non_existent')
^^^^^^^^^^^^^^^^^^^^^^^
Пример: Перехват ошибок в корутине
@gen.coroutine
def catch_bad_table(connection):
try:
yield r.table('non_existent').run(connection)
except r.ReqlRuntimeError:
print("Saw error")
# Output
Saw error
Подписка на changefeeds
Асинхронный API базы данных позволяет вам одновременно обрабатывать несколько changefeeds путём планирования фоновых корутин. Рассмотрим, например, такой обработчик changefeed:
@gen.coroutine
def print_cfeed_data(connection_future, table):
connection = yield connection_future
feed = yield r.table(table).changes().run(connection)
while (yield feed.fetch_next()):
item = yield feed.next()
print(item)
Мы можем запланировать его в цикле событий Tornado с помощью этого кода:
ioloop.IOLoop.current().add_callback(print_cfeed_data, connection, table)
Теперь корутина будет выполняться в фоновом режиме, вызывая изменения. Когда мы изменим таблицу, изменения будут замечены.
Теперь рассмотрим более сложный пример.
class ChangefeedNoticer(object):
def __init__(self, connection):
self._connection = connection
self._sentinel = object()
self._cancel_future = Future()
@gen.coroutine
def print_cfeed_data(self, table):
feed = yield r.table(table).changes().run(self._connection)
self._feeds_ready[table].set_result(True)
while (yield feed.fetch_next()):
cursor = feed.next()
chain_future(self._cancel_future, cursor)
item = yield cursor
if item is self._sentinel:
return
print("Seen on table %s: %s" % (table, item))
@gen.coroutine
def table_write(self, table):
for i in range(10):
yield r.table(table).insert({'id': i}).run(self._connection)
@gen.coroutine
def exercise_changefeeds(self):
self._feeds_ready = {'a': Future(), 'b': Future()}
loop = ioloop.IOLoop.current()
loop.add_callback(self.print_cfeed_data, 'a')
loop.add_callback(self.print_cfeed_data, 'b')
yield self._feeds_ready
yield [self.table_write('a'), self.table_write('b')]
self._cancel_future.set_result(self._sentinel)
@classmethod
@gen.coroutine
def run(cls, connection_future):
connection = yield connection_future
if 'a' in (yield r.table_list().run(connection)):
yield r.table_drop('a').run(connection)
yield r.table_create('a').run(connection)
if 'b' in (yield r.table_list().run(connection)):
yield r.table_drop('b').run(connection)
yield r.table_create('b').run(connection)
noticer = cls(connection)
yield noticer.exercise_changefeeds()
# Output
Seen on table a: {u'old_val': None, u'new_val': {u'id': 0}}
Seen on table b: {u'old_val': None, u'new_val': {u'id': 0}}
Seen on table a: {u'old_val': None, u'new_val': {u'id': 1}}
Seen on table b: {u'old_val': None, u'new_val': {u'id': 1}}
Seen on table a: {u'old_val': None, u'new_val': {u'id': 2}}
Seen on table b: {u'old_val': None, u'new_val': {u'id': 2}}
Seen on table a: {u'old_val': None, u'new_val': {u'id': 3}}
Seen on table b: {u'old_val': None, u'new_val': {u'id': 3}}
Seen on table a: {u'old_val': None, u'new_val': {u'id': 4}}
Seen on table b: {u'old_val': None, u'new_val': {u'id': 4}}
Seen on table a: {u'old_val': None, u'new_val': {u'id': 5}}
Seen on table a: {u'old_val': None, u'new_val': {u'id': 6}}
Seen on table b: {u'old_val': None, u'new_val': {u'id': 5}}
Seen on table a: {u'old_val': None, u'new_val': {u'id': 7}}
Seen on table b: {u'old_val': None, u'new_val': {u'id': 6}}
Seen on table a: {u'old_val': None, u'new_val': {u'id': 8}}
Seen on table b: {u'old_val': None, u'new_val': {u'id': 7}}
Seen on table a: {u'old_val': None, u'new_val': {u'id': 9}}
Seen on table b: {u'old_val': None, u'new_val': {u'id': 8}}
Seen on table b: {u'old_val': None, u'new_val': {u'id': 9}}
Здесь мы следим за изменениями в нескольких таблицах одновременно. Мы одновременно записываем в таблицы и наблюдаем, как наши записи появляются в changefeeds. После записи 10 элементов в каждую таблицу мы отменяем changefeeds.
Twisted
Основные примеры использования
Перед connect, используйте команду set_loop_type("twisted"), чтобы настроить RethinkDB на использование асинхронных циклов событий, совместимых с реактором Twisted.
from rethinkdb import RethinkDB
from twisted.internet import reactor, defer
from twisted.internet.defer import inlineCallbacks, returnValue
r = RethinkDB()
r.set_loop_type('twisted')
connection = r.connect(host='localhost', port=28015)
После выполнения set_loop_type, r.connect вернёт объект Twisted Deferred, как и r.run.
Пример: Простое использование
@inlineCallbacks
def single_row(conn_deferred):
# Wait for the connection to be ready
conn = yield conn_deferred
# Insert some data
yield r.table('test').insert([{"id": 0}, {"id": 1}, {"id": 2}]).run(conn)
# Print the first row in the table
row = yield r.table('test').get(0).run(conn)
print(row)
# Output
{u'id': 0}
Пример: Использование курсора
@inlineCallbacks
def use_cursor(conn):
# Insert some data
yield r.table('test').insert([{"id": 0}, {"id": 1}, {"id": 2}]).run(conn)
# Print every row in the table.
cursor = yield r.table('test').order_by(index="id").run(conn)
while (yield cursor.fetch_next()):
item = yield cursor.next()
print(item)
# Output:
{u'id': 0}
{u'id': 1}
{u'id': 2}
Обратите внимание, что для итерации по курсору необходимо использовать while и fetch_next, а не цикл for x in cursor.
Обработка ошибок
Если во время асинхронной операции произошла ошибка, оператор yield выбросит исключение, как обычно. Это может произойти немедленно (например, вы можете обратиться к таблице, которой не существует), но ваше приложение может получить большое количество данных перед ошибкой (например, ваша сеть может быть прервана после установления соединения).
Одна ошибка заслуживает особого внимания. Если у вас есть задача, которая бесконечно потребляет изменение, и соединение закрывается, задача столкнётся с ReqlRuntimeError.
Пример: Переброс ошибок
@inlineCallbacks
def bad_table(conn):
yield r.table('non_existent').run(conn)
Unhandled error in Deferred:
Traceback (most recent call last):
Failure: rethinkdb.errors.ReqlOpFailedError: Table `test.non_existent` does not exist in:
r.table('non_existent')
^^^^^^^^^^^^^^^^^^^^^^^
Пример: Обработка ошибок во время выполнения
@inlineCallbacks
def catch_bad_table(conn):
try:
yield r.table('non_existent').run(conn)
except r.ReqlRuntimeError:
print("Saw error")
# Output
Saw error
Подписка на changefeeds
Асинхронный API базы данных позволяет обрабатывать несколько changefeeds одновременно, запуская несколько фоновых задач. В качестве примера рассмотрим обработчик changefeed:
@inlineCallbacks
def print_feed(conn, table):
feed = yield r.table(table).changes().run(conn)
while (yield feed.fetch_next()):
item = yield feed.next()
print("Seen on table %s: %s" % (table, str(item)))
Мы можем запланировать его в Twisted reactor с помощью этого кода:
reactor.callLater(0, print_cfeed_data, conn, table)
Теперь задача будет выполняться в фоновом режиме, выводя изменения. При изменении таблицы изменения будут замечены.
Теперь рассмотрим более крупный пример:
@inlineCallbacks
def print_feed(conn, table, ready, cancel):
def errback_feed(feed, err):
feed.close()
return err
feed = yield r.table(table).changes().run(conn)
cancel.addErrback(lambda err: errback_feed(feed, err))
ready.callback(None)
while (yield feed.fetch_next()):
item = yield feed.next()
print("Seen on table %s: %s" % (table, str(item)))
@inlineCallbacks
def table_write(conn, table):
for i in range(10):
yield r.table(table).insert({'id': i}).run(conn)
@inlineCallbacks
def notice_changes(conn, *tables):
# Reset the state of the tables on the server
if len(tables) > 0:
table_list = yield r.table_list().run(conn)
yield defer.DeferredList([r.table_drop(t).run(conn) for t in tables if t in table_list])
yield defer.DeferredList([r.table_create(t).run(conn) for t in tables])
readies = [defer.Deferred() for t in tables]
cancel = defer.Deferred()
feeds = [print_feed(conn, table, ready, cancel) for table, ready in zip(tables, readies)]
# Wait for the feeds to become ready
yield defer.gatherResults(readies)
yield defer.gatherResults([table_write(conn, table) for table in tables])
# Cancel the feeds and wait for them to exit
cancel.addErrback(lambda err: None)
cancel.cancel()
yield defer.DeferredList(feeds)
yield notice_changes(conn, 'a', 'b')
# Output
Seen on table b: {u'old_val': None, u'new_val': {u'id': 0}}
Seen on table a: {u'old_val': None, u'new_val': {u'id': 0}}
Seen on table b: {u'old_val': None, u'new_val': {u'id': 1}}
Seen on table a: {u'old_val': None, u'new_val': {u'id': 1}}
Seen on table b: {u'old_val': None, u'new_val': {u'id': 2}}
Seen on table a: {u'old_val': None, u'new_val': {u'id': 2}}
Seen on table b: {u'old_val': None, u'new_val': {u'id': 3}}
Seen on table a: {u'old_val': None, u'new_val': {u'id': 3}}
Seen on table b: {u'old_val': None, u'new_val': {u'id': 4}}
Seen on table a: {u'old_val': None, u'new_val': {u'id': 4}}
Seen on table b: {u'old_val': None, u'new_val': {u'id': 5}}
Seen on table a: {u'old_val': None, u'new_val': {u'id': 5}}
Seen on table b: {u'old_val': None, u'new_val': {u'id': 6}}
Seen on table a: {u'old_val': None, u'new_val': {u'id': 6}}
Seen on table b: {u'old_val': None, u'new_val': {u'id': 7}}
Seen on table a: {u'old_val': None, u'new_val': {u'id': 7}}
Seen on table b: {u'old_val': None, u'new_val': {u'id': 8}}
Seen on table a: {u'old_val': None, u'new_val': {u'id': 8}}
Seen on table a: {u'old_val': None, u'new_val': {u'id': 9}}
Seen on table b: {u'old_val': None, u'new_val': {u'id': 9}}
Здесь мы слушаем изменения на нескольких таблицах одновременно. Мы одновременно записываем в таблицы и наблюдаем, как наши записи появляются в changefeeds. Затем мы отменяем changefeeds после записи 10 элементов в каждую из таблиц.
© RethinkDB contributors
Licensed under the Creative Commons Attribution-ShareAlike 3.0 Unported License.
https://rethinkdb.com/docs/async-connections/