代码之家  ›  专栏  ›  技术社区  ›  Jack Skeletron

Amqp、rabbit mq和socket.io重新连接到队列,即使客户端已关闭

  •  0
  • Jack Skeletron  · 技术社区  · 8 年前

    当我用用户登录到我的系统时,它会创建一个通知UID- 用户ID queue(目前queueName是由query oaraeter发送的,我将在解决问题后立即实现更详细的方法)

    第二用户ID .

    如果我注销其中一个用户,队列将消失(因为它不持久)。

    问题是,当我在另一个会话中刷新或加载另一个页面时,即使参数queuename未发送,它也会重新创建第二个队列。

    服务器.js

    var amqp = require('amqp');
    var app = require('express')();
    var http = require('http').Server(app);
    var io = require('socket.io')(http);
    
    var rabbitMqConnection = null;
    var _queue = null;
    var _consumerTag = null;
    
    
    io.use(function (socket, next) {
        var handshakeData = socket.handshake;
        // Here i will implement token verification
        console.log(socket.handshake.query.queueName);
        next();
    });
    
    
    // Gets the connection event form client
    io.sockets.on('connection', function (socket) {
    
        var queueName = socket.handshake.query.queueName;
    
        console.log("Socket Connected");
    
        // Connects to rabbiMq
        rabbitMqConnection = amqp.createConnection({host: 'localhost', reconnect: false});
    
        // Update our stored tag when it changes
        rabbitMqConnection.on('tag.change', function (event) {
            if (_consumerTag === event.oldConsumerTag) {
                _consumerTag = event.consumerTag;
                // Consider unsubscribing from the old tag just in case it lingers
                _queue.unsubscribe(event.oldConsumerTag);
            }
        });
    
        // Listen for ready event
        rabbitMqConnection.on('ready', function () {
            console.log('Connected to rabbitMQ');
    
            // Listen to the queue
            rabbitMqConnection.queue(queueName, {
                    closeChannelOnUnsubscribe: true,
                    durable: false,
                    autoClose: true
                },
                function (queue) {
                    console.log('Connected to ' + queueName);
                    _queue = queue;
    
                    // Bind to the exchange
                    queue.bind('users.direct', queueName);
    
                    queue.subscribe({ack: false, prefetchCount: 1}, function (message, headers, deliveryInfo, ack) {
                        console.log("Received a message from route " + deliveryInfo.routingKey);
                        socket.emit('notification', message);
                        //ack.acknowledge();
                    }).addCallback(function (res) {
                        // Hold on to the consumer tag so we can unsubscribe later
                        _consumerTag = res.consumerTag;
                    });
                });
        });
    
    
        // Listen for disconnection
        socket.on('disconnect', function () {
            _queue.unsubscribe(_consumerTag);
            rabbitMqConnection.disconnect();
            console.log("Socket Disconnected");
        });
    
    });
    
    http.listen(8080);
    

    客户端.js

    var io = require('socket.io-client');
    
    $(document).ready(function () {
    
        var socket = io('http://myserver.it:8080/', {
             query:  { queueName: 'notification-UID-' + UID},
            'sync disconnect on unload': true,
            });
    
        socket.on('notification', function (data) {
            console.log(data);
        });
    })
    

    知道吗?

    1 回复  |  直到 8 年前
        1
  •  1
  •   Jack Skeletron    8 年前

    所以我解决了我的问题,这是一个可变范围的问题,把事情搞砸了。 让我解释一下我在做什么,也许对某人有用。

    基本上,我正在尝试创建一个浏览器通知系统,这意味着我的应用程序将包含一些信息(如主题、链接和消息)的通知对象发布到exchange(生产者端)。

    交易所是 扇形分叉 ( 用户.notification.fanout 用户.direct 用户.notification.store (扇出型)。

    用户ID

    通知对象同时指向users.direct和users.notification.store 最后一个用户将通知写入数据库以防用户未登录,第一个用户将通知发布到浏览器。

    那么浏览器消费者是如何工作的呢?

    我使用了socket.io、node server和amqplib的经典组合。

    用户ID 并将其绑定到用户。直接交换。

    同时,我在我的服务器上添加了https,所以与第一个版本相比有些变化。

    所以我的 服务器.js

    var amqp = require('amqp');
    var fs = require('fs');
    var app = require('express')();
    // Https server, certificates and private key added
    var https = require('https').Server({
        key: fs.readFileSync('/home/www/site/privkey.pem'),
        cert: fs.readFileSync('/home/www/site/fullchain.pem')},app);
    var io = require('socket.io')(https);
    
    // Used to verify if token is valid
    // If not it will discard connection
    io.use(function (socket, next) {
        var handshakeData = socket.handshake;
        // Here i will implement token verification
        console.log("Check this token: " + handshakeData.query.token);
        next();
    });
    // Gets the connection event from client
    io.sockets.on('connection', function (socket) {
        // Connection log
        console.log("Socket Connected with ID: " + socket.id);
        // THIS WAS THE PROBLEM
        // Local variables for connections
        // Former i've put these variables outside the connection so at 
        // every client they were "overridden". 
        // RabbitMq Connection (Just for current client)
        var _rabbitMqConnection = null;
        // Queue (just for current client)
        var _queue = null;
        // Consumer tag (just for current client)
        var _consumerTag = null;
        // Queue name and routing key for current user
        var queueName = socket.handshake.query.queueName;
        // Connects to rabbiMq with default data to localhost guest guest
        _rabbitMqConnection = amqp.createConnection();
        // Connection ready
        _rabbitMqConnection.on('ready', function () {
            // Connection log
            console.log('#' + socket.id + ' - Connected to RabbitMQ');
            // Creates the queue (default is transient and autodelete)
            // https://www.npmjs.com/package/amqp#connectionqueuename-options-opencallback
            _rabbitMqConnection.queue(queueName, function (queue) {
                // Connection log
                console.log('#' + socket.id + ' - Connected to ' + queue.name + ' queue');
                // Stores local queue
                _queue = queue;
                // Bind to the exchange (default)
                queue.bind('users.direct', queueName, function () {
                    // Binding log
                    console.log('#' + socket.id + ' - Binded to users.direct exchange');
                    // Consumer definition
                    queue.subscribe({ack: false}, function (message, headers, deliveryInfo, messageObject) {
                        // Message log
                        console.log('#' + socket.id + ' - Received a message from route ' + deliveryInfo.routingKey);
                        // Emit the message to the client
                        socket.emit('notification', message);
                    }).addCallback(function (res) {
                        // Hold on to the consumer tag so we can unsubscribe later
                        _consumerTag = res.consumerTag;
                        // Consumer tag log
                        console.log('#' + socket.id + ' - Consumer ' + _consumerTag + ' created');
                    })
                });
    
            });
        });
        // Update our stored tag when it changes
        _rabbitMqConnection.on('tag.change', function (event) {
            if (_consumerTag === event.oldConsumerTag) {
                _consumerTag = event.consumerTag;
                // Unsubscribe from the old tag just in case it lingers
                _queue.unsubscribe(event.oldConsumerTag);
            }
        });
        // Listen for disconnection
        socket.on('disconnect', function () {
            _queue.unsubscribe(_consumerTag);
            _rabbitMqConnection.disconnect();
            console.log('#' + socket.id + ' - Socket Disconnected');
        });
    });
    
    https.listen(8080);
    

    那么我的 客户端.js

    var io=需要('socket.io client');

    $(document).ready(function () {
    
        var socket = io('https://myserver.com:8080/', {
            secure: true, // for ssl connections
            query:  { queueName: 'notification-UID-' + UID, token: JWTToken}, // params sent to server, JWTToken for authentication
            'sync disconnect on unload': true // Every time the client unload, socket disconnects
        });
    
        socket.on('notification', function (data) {
            // Do what you want with your data
            console.log(data);
        });
    })