Spec-Zone.ru › RethinkDB ruby

Асинхронные соединения

Некоторые драйверы RethinkDB поддерживают асинхронные подключения, интегрируясь с популярными асинхронными библиотеками. Это особенно полезно при работе с changefeeds и другими приложениями реального времени.

Благодаря своей событийно-ориентированной природе, JavaScript может легко выполнять запросы RethinkDB асинхронно. Официальные драйверы RethinkDB в настоящее время поддерживают интеграцию с EventMachine для Ruby и Tornado и Twisted для Python.

  • JavaScript
  • Ruby с EventMachine
  • Python с Tornado или Twisted

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/

Spec-Zone.ru

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