RabbitMQ Messaging & Async Systems · 课时

用于发布/订阅的扇出交换器

了解并实现扇出交换器,将消息广播到所有已绑定的队列,适用于简单的发布/订阅场景。

第 1 / 4 课11 个步骤

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

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

Pub/Sub & Fanout Explained

Imagine you want to broadcast a message to everyone interested, without knowing who they are. This is the idea behind the Publish/Subscribe (Pub/Sub) messaging pattern.

In RabbitMQ, the Fanout exchange is perfect for this. It acts like a megaphone, shouting your message to all connected listeners.

How Fanout Exchanges Work

A Fanout exchange is the simplest type of exchange. When a message arrives at a Fanout exchange, it doesn't care about routing keys.

  • It takes the message.
  • It duplicates it for every queue that is bound to it.
  • Then, it sends a copy of the message to each of those bound queues.

Think of it as a broadcast to all subscribers.

Key Components

Let's quickly recap the main players:

  • Producer: Sends the message.
  • Exchange: Receives messages from producers and routes them to queues. Fanout is one type.
  • Queue: A buffer that stores messages until a consumer picks them up.
  • Consumer: Receives messages from queues and processes them.

With Fanout, the exchange ensures all bound queues get the message.

Routing Keys Are Ignored

A crucial detail about Fanout exchanges: they completely ignore routing keys!

When a producer sends a message to a Fanout exchange, it might still provide a routing key (often an empty string), but the exchange simply disregards it.

Its only job is to broadcast to all queues bound to it, regardless of any key.

Setting Up the Fanout Exchange

First, our producer needs to declare the Fanout exchange. This tells RabbitMQ to create or ensure this exchange exists.

Notice the "fanout" type parameter. Try running this code to declare your exchange!

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

public class FanoutProducerSetup {
    private final static String EXCHANGE_NAME = "fanout_logs";

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

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

            // Declare a fanout exchange
            channel.exchangeDeclare(EXCHANGE_NAME, "fanout");
            System.out.println("Fanout exchange '" + EXCHANGE_NAME + "' declared.");
        }
    }
}

Binding a Queue to Fanout

Consumers don't receive directly from exchanges. They receive from queues. Each consumer needs its own queue, and that queue must be bound to the Fanout exchange.

queueDeclare() with no arguments creates a unique, exclusive, auto-delete queue. The routing key for binding is an empty string, as it's ignored by Fanout.

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

public class FanoutConsumerSetup {
    private final static String EXCHANGE_NAME = "fanout_logs";

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

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

        channel.exchangeDeclare(EXCHANGE_NAME, "fanout");
        String queueName = channel.queueDeclare().getQueue(); // A unique, auto-delete queue
        channel.queueBind(queueName, EXCHANGE_NAME, ""); // Bind with empty routing key

        System.out.println("Queue '" + queueName + "' declared and bound to '" + EXCHANGE_NAME + "'.");
        System.out.println("Ready for messages (but not consuming yet).");
    }
}

Publishing to Fanout Exchange

Once the exchange is declared, the producer can send messages to it. Notice how the basicPublish method specifies the exchange name, but the routing key is an empty string.

Run this code after you've set up your exchange. It will publish one message.

import com.rabbitmq.client.Channel;
import com.rabbitmq.client.Connection;
import com.rabbitmq.client.ConnectionFactory;
import java.nio.charset.StandardCharsets;

public class FanoutPublisher {
    private final static String EXCHANGE_NAME = "fanout_logs";

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

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

            channel.exchangeDeclare(EXCHANGE_NAME, "fanout"); // Ensure exchange exists

            String message = "Hello everyone, this is a broadcast!";
            // Publish to the exchange, routing key is ignored for fanout
            channel.basicPublish(EXCHANGE_NAME, "", null, message.getBytes(StandardCharsets.UTF_8));
            System.out.println(" [x] Sent '" + message + "'");
        }
    }
}

Consuming Broadcasts

Now, let's make our consumer actually receive messages. The DeliverCallback defines what happens when a message arrives. Remember, each consumer has its own queue!

To see the broadcast in action, first run two separate instances of this consumer code. Then, run the publisher code from the previous scene. Both consumers should receive the message!

import com.rabbitmq.client.Channel;
import com.rabbitmq.client.Connection;
import com.rabbitmq.client.ConnectionFactory;
import com.rabbitmq.client.DeliverCallback;
import java.nio.charset.StandardCharsets;

public class FanoutSubscriber {
    private final static String EXCHANGE_NAME = "fanout_logs";

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

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

        channel.exchangeDeclare(EXCHANGE_NAME, "fanout");
        String queueName = channel.queueDeclare().getQueue(); // Exclusive, auto-delete queue
        channel.queueBind(queueName, EXCHANGE_NAME, ""); // Bind to fanout exchange

        System.out.println(" [*] Waiting for messages in queue '" + queueName + "'. To exit press CTRL+C");

        DeliverCallback deliverCallback = (consumerTag, delivery) -> {
            String message = new String(delivery.getBody(), StandardCharsets.UTF_8);
            System.out.println(" [x] Received '" + message + "'");
        };
        // Auto-ack set to true for simplicity in this example
        channel.basicConsume(queueName, true, deliverCallback, consumerTag -> {});
    }
}

Fanout in Action

When you ran two consumers and then the publisher, you observed the core concept of Fanout: every active consumer received a copy of the message.

This is because each consumer had its own unique queue, and both queues were bound to the same fanout_logs exchange. The exchange simply duplicated the message to all bound queues.

This makes Fanout ideal for scenarios like real-time logging, notifications, or broadcasting updates.

Test Your Knowledge

Time for a quick check on Fanout exchanges!

Fanout Recap

You've successfully learned about the Fanout exchange!

  • It implements the Pub/Sub pattern.
  • It broadcasts messages to all bound queues.
  • It ignores routing keys when distributing messages.
  • It's perfect for scenarios where multiple consumers need to receive the same message.

Next, we'll explore the Direct exchange, which uses routing keys for more precise message delivery!

免费开始

用 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 课程适合初学者到高级学习者,你可以从这里开始或从头开始,按照自己的节奏学习。 这是第 1 节课,共 4 节。

「用于发布/订阅的扇出交换器」课时需要多长时间?

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

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

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

此课程中的所有课时

  1. 用于发布/订阅的扇出交换器
  2. 用于路由的直连交换器
  3. 灵活路由的主题交换器
  4. 默认交换器与隐式绑定
← 返回 RabbitMQ Messaging & Async Systems