【问题标题】:How to convert RabbitMQ Messages into object list in c#如何在 c# 中将 RabbitMQ 消息转换为对象列表
【发布时间】:2022-02-21 02:24:32
【问题描述】:

我正在将 json 消息发布到 rabbitmq 中的队列中,并且它工作正常。但是面临一个问题,我想消耗已发布队列中的所有数据(作为聊天应用程序)并且我必须使用所有消息。

例如,我在队列中有 9 个项目,如下所示

{"Sender":123,"Message":"Test Message-1","Group":1}
{"Sender":123,"Message":"Test Message-2","Group":1}
{"Sender":123,"Message":"Test Message-3","Group":1}
{"Sender":123,"Message":"Test Message-4","Group":1}
{"Sender":567,"Message":"Test Message-5","Group":21}
{"Sender":123,"Message":"Test Message-6","Group":1}
{"Sender":456,"Message":"Test Message-7","Group":1}
{"Sender":456,"Message":"Test Message-8","Group":1}
{"Sender":123,"Message":"Test Message-9","Group":1} 

这些所有消息都按我的意愿存储在队列中。但是,当我尝试使用下面的 api 调用来收集它们时,它将无法正常工作。有时获取数据但有时未获取任何数据并确认列表。那么有什么方法可以将所有或有限的数据放入 c# 中的对象或数组中。因为所有示例都将消息消费到控制台中。我需要收藏。

public IList<string> GetMessageFromQueue(string _key, bool AutoAck = false)
        {
            var _list = new List<string>();
            var factory = new ConnectionFactory() { HostName = "localhost" };
            using (var connection = factory.CreateConnection())
            using (var channel = connection.CreateModel())
            {
                channel.QueueDeclare(queue: _key,
                                     durable: false,
                                     exclusive: false,
                                     autoDelete: false,
                                     arguments: null);

                var response = channel.QueueDeclarePassive(_key);
                var _test= response.MessageCount;
                var _test2 = response.ConsumerCount;

                var consumer = new EventingBasicConsumer(channel);
                consumer.Received += (model, ea) =>
                {
                    var body = ea.Body.ToArray();
                    var message = Encoding.UTF8.GetString(body);
                    _list.Add(message); 
                };
                //if (_list.Count == 0)
                //    AutoAck = false;
                channel.BasicConsume(queue: _key,
                                     autoAck: AutoAck,
                                     consumer: consumer);

            }
            return _list;
        }

还有我的控制器

public IActionResult Collect(){
    _queueClient.GetMessageFromQueue("myKey",true);
}

由于 BasicConsume 的 autoack 属性,此方法 olsa 会清除队列。我也尝试使用 basicAck。

在rabbitmq/c#中将消息发送到对象数组以进行下一步操作的最佳方法是什么。

【问题讨论】:

    标签: c# rabbitmq message-queue amqp


    【解决方案1】:

    在我看来,您的函数 GetMessageFromQueue 正在完成所有设置的动作,但随后立即退出该函数,而无需等待您的 Received 函数收集所有消息。

    例如,这是您设置用于从队列中收集消息的内联函数:

    consumer.Received += (model, ea) =>
    {
        var body = ea.Body.ToArray();
        var message = Encoding.UTF8.GetString(body);
        _list.Add(message); 
    };
    

    ...但是 2 行后,您只需立即退出该函数,而无需等待您的 Recieved 函数将所有消息添加到您的列表中。

    // exit function straight away!
    return _list;
    

    我注意到在您的示例代码中,您可以获得队列中保留的消息的计数。这很好,因为这意味着您知道预计会收到多少条消息。

    var _test= response.MessageCount;
    

    所以,您可以尝试做的一件事是添加 ManualResetEventSlimSemaphoreSlim 以在函数底部等待,直到它发出信号然后返回(可能有更好的方法来做到这一点,但就是这样的想法现在突然出现在我的脑海中)

    例如,在函数顶部创建一个ManualResetEventSlim 事件

    var msgsRecievedGate = new ManualResetEventSlim(false);
    

    然后在退出函数之前等待它被设置。

    msgsRecievedGate.Wait();
    

    类似这样的:

    public IList<string> GetMessageFromQueue(string _key, bool AutoAck = false)
    {
        var _list = new List<string>();
        
        // Setup synchronization event. 
        var msgsRecievedGate = new ManualResetEventSlim(false);
        
        var factory = new ConnectionFactory() { HostName = "localhost" };
        using (var connection = factory.CreateConnection())
        using (var channel = connection.CreateModel())
        {
            channel.QueueDeclare(queue: _key,
                                 durable: false,
                                 exclusive: false,
                                 autoDelete: false,
                                 arguments: null);
    
            var response = channel.QueueDeclarePassive(_key);
            
            var msgCount = response.MessageCount;
            var msgRecieved = 0;
            
            var consumer = new EventingBasicConsumer(channel);
            consumer.Received += (model, ea) =>
            {
                msgRecieved++;
                
                var body = ea.Body.ToArray();
                var message = Encoding.UTF8.GetString(body);
                _list.Add(message); 
                
                if ( msgRecieved == msgCount )
                {
                    // Set signal here
                    msgsRecievedGate.Set();
                    
                    // exit function 
                    return;
                }
            };
            
            
            channel.BasicConsume(queue: _key,
                                 autoAck: AutoAck,
                                 consumer: consumer);
    
        }
        
        // Wait here until all messages are retrieved
        msgsRecievedGate.Wait();
        
        // now exit function! 
        return _list;
    }
    

    请注意。我没有测试上面的代码,所以你的里程我会有所不同。

    【讨论】:

      猜你喜欢
      • 2019-05-01
      • 2015-11-26
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多