您的位置:首页 > 运维架构

RabbitMQ 学习笔记(五):Topics

2017-07-14 12:44 357 查看

Topics

在之前的笔记中,我们改进了日志系统。我们使用了一个 direct 交换器,而不是使用一个只能够进行虚拟广播的 fanout 交换器,实现了有选择性地接收日志。

尽管使用 direct 交换器改进了我们的系统,但它仍然有局限性——它不能基于多个标准进行选择指定路线(routing)。

在我们的日志系统中,我们可能希望订阅的不仅是基于性质的日志,还需要基于发出日志的源。您可能从名为 syslog 的unix工具中了解这个概念,该工具根据日志的性质(info /warning/ crit…)和设施(auth / cron / kern…)来选择路线发送日志。

这将给我们带来很大的灵活性——我们可能想要倾听来自“cron”的一些关键错误,但也有来自“kern”的所有日志信息。

为了在我们的日志系统中实现这一点,我们需要了解一个更复杂的交换器 —— topic交换器

topic交换器

发送到topic交换器的消息不能有一个任意的routing key——它必须是一个单词列表,由点“.”分隔。单词可以是任何东西,但通常它们指定与消息连接的一些特性。一些有效的routing key示例:“stock.usd.nyse”, “nyse.vmw”, “quick.orange.rabbit”。在routing key中可以有任意数量的单词,最多可以达到255个字节

binding key也必须以相同的形式出现。topic 交换器背后的逻辑类似于direct 交换器——发送带有特定 routing_key 的消息将被传送到绑定具有匹配 binding key 的所有队列。但是有两种重要的 binding key 特殊情况:

* 可以代替一个词。

# 可以替代零个或多个单词。

在如下一个例子中比较容易解释:



在本例中,我们将发送所有描述动物的消息。消息将发送带有三个单词(两个点)的 routing key。routing key中的第一个词将描述速度、第二个颜色和第三个物种:“ speed . color . species ”。

我们创建了三个绑定,Q1与binding key “*.orange.*”绑定,Q2与binding key “*.*.rabbit” 和 “lazy.#”绑定

这些绑定可以概括为:

Q1对所有的橙色动物都感兴趣。

Q2想接收关于兔子的一切,关于 lazy 动物的一切。

一个routing key为“quick. orange.rabbit”的消息,将被送到两个队列中。消息“lazy.orange.elephant”也会被送到两个队列中。另一方面,“quick.orange.fox”只会进入第一个队列,然后“lazy.brown.fox”只会进入第二个队列。“lazy.pink.rabbit”只会被发送到第二个队列,即使它匹配两个绑定。“quick.brown.fox”不匹配任何绑定,所以它将被丢弃。

如果我们违反规定,用一个或四个单词作为routing key发送信息,比如“orange”或“quick. orange . male.rabbit”,会发生什么情况呢?答案是:这些消息不匹配任何绑定,将被丢弃。

另一方面,“lazy.orange.male.rabbit”尽管它有四个单词,但它将匹配最后一个绑定,并将被传递到第二个队列。

Topic 交换器

Topic 交换器是强大的,可以表现得像其他交换器一样(通用)。

队列与“#”绑定键绑定时,它将接收所有消息,不管routing key是什么,比如在fanout 交换器中。

特殊字符“*”和“#”不用于绑定时,topic 交换器就会像 direct 一样。

编译运行

我们将在我们的日志系统中使用一个topic交换器。我们将从一个工作假设开始,即日志的routing key将有两个单词:“<设施> . <严重性>”。

代码几乎和前一笔记中的一样。

EmitLogTopic.java 的代码如下:

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

public class EmitLogTopic {

private static final String EXCHANGE_NAME = "topic_logs";

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

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

channel.exchangeDeclare(EXCHANGE_NAME, BuiltinExchangeType.TOPIC);

String routingKey = getRouting(argv);
String message = getMessage(argv);

channel.basicPublish(EXCHANGE_NAME, routingKey, null, message.getBytes("UTF-8"));
System.out.println(" [x] Sent '" + routingKey + "':'" + message + "'");

}
catch  (Exception e) {
e.printStackTrace();
}
finally {
if (connection != null) {
try {
connection.close();
}
catch (Exception ignore) {}
}
}
}

private static String getRouting(String[] strings){
if (strings.length < 1)
return "anonymous.info";
return strings[0];
}

private static String getMessage(String[] strings){
if (strings.length < 2)
return "Hello World!";
return joinStrings(strings, " ", 1);
}

private static String joinStrings(String[] strings, String delimiter, int startIndex) {
int length = strings.length;
if (length == 0 ) return "";
if (length < startIndex ) return "";
StringBuilder words = new StringBuilder(strings[startIndex]);
for (int i = startIndex + 1; i < length; i++) {
words.append(delimiter).append(strings[i]);
}
return words.toString();
}
}


ReceiveLogsTopic.java 的代码如下:

import com.rabbitmq.client.*;

import java.io.IOException;

public class ReceiveLogsTopic {

private static final String EXCHANGE_NAME = "topic_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, BuiltinExchangeType.TOPIC);
String queueName = channel.queueDeclare().getQueue();

if (argv.length < 1) {
System.err.println("Usage: ReceiveLogsTopic [binding_key]...");
System.exit(1);
}

for (String bindingKey : argv) {
channel.queueBind(queueName, EXCHANGE_NAME, bindingKey);
}

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

Consumer consumer = new DefaultConsumer(channel) {
@Override
public void handleDelivery(String consumerTag, Envelope envelope,
AMQP.BasicProperties properties, byte[] body) throws IOException {
String message = new String(body, "UTF-8");
System.out.println(" [x] Received '" + envelope.getRoutingKey() + "':'" + message + "'");
}
};
channel.basicConsume(queueName, true, consumer);
}
}


编译并运行示例,包括笔记一中的类路径

编译:

javac -cp $CP ReceiveLogsTopic.java EmitLogTopic.java

接收所有的日志:

java -cp $CP ReceiveLogsTopic “#”

接收来自设备“kern”的所有日志:

java -cp $CP ReceiveLogsTopic “kern.*”

或者,如果你只想收到一些“critical”的日志:

java -cp $CP ReceiveLogsTopic “*.critical”

你可以创建多重绑定:

java -cp $CP ReceiveLogsTopic “kern.” “.critical”

并发送routing key为“kern.critical”的日志:

java -cp $CP EmitLogTopic “kern.critical” “A critical kernel error”

请注意,代码没有对 routing key 或 binding key 做出任何约束,你可以使用两个以上的 routing key 参数。

相关链接

rabbitmq-c++(SimpleAmqpClient) 笔记代码系列:

rabbitmq-c++(SimpleAmqpClient) 笔记代码一

rabbitmq-c++(SimpleAmqpClient) 笔记代码二

rabbitmq-c++(SimpleAmqpClient) 笔记代码三

rabbitmq-c++(SimpleAmqpClient) 笔记代码四

rabbitmq-c++(SimpleAmqpClient) 笔记代码五

rabbitmq-c++(SimpleAmqpClient) 笔记代码六

RabbitMQ 学习笔记系列:

RabbitMQ 学习笔记(一):简单介绍及”Hello World”

RabbitMQ 学习笔记(二):work queues

RabbitMQ 学习笔记(三):Publish/Subscribe

RabbitMQ 学习笔记(四):Routing

RabbitMQ 学习笔记(五):Topics

RabbitMQ 学习笔记(六):RPC
内容来自用户分享和网络整理,不保证内容的准确性,如有侵权内容,可联系管理员处理 点击这里给我发消息
标签: