前言
RabbitMQ是工作中经常使用的一个开源消息队列,它是通过Erlang语言来开发的,RabbitMQ中有一个非常常见且重要的概念叫交换器(Exchange),所有生产者提交的消息都是通过Exchange进行接收,然后通过特定的路由策略转发对特定的队列(Queue)中,然后消费者通过绑定特定的队列,将队列中的信息进行消费处理。我们如果是通过RabbitMQ的客户端来进行创建Exchange的话,一定需要你选择一个特定的Exchange类型的,这个时候如果不了解每一种类型的特性,就很容易选择错误,造成整个业务流程中的异常,接下来我们就来梳理一下Exchange对应的四种类型:direct、fanout、headers、topic分别的特点和使用场景
消息接受与发送入门
在详细梳理Exchange之前,我们先简单的了解看一下RabbitMq通过最简单的方式,来实现针对一条消息的发送与接手,在这个代码运行前,确保已经在本地安装了一个可以正常运行的RabbitMq,最好还有一个还有了一个可以RabbitMq的管理后台,用来更加方便的查看当前管理和维护的Exchange、Queue等信息,如下图所示:
依赖配置
如果你是采用maven来进行代码依赖管理的话,首先需要将RabbitMq的依赖的pom文件加入到项目中,我的代码使用的是最新版的 5.11.0,如下图所示
<dependency>
<groupId>com.rabbitmq</groupId>
<artifactId>amqp-client</artifactId>
<version>5.11.0</version>
</dependency>
发送消息到指定队列(Queue)
import com.rabbitmq.client.Channel;
import com.rabbitmq.client.Connection;
import com.rabbitmq.client.ConnectionFactory;
import java.nio.charset.StandardCharsets;
public class PublishTest {
public static final String TEST_QUEUE = "test-queue";
public static void main(String[] args) {
ConnectionFactory factory = new ConnectionFactory();
factory.setHost("localhost");
try (Connection connection = factory.newConnection();
Channel channel = connection.createChannel()) {
channel.queueDeclare(TEST_QUEUE, false, false, false, null);
String message = "你好,朋友 - 1 ";
channel.basicPublish("",TEST_QUEUE,null, message.getBytes(StandardCharsets.UTF_8));
System.out.println(" 发送消息 : " + message);
} catch (Exception e) {
e.printStackTrace();
}
}
}
消费指定队列消息
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 ConsumerTest {
public static void main(String[] args) throws Exception{
ConnectionFactory factory = new ConnectionFactory();
factory.setHost("localhost");
Connection connection = factory.newConnection() ;
Channel channel = connection.createChannel();
System.out.println(" 等待消息到达,通过 CTRL+C 关闭");
DeliverCallback deliverCallback = (consumerTag, delivery) -> {
String message = new String(delivery.getBody(), StandardCharsets.UTF_8);
System.out.println(" 接收到的消息 : " + message );
};
channel.basicConsume(PublishTest.TEST_QUEUE, true, deliverCallback, consumerTag -> {
});
}
}
Exchange类型介绍
direct
Exchange默认的类型就是direct,该类型的Exchange消息分发机制非常简单,需要和routingKey配合使用,首先在配置阶段,direct类型的Exchange通过一个rountingKey和Queue进行绑定,绑定完成后,消息通过使用Exchange名称和rountingKey进行发送,消息到达RabbitMq服务器后,rountingKey需要和Queue绑定的rountingKey完全相同才会进行消息的路由,路由到对应的Queue,流程如下图所示:
代码功能验证
发送Direct类型Exchange消息
import com.rabbitmq.client.BuiltinExchangeType;
import com.rabbitmq.client.Channel;
import com.rabbitmq.client.Connection;
import com.rabbitmq.client.ConnectionFactory;
import java.nio.charset.StandardCharsets;
public class DirectExchange {
public static final String DIRECT_EXCHANGE_NAME = "direct-exchange-test";
public static final String ROUTING_KEY_A = "A.D";
public static final String ROUTING_KEY_B = "B.D";
public static final String ROUTING_KEY_C = "C";
public static final String ROUTING_KEY_D = "D";
public static final String ROUTING_KEY_E = "A.*";
public static final String MESSAGE_A = "你好,朋友-A";
public static final String MESSAGE_B = "你好,朋友-B";
public static final String MESSAGE_C = "你好,朋友-C";
public static final String MESSAGE_D = "你好,朋友-D";
public static final String MESSAGE_E = "你好,朋友-E";
public static void main(String[] args) {
ConnectionFactory factory = new ConnectionFactory();
factory.setHost("localhost");
try (Connection connection = factory.newConnection();
Channel channel = connection.createChannel()) {
channel.exchangeDeclare(DIRECT_EXCHANGE_NAME, BuiltinExchangeType.DIRECT);
channel.basicPublish(DIRECT_EXCHANGE_NAME, ROUTING_KEY_A, null, MESSAGE_A.getBytes(StandardCharsets.UTF_8));
channel.basicPublish(DIRECT_EXCHANGE_NAME, ROUTING_KEY_B, null, MESSAGE_B.getBytes(StandardCharsets.UTF_8));
channel.basicPublish(DIRECT_EXCHANGE_NAME, ROUTING_KEY_C, null, MESSAGE_C.getBytes(StandardCharsets.UTF_8));
channel.basicPublish(DIRECT_EXCHANGE_NAME, ROUTING_KEY_D, null, MESSAGE_D.getBytes(StandardCharsets.UTF_8));
channel.basicPublish(DIRECT_EXCHANGE_NAME, ROUTING_KEY_E, null, MESSAGE_E.getBytes(StandardCharsets.UTF_8));
} catch (Exception e) {
e.printStackTrace();
}
}
}
消费Direct类型Exchange消息
import com.example.druiddemo.publisher.DirectExchange;
import com.rabbitmq.client.*;
import java.nio.charset.StandardCharsets;
public class DirectConsumerTest {
public static final String DIRECT_QUEUE_NAME = "direct-queue-test";
public static final String BINDING_KEY_A = "A.*";
public static final String BINDING_KEY_B = "*.D";
public static final String BINDING_KEY_C = "C";
public static void main(String[] args) throws Exception {
ConnectionFactory factory = new ConnectionFactory();
factory.setHost("localhost");
Connection connection = factory.newConnection() ;
Channel channel = connection.createChannel();
channel.queueDeclare(DIRECT_QUEUE_NAME, false, false, false, null);
channel.queueBind(DIRECT_QUEUE_NAME, DirectExchange.DIRECT_EXCHANGE_NAME, BINDING_KEY_A);
channel.queueBind(DIRECT_QUEUE_NAME, DirectExchange.DIRECT_EXCHANGE_NAME, BINDING_KEY_B);
channel.queueBind(DIRECT_QUEUE_NAME, DirectExchange.DIRECT_EXCHANGE_NAME, BINDING_KEY_C);
System.out.println(" 等待消息到达,通过 CTRL+C 关闭");
DeliverCallback deliverCallback = (consumerTag, delivery) -> {
String message = new String(delivery.getBody(), StandardCharsets.UTF_8);
System.out.println(" 接收到的消息 : " + message );
};
channel.basicConsume(DIRECT_QUEUE_NAME, true, deliverCallback, consumerTag -> {
});
}
}
fanout
Exchange中的fanout类型是我们一般使用的偏少的一种类型,采用的是广播的形式来进行消息转发处理,同时该模式下的不需要使用RountingKey,只要是与该Exchange绑定的所有的队列均可以收到发送的消息,该模式主要用于当某个事件发生是对相关业务的通知,这样只要想要知道该事件的所有业务只需要绑定到这个Exchange即可平等收到转发的消息
代码功能验证
发送Fanout类型Exchange消息
import com.rabbitmq.client.BuiltinExchangeType;
import com.rabbitmq.client.Channel;
import com.rabbitmq.client.Connection;
import com.rabbitmq.client.ConnectionFactory;
import java.nio.charset.StandardCharsets;
public class FanoutExchange {
public static final String FANOUT_EXCHANGE_NAME = "fanout-exchange-test";
public static final String MESSAGE_A = "你好,朋友-A";
public static final String MESSAGE_B = "你好,朋友-B";
public static final String MESSAGE_C = "你好,朋友-C";
public static final String MESSAGE_D = "你好,朋友-D";
public static final String MESSAGE_E = "你好,朋友-E";
public static void main(String[] args) {
ConnectionFactory factory = new ConnectionFactory();
factory.setHost("localhost");
try (Connection connection = factory.newConnection();
Channel channel = connection.createChannel()) {
channel.exchangeDeclare(FANOUT_EXCHANGE_NAME, BuiltinExchangeType.FANOUT);
channel.basicPublish(FANOUT_EXCHANGE_NAME, "", null, MESSAGE_A.getBytes(StandardCharsets.UTF_8));
channel.basicPublish(FANOUT_EXCHANGE_NAME, "", null, MESSAGE_B.getBytes(StandardCharsets.UTF_8));
channel.basicPublish(FANOUT_EXCHANGE_NAME, "", null, MESSAGE_C.getBytes(StandardCharsets.UTF_8));
channel.basicPublish(FANOUT_EXCHANGE_NAME, "", null, MESSAGE_D.getBytes(StandardCharsets.UTF_8));
channel.basicPublish(FANOUT_EXCHANGE_NAME, "", null, MESSAGE_E.getBytes(StandardCharsets.UTF_8));
} catch (Exception e) {
e.printStackTrace();
}
}
}
消费Fanout类型Exchange消息
import com.example.druiddemo.publisher.FanoutExchange;
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 FanoutConsumerTest {
public static final String FANOUT_QUEUE_NAME = "fanout-queue-test";
public static void main(String[] args) throws Exception {
ConnectionFactory factory = new ConnectionFactory();
factory.setHost("localhost");
Connection connection = factory.newConnection() ;
Channel channel = connection.createChannel();
channel.queueDeclare(FANOUT_QUEUE_NAME, false, false, false, null);
channel.queueBind(FANOUT_QUEUE_NAME, FanoutExchange.FANOUT_EXCHANGE_NAME, "");
System.out.println("fanout-queue-test 等待消息到达,通过 CTRL+C 关闭");
DeliverCallback deliverCallback = (consumerTag, delivery) -> {
String message = new String(delivery.getBody(), StandardCharsets.UTF_8);
System.out.println("fanout-queue-test 接收到的消息 : " + message );
};
channel.basicConsume(FANOUT_QUEUE_NAME, true, deliverCallback, consumerTag -> {
});
}
}
消费Fanout类型Exchange消息
import com.example.druiddemo.publisher.FanoutExchange;
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 FanoutConsumerTwoTest {
public static final String FANOUT_QUEUE_NAME_COPY = "fanout-queue-test-copy";
public static void main(String[] args) throws Exception {
ConnectionFactory factory = new ConnectionFactory();
factory.setHost("localhost");
Connection connection = factory.newConnection() ;
Channel channel = connection.createChannel();
channel.queueDeclare(FANOUT_QUEUE_NAME_COPY, false, false, false, null);
channel.queueBind(FANOUT_QUEUE_NAME_COPY, FanoutExchange.FANOUT_EXCHANGE_NAME, "");
System.out.println("fanout-queue-test-copy 等待消息到达,通过 CTRL+C 关闭");
DeliverCallback deliverCallback = (consumerTag, delivery) -> {
String message = new String(delivery.getBody(), StandardCharsets.UTF_8);
System.out.println("fanout-queue-test-copy 接收到的消息 : " + message );
};
channel.basicConsume(FANOUT_QUEUE_NAME_COPY, true, deliverCallback, consumerTag -> {
});
}
}
headers
该类型的Exchange是我们在日常开发中使用的最少的一种类型,该类型的Exhange转发消息并不是通过RountingKey来进行路由,而是通过它的绑定参数(Arguments)来进行转发的,它的配置位置是配置信息的headers中,该模式下的消息转发路径如下:
代码功能验证
发送Headers类型Exchange消息
import com.rabbitmq.client.*;
import java.nio.charset.StandardCharsets;
import java.util.Hashtable;
import java.util.Map;
public class HeadersExchange {
public static final String HEADERS_EXCHANGE_NAME = "headers-exchange-test";
public static final String MESSAGE_A = "你好,朋友-A";
public static final String MESSAGE_B = "你好,朋友-B";
public static final String MESSAGE_C = "你好,朋友-C";
public static void main(String[] args) {
ConnectionFactory factory = new ConnectionFactory();
factory.setHost("localhost");
try (Connection connection = factory.newConnection();
Channel channel = connection.createChannel()) {
channel.exchangeDeclare(HEADERS_EXCHANGE_NAME, BuiltinExchangeType.HEADERS);
Map<String, Object> headersOne = new Hashtable<>();
headersOne.put("t", "s");
headersOne.put("m", "k");
AMQP.BasicProperties basicProperties = new AMQP.BasicProperties().builder().headers(headersOne).build();
System.out.println(basicProperties.toString());
channel.basicPublish(HEADERS_EXCHANGE_NAME, "", basicProperties, MESSAGE_A.getBytes(StandardCharsets.UTF_8));
Map<String, Object> headersTwo = new Hashtable<>();
headersTwo.put("t", "s");
AMQP.BasicProperties basicPropertiesTwo = new AMQP.BasicProperties().builder().headers(headersTwo).build();
System.out.println(basicPropertiesTwo.toString());
channel.basicPublish(HEADERS_EXCHANGE_NAME, "", basicPropertiesTwo, MESSAGE_B.getBytes(StandardCharsets.UTF_8));
Map<String, Object> headersThree = new Hashtable<>();
headersThree.put("m", "k");
AMQP.BasicProperties basicPropertiesThree = new AMQP.BasicProperties().builder().headers(headersThree).build();
System.out.println(basicPropertiesThree.toString());
channel.basicPublish(HEADERS_EXCHANGE_NAME, "", basicPropertiesThree, MESSAGE_C.getBytes(StandardCharsets.UTF_8));
} catch (Exception e) {
e.printStackTrace();
}
}
}
消费Headers类型Exchange消息
import com.example.druiddemo.publisher.HeadersExchange;
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;
import java.util.Hashtable;
import java.util.Map;
public class HeadersConsumerTest {
public static final String HEADERS_QUEUE_NAME = "headers-queue-test";
public static void main(String[] args) throws Exception {
ConnectionFactory factory = new ConnectionFactory();
factory.setHost("localhost");
Connection connection = factory.newConnection() ;
Channel channel = connection.createChannel();
channel.queueDeclare(HEADERS_QUEUE_NAME, false, false, false, null);
Map<String, Object> headers = new Hashtable<>();
headers.put("t", "s");
headers.put("m", "k");
headers.put("x-match", "all");
channel.queueBind(HEADERS_QUEUE_NAME, HeadersExchange.HEADERS_EXCHANGE_NAME, "", headers);
System.out.println("headers-queue-test 等待消息到达,通过 CTRL+C 关闭");
DeliverCallback deliverCallback = (consumerTag, delivery) -> {
String message = new String(delivery.getBody(), StandardCharsets.UTF_8);
System.out.println("headers-queue-test 接收到的消息 : " + message );
};
channel.basicConsume(HEADERS_QUEUE_NAME, true, deliverCallback, consumerTag -> {
});
}
}
消费Headers类型Exchange消息
import com.example.druiddemo.publisher.HeadersExchange;
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;
import java.util.Hashtable;
import java.util.Map;
public class HeadersConsumerTwoTest {
public static final String HEADERS_QUEUE_NAME = "headers-queue-test-two";
public static void main(String[] args) throws Exception {
ConnectionFactory factory = new ConnectionFactory();
factory.setHost("localhost");
Connection connection = factory.newConnection() ;
Channel channel = connection.createChannel();
channel.queueDeclare(HEADERS_QUEUE_NAME, false, false, false, null);
Map<String, Object> headers = new Hashtable<>();
headers.put("t", "s");
headers.put("m", "k");
headers.put("x-match", "any");
channel.queueBind(HEADERS_QUEUE_NAME, HeadersExchange.HEADERS_EXCHANGE_NAME, "", headers);
System.out.println("headers-queue-test-two 等待消息到达,通过 CTRL+C 关闭");
DeliverCallback deliverCallback = (consumerTag, delivery) -> {
String message = new String(delivery.getBody(), StandardCharsets.UTF_8);
System.out.println("headers-queue-test-two 接收到的消息 : " + message );
};
channel.basicConsume(HEADERS_QUEUE_NAME, true, deliverCallback, consumerTag -> {
});
}
}
消费Headers类型Exchange消息
import com.example.druiddemo.publisher.HeadersExchange;
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;
import java.util.Hashtable;
import java.util.Map;
public class HeadersConsumerThreeTest {
public static final String HEADERS_QUEUE_NAME = "headers-queue-test-three";
public static void main(String[] args) throws Exception {
ConnectionFactory factory = new ConnectionFactory();
factory.setHost("localhost");
Connection connection = factory.newConnection() ;
Channel channel = connection.createChannel();
channel.queueDeclare(HEADERS_QUEUE_NAME, false, false, false, null);
Map<String, Object> headers = new Hashtable<>();
headers.put("m", "k");
headers.put("x-match", "any");
channel.queueBind(HEADERS_QUEUE_NAME, HeadersExchange.HEADERS_EXCHANGE_NAME, "", headers);
System.out.println("headers-queue-test-two 等待消息到达,通过 CTRL+C 关闭");
DeliverCallback deliverCallback = (consumerTag, delivery) -> {
String message = new String(delivery.getBody(), StandardCharsets.UTF_8);
System.out.println("headers-queue-test-two 接收到的消息 : " + message );
};
channel.basicConsume(HEADERS_QUEUE_NAME, true, deliverCallback, consumerTag -> {
});
}
}
topic
topic类型的Exchange必须要配置一个RountingKey,该RountingKey由字符组成,通过dots(也就是 . )来进行分割,RountingKey可以是任意字符来组成,但是上线不可以超过 255 bytes;Queue通过Binding Key来将Exchange与之绑定,Exchange下RountingKey与Binding Key相匹配的消息被发送到对应的绑定的Queue中,然后消息被当前Queue的消费者进行消费处理,在Binding Key中有两个非常常用并且特殊的字符需要注意:
*:与RountingKey进行匹配时,可以精确的代替一个字符
#:与RountingKey进行匹配时,可以代替0个或多个字符
代码功能验证
发送Topic类型Exchange消息
import com.rabbitmq.client.BuiltinExchangeType;
import com.rabbitmq.client.Channel;
import com.rabbitmq.client.Connection;
import com.rabbitmq.client.ConnectionFactory;
import java.nio.charset.StandardCharsets;
public class TopicExchange {
public static final String TOPIC_EXCHANGE_NAME = "topic-exchange-test";
public static final String ROUTING_KEY_A = "A.D";
public static final String ROUTING_KEY_B = "B.D";
public static final String ROUTING_KEY_C = "C";
public static final String ROUTING_KEY_D = "D";
public static final String MESSAGE_A = "你好,朋友-A";
public static final String MESSAGE_B = "你好,朋友-B";
public static final String MESSAGE_C = "你好,朋友-C";
public static final String MESSAGE_D = "你好,朋友-D";
public static void main(String[] args) {
ConnectionFactory factory = new ConnectionFactory();
factory.setHost("localhost");
try (Connection connection = factory.newConnection();
Channel channel = connection.createChannel()) {
channel.exchangeDeclare(TOPIC_EXCHANGE_NAME, BuiltinExchangeType.TOPIC);
channel.basicPublish(TOPIC_EXCHANGE_NAME, ROUTING_KEY_A, null, MESSAGE_A.getBytes(StandardCharsets.UTF_8));
channel.basicPublish(TOPIC_EXCHANGE_NAME, ROUTING_KEY_B, null, MESSAGE_B.getBytes(StandardCharsets.UTF_8));
channel.basicPublish(TOPIC_EXCHANGE_NAME, ROUTING_KEY_C, null, MESSAGE_C.getBytes(StandardCharsets.UTF_8));
channel.basicPublish(TOPIC_EXCHANGE_NAME, ROUTING_KEY_D, null, MESSAGE_D.getBytes(StandardCharsets.UTF_8));
} catch (Exception e) {
e.printStackTrace();
}
}
}
消费Topic类型Exchange消息
import com.example.druiddemo.publisher.TopicExchange;
import com.rabbitmq.client.*;
import java.nio.charset.StandardCharsets;
public class TopicConsumerTest {
public static final String TOPIC_QUEUE_NAME = "topic-queue-test";
public static final String BINDING_KEY_A = "A.*";
public static final String BINDING_KEY_B = "*.D";
public static final String BINDING_KEY_C = "C";
public static void main(String[] args) throws Exception {
ConnectionFactory factory = new ConnectionFactory();
factory.setHost("localhost");
Connection connection = factory.newConnection() ;
Channel channel = connection.createChannel();
channel.exchangeDeclare(TopicExchange.TOPIC_EXCHANGE_NAME, BuiltinExchangeType.TOPIC);
channel.queueDeclare(TOPIC_QUEUE_NAME, false, false, false, null);
channel.queueBind(TOPIC_QUEUE_NAME, TopicExchange.TOPIC_EXCHANGE_NAME, BINDING_KEY_C);
System.out.println(" 等待消息到达,通过 CTRL+C 关闭");
DeliverCallback deliverCallback = (consumerTag, delivery) -> {
String message = new String(delivery.getBody(), StandardCharsets.UTF_8);
System.out.println(" 接收到的消息 : " + message );
};
channel.basicConsume(TOPIC_QUEUE_NAME, true, deliverCallback, consumerTag -> {
});
}
}
总结
通过上面的介绍,我们已经清楚的了解了每一种Exchange类型的特点和他们的使用场景,在我们的日常使用中一定要先明确我们的使用的特点,然后再去根据实际情况去选择最合适的Exchange来进行创建和使用,这样才是最为高效和合理的,其实在上面的topic类型的Exchange是可以通过一些变通来实现direct类型和fanout类型的Exhange的:如果topic类型的Exchange绑定的Queue的绑定Key是 “#” ,则该Exhange下的所有消息都会转发到该Queue;如果topic类型的Exchange绑定的Queue的绑定Key中不包含 "* " 或 “#”,则该Exchange和direct类型的使用场景就一致了
弦外音
这个地方主要分析的是RabbitMq的一些重要组成部分的详细的使用,如果你想要深入的理解这里的每一种类型的使用,还是需要亲自在编写相关的代码,然后进行相关的类型的运行和实际效果的观察,如果你的电脑还没有安装RabbitMq,那么你可以查看这篇文章,这里有安装的完整步骤和安装过程中容易碰到的几个异常,希望可以帮助你更好的使用RabbitMq:https://community.jiguang.cn/article/465057
0条评论