Spec-Zone.ru › RethinkDB ruby

Интеграция RethinkDB с RabbitMQ

RethinkDB поддерживает changefeeds, которые позволяют вам подписываться на изменения в таблице. База данных отправляет эти изменения вам по мере их возникновения.

Это открывает возможность мгновенного уведомления клиентских приложений об изменениях в таблице. Для приложений в реальном времени это поведение push является важным.

RabbitMQ является естественным выбором для распространения уведомлений об изменениях. Он разработан для эффективной маршрутизации сообщений многим слушателям, и существуют библиотеки клиентов для большинства популярных языков. В этом руководстве мы используем топические обмены RabbitMQ. Топические обмены позволяют клиентам подписываться на интересующие их сообщения и игнорировать остальные.

Перед началом

  • Прочитайте краткое руководство (30 секунд)
  • Убедитесь, что у вас установлен RethinkDB для вашей платформы
  • Установите Bunny, библиотеку RabbitMQ для Ruby

Отправка изменений в RabbitMQ

Давайте напишем скрипт, который отслеживает изменения в сервере RethinkDB и отправляет их в RabbitMQ.

Сначала нам нужно настроить подключение к серверу RethinkDB:

require 'bunny'
require 'rethinkdb'
include RethinkDB::Shortcuts
require 'json'

rethink_conn = r.connect(:host => 'localhost', :port => 28015)

Далее, мы подключимся к серверу RabbitMQ с помощью Bunny:

rabbit_conn = Bunny.new(:host => 'localhost', :port => 5672).start
channel = rabbit_conn.create_channel

Каналы (channels) мультиплексируют одно TCP-соединение. Все операции RabbitMQ выполняются на канале, а не напрямую на соединении. Далее мы объявим топический обмен (topic exchange), чтобы иметь место для отправки наших уведомлений об изменениях:

exchange = channel.topic("rethinkdb", :durable => false)

Это утверждает, что топический обмен с именем «rethinkdb» существует и что он недолговечен (non-durable). Если обмен не существует, он будет создан. Если он существует, но имеет другие свойства, произойдет исключение. Недолговечность означает, что он не сохранится после перезагрузки RabbitMQ (это значение по умолчанию).

В этом руководстве мы предположим, что сервер RethinkDB имеет базу данных с именем «change_example» и таблицу с именем «mytable». Вот запрос, который отслеживает изменения:

table_changes = r.db('change_example').table('mytable').changes

Вывод запроса changes соответствует следующему протоколу:

  • Если old_val равно nil, то new_val содержит только что созданный документ.
  • Если new_val равно nil, то old_val содержит удаленный документ.
  • В противном случае, документ был обновлён с new_val до old_val.

Теперь мы можем напрямую передавать наши изменения в Rabbit:

table_changes.run(rethink_conn).each do |change|
  routing_key = "mytable.#{type_of_change change}"
  exchange.publish(change.to_json, :routing_key => routing_key)
end

table_changes.run будет ожидать, пока не произойдёт изменение, после чего мы отправим его в обменник. routing_key — это тема, по которой мы будем отправлять. В этом примере у нас три разные темы: mytable.create, mytable.update, и mytable.delete. Каждая тема содержит только изменения соответствующего типа. Функция type_of_change выполняет эту сопоставление согласно описанному выше протоколу.

Приём сообщений из RabbitMQ

Слушатель — это другая сторона взаимодействия: он подключается к RabbitMQ, регистрируется для получения уведомлений о сообщениях, которые его интересуют, и выполняет какое-либо действие при получении сообщения.

Как и прежде, нам нужно создать соединение и канал RabbitMQ, а также убедиться, что обменник существует:

require 'bunny'
require 'json'

rabbit_conn = Bunny.new(:host => 'localhost', :port => 5672).start
channel = rabbit_conn.create_channel
exchange = channel.topic("rethinkdb", :durable => false)

В отличие от скрипта, который отправляет данные в Rabbit, для прослушивания нам нужна очередь. Очереди — это по сути почтовые ящики. Вы подключаетесь к обмену и подписываете очередь на различные темы из этого обмена:

queue = channel.queue('', :exclusive => true)

Вы можете присвоить очереди имя, если хотите, но поскольку мы передали пустую строку в queue, она создаст для нас случайное имя.

Теперь нам нужно «связать» очередь с темами, которые нас интересуют. Другие слушатели могут подписаться на ту же тему, и Rabbit скопирует сообщение для каждой очереди. Здесь мы просто будем подписываться на все события из «mytable»:

queue.bind(exchange, :routing_key => 'mytable.*')

Наконец, для прослушивания очереди мы используем метод queue.subscribe. Аналогично курсору changefeed из RethinkDB, subscribe будет ожидать, пока в очереди не появится сообщение.

queue.subscribe(:block => true) do |delivery_info, metadata, payload|
  change = JSON.parse(payload)
  tablename, change_type = delivery_info.routing_key.split('.')

  puts tablename, 'got a change of type:', change_type
  puts JSON.pretty_generate(change)
end

Это десериализует сообщение об изменении и красиво выведет его, вместе с кратким описанием типа изменения.

Дополнительные материалы

  • Полный исходный код этого руководства
  • Подробное описание модели RabbitMQ
  • Changefeeds RethinkDB

© RethinkDB contributors
Licensed under the Creative Commons Attribution-ShareAlike 3.0 Unported License.
https://rethinkdb.com/docs/rabbitmq/ruby/

Spec-Zone.ru

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