Spec-Zone.ru › RethinkDB javascript

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

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

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

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

Прежде чем начать

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

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

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

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

var r = require('rethinkdb');
var amqp = require('amqplib');

var rethinkConn = null;
var rabbitConn = null;
var channel = null;

var promise = r.connect({host: 'localhost', port: 28015}).then(function(conn){
   rethinkConn = conn;
})

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

promise = promise.then(function(){
    return amqp.connect('amqp://localhost:5672');
}).then(function(conn){
    rabbitConn = conn;
    return rabbitConn.createChannel();
}).then(function(ch){
    channel = ch;
})

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

promise = promise.then(function(){
    return channel.assertExchange('rethinkdb', 'topic', {durable: false});
})

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

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

var tableChanges = r.db('change_example').table('mytable').changes();

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

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

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

promise = promise.then(function(){
    return tableChanges.run(rethinkConn);
}).then(function(changeCursor){
    changeCursor.each(function(err, change){
        var routingKey = 'mytable.' + typeOfChange(change);
        var payload = new Buffer(JSON.stringify(change));
        channel.publish('rethinkdb', routingKey, payload);
    })
})

Каждый раз, когда происходит изменение, changeCursor.each отправит сообщение в обмен. routingKey — это тема, на которой мы его отправим. В этом примере у нас есть три разные темы: mytable.create, mytable.update, и mytable.delete. Каждая тема содержит только изменения соответствующего типа. Функция typeOfChange выполняет эту сопоставление, используя описанный выше протокол.

Прослушивание сообщений RabbitMQ

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

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

amqp = require('amqplib');

var rabbit_conn = null;
var channel = null;
var queue = null;

var promise = amqp.connect('amqp://localhost:5672').then(function(conn){
    rabbitConn = conn;
    return rabbitConn.createChannel();
}).then(function(ch){
    channel = ch;
    return channel.assertExchange('rethinkdb', 'topic', {durable: false});
})

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

promise = promise.then(function(){
    return channel.assertQueue('', {exclusive: true});
}).then(function(q){
    queue = q.queue;
})

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

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

promise = promise.then(function(){
    return channel.bindQueue(queue, 'rethinkdb', 'mytable.*');
})

Наконец, для прослушивания очереди мы используем генератор channel.consume. Подобно курсору changefeed из RethinkDB, consume будет вызывать свой обратный вызов всякий раз, когда в очереди появится сообщение.

promise = promise.then(function(){
    channel.consume(queue, function(msg){
        var change = JSON.parse(msg.content);
        var tablename = msg.fields.routingKey.split('.')[0];
        var changeType = msg.fields.routingKey.split('.')[1];

        console.log(tablename, 'got a change of type:', changeType);
        console.log(JSON.stringify(change, undefined, 2));
    })
})

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

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

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

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

Spec-Zone.ru

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