Интеграция 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
Это десериализует сообщение об изменении и красиво выведет его, вместе с кратким описанием типа изменения.
Дополнительные материалы
© RethinkDB contributors
Licensed under the Creative Commons Attribution-ShareAlike 3.0 Unported License.
https://rethinkdb.com/docs/rabbitmq/ruby/