跳至主内容

使用 RabbitMQ 实现分布式信号量

·8 分钟阅读
Álvaro Videla

在这篇博文中,我们将解决在分布式系统中控制对特定资源的访问的问题。解决这个问题的方法在计算机科学中是众所周知的,它被称为信号量(Semaphore),由 Dijkstra 在 1965 年的论文“Cooperating Sequential Processes”中发明。我们将看到如何使用 AMQP 的构建块(如消费者、生产者和队列)来实现它。

信号量的需求

在探讨具体解决方案之前,让我们先看看在什么场景下会真正需要这样的机制。

假设我们的应用程序有多个进程从队列中获取作业,然后将记录插入数据库,我们可能需要限制同时执行这些操作的进程数量。

同样,如果工作进程正在调整图像大小,并在准备好后通过网络将其存储到远程服务器上。我们希望防止因同时传输过多的图像而导致网络链路拥塞,因此我们也限制了同时能传输图像的工作进程数量。通过这种方式,虽然工作进程会尽可能快地调整图像大小,但当轮到它们使用网络链路时,它们会批量地将图像移动到最终目的地。

另一个与 RabbitMQ 相关的例子是:您的应用程序可能只需要一组生产者中的一个向交换机发送消息,但一旦该进程停止,您希望该组中的下一个生产者立即开始发送消息。这种需求有很多原因,容量控制就是其中之一。

另一方面,可能存在消费者竞争访问队列的需求。虽然 AMQP 提供了独占队列和独占消费者的方式,但空闲的消费者无法感知队列访问权何时被释放。因此,使用与上述类似的方法,我们可以让消费者轮流访问队列。

值得注意的是,并没有什么限制阻止我们让多个进程访问特定资源。假设我们有十个生产者,但我们只希望其中五个同时发布消息。使用信号量,我们也可以实现这一点。

上述例子都有一个共同的额外要求:竞争资源的进程不应该为了知道何时可以开始工作而不断轮询 RabbitMQ 或其他协调器。理想情况下,它们应该处于空闲等待状态,一旦资源被释放,RabbitMQ 会通知下一个进程,以便它可以自动开始工作。

现在让我们进入实现阶段。

信号量的实现

我们的信号量将使用队列和消息来实现。没错,就是这么简单!

首先,我们声明一个名为 resource.semaphore 的队列,其中 resource 是信号量要控制的资源名称,它可以是 "images"、"database"、"file_server" 或任何适合您特定应用的名称。

我们向 resource.semaphore 队列发布一条消息。然后,我们启动那些寻求获取访问权限的进程。每个进程都将从 resource.semaphore 队列中进行消费;第一个到达的进程将获取该消息,而所有其他进程则会处于空闲等待状态。诀窍在于,这些进程永远不会确认 (acknowledge) 消息,但它们会以 ack_mode=on 的方式从 resource.semaphore 队列进行消费。因此,RabbitMQ 会跟踪该消息,如果进程崩溃或退出,消息将回到队列中,并被发送给信号量队列中监听的下一个进程。

通过这种简单的技术,同一时刻只能有一个进程访问资源,并且我们可以确保如果进程崩溃,它不会占用资源。当然,我们假设所有访问信号量的进程表现良好,即:它们永远不会确认消息。如果它们确认了,RabbitMQ 将删除该消息,组内所有其他进程都将陷入饥饿状态。

当某个进程希望停止时,它该如何交还“令牌”?当然,该进程可以突然关闭通道,RabbitMQ 会自动处理该消息,但还有一种更友好的方式。进程可以 basic.reject 该消息,告诉 RabbitMQ 重新排队 (re-enqueue),使其返回到信号量队列。

让我们看看代码实现,假设我们已经获得了一个连接和一个通道。

以下是设置信号量的代码:

channel.queueDeclare("resource.semaphore", true, false, false, null);
String message = "resource";
channel.basicPublish("", "resource.semaphore", null, message.getBytes());

我们创建一个名为 "resource.semaphore" 的持久化队列,然后使用默认交换机向其发布一条消息。

以下是进程访问信号量时会使用的代码:

QueueingConsumer consumer = new QueueingConsumer(channel);
channel.basicQos(1);
channel.basicConsume("resource.semaphore", false, consumer);

while (true) {
QueueingConsumer.Delivery delivery = consumer.nextDelivery();

// here we access the resource controlled by the semaphore.

if(shouldStopProcessing()) {
channel.basicReject(delivery.getEnvelope().getDeliveryTag(), true);
}
}

我们在那里创建了一个 QueueingConsumer,等待来自 "resource.semaphore" 队列的消息。通过在 basicQos 调用中设置 prefetch-count 为 1,我们确保我们的进程每次只从队列中获取一条消息。一旦收到消息,该进程就会开始使用该资源。当 shouldStopProcessing() 条件满足时,进程将 basicReject 该消息,告诉 RabbitMQ 将其重新入队。请记住,消费者是在确认模式下启动的,它永远不会确认从信号量队列收到的消息。如果它确认了,则被视为存在 Bug。

信号量访问优先级

是否可以为信号量访问设置优先级?是的,自 3.2.0 版本起,RabbitMQ 支持 消费者优先级 (Consumer Priorities)。通过使用消费者优先级,我们可以告诉 RabbitMQ 在传递信号量令牌消息时优先考虑哪些进程。

二进制信号量 vs 计数信号量

到目前为止,我们实现的是所谓的二进制信号量,即同一时间只允许一个进程访问资源的信号量。如果我们允许同一时间有多个进程访问同一资源,但仍需要限制该操作的并发数,那么我们可以实现计数信号量。为此,在设置信号量时,不是发布一条消息,而是发布与允许同时工作的进程数量相等的消息。我们需要确保进程像之前一样将 prefetch-count 值设置为 1。

修改计数

请注意,设置信号量队列的进程可以随着时间的推移添加额外消息以增加处理容量。如果我们想减少可以同时访问资源的进程数量,则必须停止正在运行的进程并清除队列。另一种方法是启动一个具有非常 高优先级 的额外消费者,这样它就会从信号量队列中获取所需数量的消息并确认它们,从而将其从系统中移除。

扩展阅读

如你所见,使用 AMQP 基本结构实现信号量非常简单,并且通过 RabbitMQ,我们还可以为资源访问设定优先级。

最后,我想分享一些关于信号量作为并发构件的文章。首先是 Dijkstra 的开创性论文 Cooperating Sequential Processes。最后是维基百科关于信号量的文章,其中解释了许多定义:信号量

编辑:2014.02.20

正如这里及其他地方与同事讨论的那样,这种设置对网络分区不具有弹性,因此请谨慎使用。感谢 @aphyr 和其他人为我提供的反馈。在 RabbitMQ 团队,我们始终保持诚实,并告知用户服务器能做什么和不能做什么。

编辑:2014.03.10

值得注意的是,这种设置在一般情况下对网络故障也不具备弹性。例如,可能出现工作进程持有令牌时与服务器的连接突然中断。此时服务器会收回令牌并将其重新入队,以便分发给另一个工作进程。在此期间,网络连接中断的那个工作进程仍会认为自己拥有令牌,因此会继续访问它本不该访问的资源。

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