Асинхронные соединения
Некоторые драйверы RethinkDB поддерживают асинхронные соединения, интегрируясь с популярными асинхронными библиотеками. Это особенно полезно при работе с изменениями в данных и других приложениях реального времени.
Благодаря своей событийно-ориентированной природе, 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]] может быть напечатан первым.
Изменения в данных
Изменения в данных обрабатываются как любой другой поток; при передаче блока в em_run, блок вызывается с каждым документом, полученным в потоке. Если вы передаёте Handler, который определяет on_stream_val (или on_val), эти методы будут вызваны с каждым документом.
Кроме того, существуют методы, специфичные для изменений в данных.
-
on_initial_val: если изменения в данных возвращают начальные значения (include_initialбыл указан в качестве параметра для changes), эти значения будут переданы в этот метод. -
on_uninitial_val: изменения в данных, которые возвращают начальные значения, могут также возвращать «не-начальные» значения, чтобы указать, что документ, уже отправленный как начальное значение, был изменён (см. документацию по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, он прекратит обработку изменений, и открытые потоки, использующие этот обработчик, будут закрыты. Зарегистрированные с этим экземпляром обработчика запросы не будут прерваны, если они в данный момент обрабатываются (например, пакетная запись), но закроются вместо выполнения после остановки обработчика.
Пример: Вывести первые пять изменений в таблице. После остановки обработчика запрос к потоку изменений будет закрыт при следующем изменении таблицы, а не возвращая значение.
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, так и с сетевым движком Twisted. Используя команду 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 сгенерирует исключение, как обычно. Это может произойти немедленно (например, вы можете обратиться к несуществующей таблице), но ваше приложение может получить большое количество данных до возникновения ошибки (например, ваше сетевое соединение может прерваться после установления соединения).
Одна ошибка, в частности, заслуживает внимания. Если у вас есть сопрограмма, настроенная на постоянное использование канала изменений, и соединение закрывается, сопрограмма столкнётся с 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
Подписка на изменения в данных
Асинхронный API базы данных позволяет обрабатывать несколько каналов изменений одновременно, планируя фоновые сопрограммы. В качестве примера рассмотрим этот обработчик каналов изменений:
@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}}
Здесь мы подписываемся на изменения в нескольких таблицах одновременно. Мы одновременно записываем данные в таблицы и наблюдаем появление наших записей в каналах изменений. Затем мы отменяем каналы изменений после записи 10 элементов в каждую таблицу.
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
Подписка на изменения
Асинхронный API базы данных позволяет обрабатывать несколько изменении одновременно, запуская несколько фоновых задач. В качестве примера рассмотрим обработчик изменений:
@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.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}}
Здесь мы прослушиваем изменения в нескольких таблицах одновременно. Мы одновременно записываем в таблицы и наблюдаем, как наши записи появляются в потоках изменений. Затем мы отменяем потоки изменений после записи 10 элементов в каждую из таблиц.
© RethinkDB contributors
Licensed under the Creative Commons Attribution-ShareAlike 3.0 Unported License.
https://rethinkdb.com/docs/async-connections/