跳至主内容

消费者优先级与 RabbitMQ

·阅读时长4分钟
Álvaro Videla

在 RabbitMQ 3.2.0 中,我们引入了 消费者优先级,顾名思义,它允许我们为消费者设置优先级。这让我们能够在一定程度上控制 RabbitMQ 如何将消息传递给消费者,以实现对应用程序有益的另一种调度方式。

您会在什么时候在代码中使用消费者优先级?

异构集群

假设我们的工作节点集群并非运行在完全相同的硬件上。根据我们执行任务的类型,某些机器具备的硬件特性使其在集群中占据优势。例如,有些机器配备了 SSD,而我们的任务需要大量的 I/O 操作;或者任务需要更快的 CPU 来执行计算;又或者需要更多的内存来缓存结果以备后续计算使用。无论如何,如果现在有两个消费者准备好接收更多消息,且其中一个处于性能更好的机器上,那么 RabbitMQ 理应选择性能更好的机器上的消费者来派发消息,而不是选择那台较差的机器。请记住,消费者优先级*仅*对准备好接收消息的消费者生效。因此,如果性能较差的机器上的消费者已就绪,而较好的机器上没有就绪的消费者,RabbitMQ 将直接把消息发送给该消费者,而不会等待更快的消费者变为可用状态。

数据本地性

消费者优先级的另一个用途是利用数据本地性。在 RabbitMQ 中,队列内容驻留在最初声明该队列的节点上;在镜像队列的情况下,会有一个主节点(master node)来协调该队列。因此,虽然消费者可以连接到集群中的不同节点并从镜像中获取消息,但归根结底,关于谁消费了哪些消息的信息最终还是会传回主节点。在这种情况下,我们可以使用消费者优先级来告知 RabbitMQ 优先将消息派发给连接到主节点的消费者。为此,连接到主节点的消费者在发出 basic.consume 命令时,需要为自己设置更高的优先级(前提是它有办法知道自己是否连接到了主节点)。

声明消费者优先级

以下示例代码展示了如何使用 RabbitMQ Java 客户端声明消费者优先级。

import java.util.*;
import com.rabbitmq.client.ConnectionFactory;
import com.rabbitmq.client.Connection;
import com.rabbitmq.client.Channel;
import com.rabbitmq.client.QueueingConsumer;

public class Consumer {

private final static String EXCHANGE_NAME = "my_exchange";
private final static String QUEUE_NAME = "my_queue";

public static void main(String[] argv) throws Exception {
ConnectionFactory factory = new ConnectionFactory();
factory.setHost("localhost");
Connection connection = factory.newConnection();
Channel channel = connection.createChannel();

channel.queueDeclare(QUEUE_NAME, true, false, false, null);
channel.exchangeDeclare(EXCHANGE_NAME, "direct", true);
channel.queueBind(QUEUE_NAME, EXCHANGE_NAME, "");
System.out.println("Waiting for messages. To exit press CTRL+C");

QueueingConsumer consumer = new QueueingConsumer(channel);

Map<String, Object> args = new HashMap<String, Object>();
args.put("x-priority", 10);
channel.basicConsume(QUEUE_NAME, false, "", false, false, args, consumer);

while (true) {
QueueingConsumer.Delivery delivery = consumer.nextDelivery();
String message = new String(delivery.getBody());
System.out.println("Received '" + message + "'");
channel.basicAck(delivery.getEnvelope().getDeliveryTag(), false);
}
}
}

这段代码基于教程 1 中的示例实现了一个非常简单的消费者。有趣的部分在第 25 到 27 行,我们首先创建了一个 HashMap 来保存传给 basicConsume 的参数。我们创建了一个名为 x-priority、值为 10 的参数(值越高,优先级越高)。当我们调用 basicConsume 时,将这些参数传给 RabbitMQ,这就完成了!这是一个非常强大且易于使用的功能。像往常一样,明智的做法是进行性能测试,以确定最适合我们消费者的优先级策略。

© . This site is unofficial and not affiliated with VMware.