RabbitMQ Messaging & Async Systems · 课时

工作队列:公平分发

使用轮询分发策略,在多个消费者之间分配任务,实现工作队列,并学习如何异步处理耗时任务。

第 2 / 4 课11 个步骤

工作队列:公平分发 是 CoddyKit 上的免费 RabbitMQ Messaging & Async Systems 课时。 这是第 2 节课,共 4 节。 你可以在下方免费阅读本课时的完整内容 — 然后在浏览器中使用内置代码编辑器和全天候 AI 导师进行实践。 这是 RabbitMQ Messaging & Async Systems 学习路径的一部分,你的进度在网页和 CoddyKit 应用中同步。 RabbitMQ Messaging & Async Systems 课程共包含 4 节课。

本课时的部分内容尚未翻译,以英文显示。

Meet Work Queues

Welcome to Work Queues! In distributed systems, you often have tasks that take time to complete, like processing an image or generating a report.

Work Queues are a pattern that helps distribute these time-consuming tasks among multiple workers (consumers) efficiently, preventing any single worker from getting overloaded.

Why Use Work Queues?

Imagine you have many jobs to do, but only one employee. If all jobs go to that single employee, they'll get overwhelmed and tasks will pile up.

  • Load Balancing: Work queues allow you to add more workers to share the load.
  • Asynchronous Processing: The producer doesn't wait for a task to finish, it just adds it to the queue.
  • Reliability: If one worker fails, others can pick up tasks.

How Work Queues Operate

The setup for a work queue is simple:

  • One producer sends messages (tasks) to a single queue.
  • Multiple consumers listen to this same queue.
  • RabbitMQ ensures that each message is delivered to only one of the waiting consumers.

This way, tasks are never duplicated and are processed in parallel.

Fair Dispatch: Round-Robin

By default, RabbitMQ distributes messages to consumers using a round-robin strategy. This means messages are sent to consumers in a rotating fashion:

  • Consumer 1 gets the first message.
  • Consumer 2 gets the second message.
  • Consumer 1 gets the third message, and so on.

This aims for an even distribution of tasks among all active consumers.

The Task Producer

Let's create a producer that sends 10 tasks to our queue. Each task will be a simple string like 'Processing image 1'.

Run this code once to populate the queue with tasks.

import com.rabbitmq.client.Channel;
import com.rabbitmq.client.Connection;
import com.rabbitmq.client.ConnectionFactory;

public class TaskProducer {
    private final static String QUEUE_NAME = "task_queue";

    public static void main(String[] argv) throws Exception {
        ConnectionFactory factory = new ConnectionFactory();
        factory.setHost("localhost"); // Assuming RabbitMQ is local

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

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

            for (int i = 0; i < 10; i++) {
                String message = "Processing image " + (i + 1);
                channel.basicPublish("", QUEUE_NAME, null, message.getBytes("UTF-8"));
                System.out.println(" [x] Sent '" + message + "'");
            }
            System.out.println(" [x] All tasks sent.");
        }
    }
}

Producer Code Breakdown

What's happening in our producer code?

  • We establish a connection and a channel to interact with RabbitMQ.
  • channel.queueDeclare(QUEUE_NAME, false, false, false, null); ensures the queue exists. The false flags keep it non-durable and non-exclusive for simplicity here.
  • A loop sends 10 messages to the task_queue. Each message represents a distinct task.

Our First Task Consumer

Now, let's create a consumer, which we'll call a 'worker'. This worker will listen for tasks from the task_queue.

We'll simulate a long-running task using Thread.sleep(). Run this code in one terminal window.

import com.rabbitmq.client.Channel;
import com.rabbitmq.client.Connection;
import com.rabbitmq.client.ConnectionFactory;
import com.rabbitmq.client.DeliverCallback;

public class TaskConsumer {
    private final static String QUEUE_NAME = "task_queue";

    public static void main(String[] argv) throws Exception {
        ConnectionFactory factory = new ConnectionFactory();
        factory.setHost("localhost"); // Assuming RabbitMQ is local

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

        channel.queueDeclare(QUEUE_NAME, false, false, false, null);
        System.out.println(" [*] Waiting for messages. To exit press CTRL+C");

        DeliverCallback deliverCallback = (consumerTag, delivery) -> {
            String message = new String(delivery.getBody(), "UTF-8");
            System.out.println(" [x] Received '" + message + "'");
            try {
                // Simulate long-running task
                Thread.sleep(1000); // 1 second per task
            } catch (InterruptedException e) {
                Thread.currentThread().interrupt();
            } finally {
                System.out.println(" [x] Done with '" + message + "'");
            }
        };
        channel.basicConsume(QUEUE_NAME, true, deliverCallback, consumerTag -> {});
    }
}

Consumer Code Breakdown

Let's look at the consumer's logic:

  • Similar to the producer, it connects and declares the queue.
  • A DeliverCallback defines the actions when a message arrives.
  • Inside the callback, Thread.sleep(1000) simulates a 1-second task.
  • channel.basicConsume(QUEUE_NAME, true, deliverCallback, ...) starts consuming. The true means messages are automatically acknowledged after delivery.

Scale with Multiple Workers

Here's the core demonstration of work queues:

1. Run the TaskConsumer code in two separate terminal windows (or instances).

2. Then, run the TaskProducer code once.

You will observe that the 10 tasks are divided between your two worker instances, each processing roughly 5 tasks due to RabbitMQ's round-robin dispatch.

Work Queue Quiz

Imagine you have a single RabbitMQ queue and two consumers (Worker A and Worker B) listening to it. A producer sends 4 messages (M1, M2, M3, M4) to this queue.

Which statement accurately describes how the messages are typically distributed using RabbitMQ's default fair dispatch?

Work Queues: Key Takeaways

You've successfully implemented Work Queues, a fundamental pattern for distributing tasks across multiple consumers!

  • Work queues enable asynchronous processing and prevent single points of failure.
  • RabbitMQ's default round-robin dispatch strategy ensures tasks are distributed fairly.
  • By running multiple consumer instances, you can easily scale your task processing capacity.

Next, we'll dive into making your message handling even more robust with acknowledgements and message durability!

免费开始

用 AI 导师学习 RabbitMQ Messaging & Async Systems — 免费

在浏览器中编写并运行真实代码,获得全天候 AI 导师的即时帮助,并在网页或应用中继续学习。

课程
11
课程
44

常见问题解答

「工作队列:公平分发」课时是免费的吗?

是的 — 「工作队列:公平分发」的完整文本可在网页上免费阅读。要进行交互式练习(内置代码编辑器和全天候 AI 导师)并解锁 RabbitMQ Messaging & Async Systems 课程的其余内容,请升级到 CoddyKit PRO。 RabbitMQ Messaging & Async Systems 课程共包含 4 节课。

「工作队列:公平分发」这节课中我会学到什么?

使用轮询分发策略,在多个消费者之间分配任务,实现工作队列,并学习如何异步处理耗时任务。 你通过在浏览器中直接运行的动手代码来练习 RabbitMQ Messaging & Async Systems,全天候 AI 导师会在你学习这节课的过程中回答你的问题。

学习 RabbitMQ Messaging & Async Systems 需要有经验吗?

无需任何先前经验。CoddyKit 上的 RabbitMQ Messaging & Async Systems 课程适合初学者到高级学习者,你可以从这里开始或从头开始,按照自己的节奏学习。 这是第 2 节课,共 4 节。

「工作队列:公平分发」课时需要多长时间?

大多数 CoddyKit 课程大约需要 5–10 分钟。每节课都很精短且互动,所以你能稳步进步,并在网页和应用中从离开的地方继续。

我能在这节 RabbitMQ Messaging & Async Systems 课中编写并运行代码吗?

能。每节 RabbitMQ Messaging & Async Systems 课都包含内置代码编辑器,你可以在浏览器中直接编写并运行真实代码,并获得即时 AI 反馈 — 无需本地设置。

此课程中的所有课时

  1. Hello World:简单队列
  2. 工作队列:公平分发
  3. 消息确认与持久化
  4. 使用扇出交换器实现发布/订阅
← 返回 RabbitMQ Messaging & Async Systems