【问题标题】:How to consume one message?如何消费一条消息?
【发布时间】:2015-04-24 07:49:23
【问题描述】:

在rabbitmq中使用example,消费者一次从队列中获取所有消息。如何消费一条消息并退出?

QueueingConsumer consumer = new QueueingConsumer(channel);
channel.basicConsume(QUEUE_NAME, true, consumer);

while (true) {
  QueueingConsumer.Delivery delivery = consumer.nextDelivery();
  String message = new String(delivery.getBody());
  System.out.println(" [x] Received '" + message + "'");
}

【问题讨论】:

  • 如果不使用循环,所有消息都会丢失,除了一个。
  • 这不是你想要的吗?使用一条消息并退出。
  • QueueingConsumer.Delivery 交付 = consumer.nextDelivery();一次读取队列中的所有消息

标签: java rabbitmq


【解决方案1】:

您必须声明 basicQos 设置才能一次获取一条消息,从 ACK 状态变为 NACK 状态,并禁用自动 ACK 以明确给予确认。

ConnectionFactory factory = new ConnectionFactory();
    factory.setHost("localhost");
    Connection connection = factory.newConnection();
    Channel channel = connection.createChannel();
    channel.basicQos(1);
    channel.queueDeclare(QUEUE_NAME, true, false, false, null);
    System.out.println("[*] waiting for messages. To exit press CTRL+C");

    QueueingConsumer consumer = new QueueingConsumer(channel);
    channel.basicConsume(QUEUE_NAME, consumer);
    while(true) {
        QueueingConsumer.Delivery delivery = consumer.nextDelivery();
        int n = channel.queueDeclarePassive(QUEUE_NAME).getMessageCount();
        System.out.println(n);
        if(delivery != null) {
            byte[] bs = delivery.getBody();
            System.out.println(new String(bs));
            //String message= new String(delivery.getBody());
            channel.basicAck(delivery.getEnvelope().getDeliveryTag(), false);
            //System.out.println("[x] Received '"+message);
        }
    }

希望对你有帮助!

【讨论】:

  • 另外,作为对 OP 的警告:RabbitMQ 团队不建议一次只阅读一个。原因:rabbitmq.com/blog/2011/09/24/sizing-your-rabbits
  • 我不认为将 BasicQOS 设置为 1 可以解决它。 “这是通过使用 basic.qos 方法设置“预取计数”值来完成的。该值定义了通道上允许的未确认传递的最大数量。一旦数量达到配置的计数,RabbitMQ 将停止传递更多消息除非至少有一个未完成的通道被确认。"。来自rabbitmq.com/confirms.html
【解决方案2】:

使用 AMQP 0.9.1 basic.get 同步获取一条消息。

ConnectionFactory factory = new ConnectionFactory();
factory.setUri(uri);

Connection connection = factory.newConnection();
Channel channel = connection.createChannel();

channel.queueDeclare(QUEUE_NAME, true, false, false, null);

GetResponse response = channel.basicGet(QUEUE_NAME, true);
if (response != null) {
    String message = new String(response.getBody(), "UTF-8");
}

channel.close();
connection.close();

【讨论】:

  • 完美,正是我正在寻找的快速集成测试。
【解决方案3】:
const consumeFromQueue = async (queueName) => {
    try {

        let data = await channel.get(queueName)// get one msg at a time
        if (data) {

            data.content ? eval("(" + data.content.toString() + ")()") : ""
            channel.ack(data)
        } else {
            //console.log("Empty Queue")
        }
    }
    catch (error) {
        //console.log("Error while consuming from rabbitmq queue", error)
        return Promise.reject(error)
    }
}

【讨论】:

  • 那不是 Java,问题被标记为 Java。
猜你喜欢
  • 2012-10-08
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2018-12-01
  • 2021-07-28
  • 2019-04-22
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多