【问题标题】:RabbitMQ: multiple messages and single consumerRabbitMQ:多条消息和单个消费者
【发布时间】:2015-07-31 09:59:57
【问题描述】:

我对 RabbitMQ 消费者有疑问。实际上我有一个消费者从三个队列中获取消息。问题是我需要从他们每个人那里获得多条消息,但是我的消费者每个队列只获得一条消息并最终获得。如果有人能帮我解决这个问题,我将不胜感激。

下面的消费者代码

        for (int i = 0; i < queueNames.size(); i++) {

        Channel channel = connection.createChannel();
        QueueingConsumer consumer = new QueueingConsumer(channel);
        channel.basicConsume(queueNames.get(i).toString(), true, consumer_tag, consumer);

        flag = true;
        while (flag) {

            QueueingConsumer.Delivery delivery = consumer.nextDelivery();
            String routingKey = delivery.getEnvelope().getRoutingKey();
            System.out.println(routingKey);
            String message = new String(delivery.getBody(), "UTF-8");

                flag = false;
        }
    }

其中 queueNames 是一个包含我的队列名称的列表(以 3 个为单位)。

【问题讨论】:

    标签: java rabbitmq messaging amqp


    【解决方案1】:

    您需要订阅队列,消费者只会按照您定义的方式消费 1 条消息

    boolean autoAck = false;
    channel.basicConsume(queueName, autoAck, "myConsumerTag",
     new DefaultConsumer(channel) {
         @Override
         public void handleDelivery(String consumerTag,
                                    Envelope envelope,
                                    AMQP.BasicProperties properties,
                                    byte[] body)
             throws IOException
         {
             String routingKey = envelope.getRoutingKey();
             String contentType = properties.getContentType();
             long deliveryTag = envelope.getDeliveryTag();
             // (process the message components here ...)
             channel.basicAck(deliveryTag, false);
         }
     });
    

    更多信息在这里:https://www.rabbitmq.com/api-guide.html

    【讨论】:

    • 当我只做一个 while(true) 循环时,来自单个队列的所有消息都会被接收到,但这里又出现了另一个问题。消费者不知道是否还有更多消息,并且仍在等待它们。如果没有更多消息可以从这个循环中获取,我想去另一个队列。有什么建议吗?
    【解决方案2】:

    好的,我这样解决问题:

    boolean flag;
        System.out.println("Rozmiar queue " + queueNames.size());
        for (int i = 0; i < queueNames.size(); i++) {
    
            Channel channel = connection.createChannel();
            QueueingConsumer consumer = new QueueingConsumer(channel);
            channel.basicConsume(queueNames.get(i).toString(), true, consumer_tag, consumer);
    
            flag = true;
            while (flag) {
    
                QueueingConsumer.Delivery delivery = consumer.nextDelivery(timeout);
                if (delivery == null) {
                    flag = false;
                } else {
    
                    String message = new String(delivery.getBody(), "UTF-8");
                    System.out.println(" [x] Message Received '" + message + "'");
                }
            }
        }
    

    我希望这个解决方案将来可以帮助某人:)

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2012-05-24
      • 1970-01-01
      • 2020-04-20
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多