Интеграция RethinkDB с RabbitMQ
RethinkDB поддерживает изменения в данных, которые позволяют вам подписываться на изменения в таблице. База данных отправляет эти изменения по мере их возникновения.
Это открывает возможность немедленного уведомления клиентских приложений при изменении таблицы. Для приложений в реальном времени такое поведение push-сообщений имеет важное значение.
RabbitMQ — естественный выбор для распространения уведомлений об изменениях. Он разработан для эффективного маршрутизации сообщений к множеству слушателей, и существуют библиотеки клиентов для большинства популярных языков программирования. В этом учебнике мы используем топические обмены RabbitMQ. Топические обмены позволяют клиентам подписываться на интересующие их сообщения и игнорировать остальные.
Перед началом
- Прочитайте краткое руководство (30 секунд)
- Убедитесь, что у вас установлена RethinkDB для вашей платформы
- Установите pika, библиотеку RabbitMQ для Python
Отправка изменений в RabbitMQ
Давайте напишем скрипт, который отслеживает изменения на сервере RethinkDB и отправляет их в RabbitMQ.
Сначала нам нужно настроить соединение с сервером RethinkDB:
from rethinkdb import RethinkDB
import pika
import json
r = RethinkDB()
rethink_conn = r.connect(host='localhost', port=28015)
Далее, мы подключимся к серверу RabbitMQ с помощью pika:
rabbit_conn = pika.BlockingConnection(
pika.ConnectionParameters(host='localhost', port=5672)
)
channel = rabbit_conn.channel()
Каналы мультиплексируют одно TCP-соединение. Все операции RabbitMQ выполняются в канале, а не непосредственно в соединении. Далее, мы объявим топический обмен, чтобы у нас было место для отправки уведомлений об изменениях:
channel.exchange_declare('rethinkdb', exchange_type='topic', durable=False)
Это гарантирует, что топический обмен с именем «rethinkdb» существует и что он не является долговременным. Если обмен не существует, он будет создан. Если он существует и имеет другие свойства, произойдет исключение. Недолговременность означает, что он не будет сохраняться после перезапуска RabbitMQ (это значение по умолчанию).
Для этого учебника мы предположим, что на сервере RethinkDB есть база данных с именем «change_example» и таблица с именем «mytable». Вот запрос, который следит за изменениями:
table_changes = r.db('change_example').table('mytable').changes()
Вывод запроса changes соответствует следующему протоколу:
- Если
old_valявляетсяNone, тоnew_valсодержит вновь созданный документ. - Если
new_valявляетсяNone, тоold_valсодержит удаленный документ. - В противном случае документ был обновлен с
new_valдоold_val
Теперь мы можем напрямую подключать наши изменения к Rabbit:
for change in table_changes.run(rethink_conn):
routing_key = 'mytable.' + type_of_change(change)
channel.basic_publish(exchange, routing_key, json.dumps(change))
table_changes.run() будет блокироваться до тех пор, пока не произойдет изменение, в этот момент мы отправляем его в обмен. routing_key — это тема, на которой мы будем его отправлять. В этом примере у нас есть три разных темы: mytable.create, mytable.update, и mytable.delete. Каждая тема содержит только изменения соответствующего типа. Функция type_of_change выполняет эту сопоставление, используя описанный выше протокол.
Прием сообщений RabbitMQ
Слушатель — другая сторона взаимодействия: он подключается к RabbitMQ, регистрируется для получения уведомлений об интересующих его сообщениях и выполняет какие-либо действия при получении сообщения.
Как и прежде, нам нужно создать соединение и канал RabbitMQ, и нам нужно убедиться, что обмен существует:
import pika
import json
rabbit_conn = pika.BlockingConnection(
pika.ConnectionParameters(host='localhost', port=5672)
channel = rabbit_conn.channel()
channel.exchange_declare('rethinkdb', exchange_type='topic', durable=False)
В отличие от скрипта, который отправляет данные в Rabbit, для прослушивания нам нужна очередь. Очереди — это фактически почтовые ящики. Вы идете в обмен и подписываете очередь на различные темы из этого обмена:
queue = channel.queue_declare(exclusive=True).method.queue
Вы можете присвоить очереди имя, если хотите, но поскольку мы не передали имя в queue_declare, она сгенерирует для нас случайное имя.
Теперь нам нужно «связать» очередь с темами, которые нас интересуют. Другие слушатели могут подписаться на ту же тему, и Rabbit скопирует сообщение для каждой очереди. Здесь мы просто будем держать это просто и связать с всеми событиями из «mytable»:
channel.queue_bind(queue, exchange='rethinkdb', routing_key='mytable.*')
Наконец, для прослушивания очереди мы используем генератор channel.consume. Подобно курсору changefeed от RethinkDB, consume будет блокироваться до тех пор, пока сообщение не поступит в очередь.
for method, properties, payload in channel.consume(queue):
change = json.loads(payload)
tablename, change_type = method.routing_key.split('.')
print tablename, 'got a change of type:', change_type
print json.dumps(change, indent=True, sort_keys=True)
Это позволит десериализовать сообщение об изменениях и красиво вывести его, а также краткое описание того, какой это тип изменения.
Дополнительные материалы
© RethinkDB contributors
Licensed under the Creative Commons Attribution-ShareAlike 3.0 Unported License.
https://rethinkdb.com/docs/rabbitmq/python/