DotNet Core中使用RabbitMQ

原文: DotNet Core中使用RabbitMQ

  上一篇随笔记录到RabbitMQ的安装,安装完成,咱们就开始使用吧。html

RabbitMQ简介算法

  AMQP,即Advanced Message Queuing Protocol,高级消息队列协议,是应用层协议的一个开放标准,为面向消息的中间件设计。消息中间件主要用于组件之间的解耦,消息的发送者无需知道消息使用者的存在,反之亦然。安全

  AMQP的主要特征是面向消息、队列、路由(包括点对点和发布/订阅)、可靠性、安全。
RabbitMQ是一个开源的AMQP实现,服务器端用Erlang语言编写,支持多种客户端,如:Python、Ruby、.NET、Java、JMS、C、PHP、ActionScript、XMPP、STOMP等,支持AJAX。用于在分布式系统中存储转发消息,在易用性、扩展性、高可用性等方面表现不俗。服务器

  RabbitMQ提供了可靠的消息机制、跟踪机制和灵活的消息路由,支持消息集群和分布式部署。适用于排队算法、秒杀活动、消息分发、异步处理、数据同步、处理耗时任务、CQRS等应用场景。异步

DotNet Core使用RabbitMQ分布式

经过nuget安装:https://www.nuget.org/packages/RabbitMQ.Client/ui

定义生产者:this

复制代码
//建立链接工厂
ConnectionFactory factory = new ConnectionFactory
{
    UserName = "guest",//用户名
    Password = "guest",//密码
    HostName = "127.0.0.1"//rabbitmq ip
};

//建立链接
var connection = factory.CreateConnection();
//建立通道
var channel = connection.CreateModel();
//声明一个队列
channel.QueueDeclare("hello", false, false, false, null);

Console.WriteLine("\nRabbitMQ链接成功,请输入消息,输入exit退出!");

string input;
do
{
    input = Console.ReadLine();

    var sendBytes = Encoding.UTF8.GetBytes(input);
    //发布消息
    channel.BasicPublish("", "hello", null, sendBytes);

} while (input.Trim().ToLower() != "exit");
channel.Close();
connection.Close();
复制代码

定义消费者:spa

复制代码
//建立链接工厂
ConnectionFactory factory = new ConnectionFactory
{
    UserName = "guest",//用户名
    Password = "guest",//密码
    HostName = "127.0.0.1"//rabbitmq ip
};

//建立链接
var connection = factory.CreateConnection();
//建立通道
var channel = connection.CreateModel();

//事件基本消费者
EventingBasicConsumer consumer = new EventingBasicConsumer(channel);

//接收到消息事件
consumer.Received += (ch, ea) =>
{
    var message = Encoding.UTF8.GetString(ea.Body);
    Console.WriteLine($"收到消息: {message}");
    //确认该消息已被消费
    channel.BasicAck(ea.DeliveryTag, false);
};
//启动消费者 设置为手动应答消息
channel.BasicConsume("hello", false, consumer);
Console.WriteLine("消费者已启动");
Console.ReadKey();
channel.Dispose();
connection.Close();
复制代码

演示以下:设计

启动了一个生产者,两个消费者,能够看见两个消费者都能接收到消息,消息投递到哪一个消费者是由RabbitMQ决定的。

RabbitMQ消费失败的处理

  RabbitMQ采用消息应答机制,即消费者收到一个消息以后,须要发送一个应答,而后RabbitMQ才会将这个消息从队列中删除,若是消费者在消费过程当中出现异常,断开链接切没有发送应答,那么RabbitMQ会将这个消息从新投递。

咱们来修改一下消费者的代码:

复制代码
 //接收到消息事件
 consumer.Received += (ch, ea) =>
 {
     var message = Encoding.UTF8.GetString(ea.Body);

     Console.WriteLine($"收到消息: {message}");

     Console.WriteLine($"收到该消息[{ea.DeliveryTag}] 延迟10s发送回执");
     Thread.Sleep(10000);
     //确认该消息已被消费
     channel.BasicAck(ea.DeliveryTag, false);
     Console.WriteLine($"已发送回执[{ea.DeliveryTag}]");
 };
复制代码

演示以下:

从图中能够看出,设置了消息应答延迟10s,若是在这10s中,该消费者断开了链接,那么消息会被RabbitMQ从新投递。

使用RabbitMQ的Exchange

前面的例子,咱们能够看到生产者将消息投递到Queue中,实际上这种方式在RabbitMQ中永远都不会发生的。实际的状况是,生产者将消息发送到Exchange(交换器),下图中的X,由Exchange(交换器)将消息路由到一个或多个Queue中(或者丢弃)。

 

AMQP协议中的核心思想就是生产者和消费者隔离,生产者从不直接将消息发送给队列。生产者一般不知道是否一个消息会被发送到队列中,只是将消息发送到一个交换机。先由Exchange来接收,而后Exchange按照特定的策略转发到Queue进行存储。同理,消费者也是如此。Exchange 就相似于一个交换机,转发各个消息分发到相应的队列中。

Exchange Types(交换器类型)

RabbitMQ经常使用的Exchange Type有Fanout、Direct、Topic、Headers这四种

一、Fanout:

  这种类型的Exchange路由规则很是简单,它会把全部发送到该Exchange的消息路由到全部与它绑定的Queue中,这时Routing key不起做用

 

 

 

Fanout Exchange 不须要处理RouteKey 。只须要简单的将队列绑定到exchange 上。这样发送到exchange的消息都会被转发到与该交换机绑定的全部队列上。相似子网广播,每台子网内的主机都得到了一份复制的消息。

因此,Fanout Exchange 转发消息是最快的。

为了演示效果,定义了两个队列,分别为hello1,hello2,每一个队列都拥有一个消费者。

复制代码
static void Main(string[] args)
{
    string exchangeName = "TestFanoutChange";
    string queueName1 = "hello1";
    string queueName2 = "hello2";
    string routeKey = "";

    //建立链接工厂
    ConnectionFactory factory = new ConnectionFactory
    {
        UserName = "guest",//用户名
        Password = "guest",//密码
        HostName = "127.0.0.1"//rabbitmq ip
    };

    //建立链接
    var connection = factory.CreateConnection();
    //建立通道
    var channel = connection.CreateModel();

    //定义一个Direct类型交换机
    channel.ExchangeDeclare(exchangeName, ExchangeType.Fanout, false, false, null);

    //定义队列1
    channel.QueueDeclare(queueName1, false, false, false, null);
    //定义队列2
    channel.QueueDeclare(queueName2, false, false, false, null);

    //将队列绑定到交换机
    channel.QueueBind(queueName1, exchangeName, routeKey, null);
    channel.QueueBind(queueName2, exchangeName, routeKey, null);

    //生成两个队列的消费者
    ConsumerGenerator(queueName1);
    ConsumerGenerator(queueName2);


    Console.WriteLine($"\nRabbitMQ链接成功,\n\n请输入消息,输入exit退出!");

    string input;
    do
    {
        input = Console.ReadLine();

        var sendBytes = Encoding.UTF8.GetBytes(input);
        //发布消息
        channel.BasicPublish(exchangeName, routeKey, null, sendBytes);

    } while (input.Trim().ToLower() != "exit");
    channel.Close();
    connection.Close();
}
复制代码
复制代码
 /// <summary>
 /// 根据队列名称生成消费者
 /// </summary>
 /// <param name="queueName"></param>
 static void ConsumerGenerator(string queueName)
 {
     //建立链接工厂
     ConnectionFactory factory = new ConnectionFactory
     {
         UserName = "guest",//用户名
         Password = "guest",//密码
         HostName = "127.0.0.1"//rabbitmq ip
     };

     //建立链接
     var connection = factory.CreateConnection();
     //建立通道
     var channel = connection.CreateModel();

     //事件基本消费者
     EventingBasicConsumer consumer = new EventingBasicConsumer(channel);

     //接收到消息事件
     consumer.Received += (ch, ea) =>
     {
         var message = Encoding.UTF8.GetString(ea.Body);

         Console.WriteLine($"Queue:{queueName}收到消息: {message}");
         //确认该消息已被消费
         channel.BasicAck(ea.DeliveryTag, false);
     };
     //启动消费者 设置为手动应答消息
     channel.BasicConsume(queueName, false, consumer);
     Console.WriteLine($"Queue:{queueName},消费者已启动");
 }
复制代码

运行效果以下:

二、Direct

  这种类型的Exchange路由规则也很简单,它会把消息路由到哪些binding key与routingkey彻底匹配的Queue中。

 

   Direct模式,可使用rabbitMQ自带的Exchange:default Exchange 。因此不须要将Exchange进行任何绑定(binding)操做 。消息传递时,RouteKey必须彻底匹配,才会被队列接收,不然该消息会被抛弃。

复制代码
static void Main(string[] args)
{
    string exchangeName = "TestChange";
    string queueName = "hello";
    string routeKey = "helloRouteKey";

    //建立链接工厂
    ConnectionFactory factory = new ConnectionFactory
    {
        UserName = "guest",//用户名
        Password = "guest",//密码
        HostName = "127.0.0.1"//rabbitmq ip
    };

    //建立链接
    var connection = factory.CreateConnection();
    //建立通道
    var channel = connection.CreateModel();

    //定义一个Direct类型交换机
    channel.ExchangeDeclare(exchangeName, ExchangeType.Direct, false, false, null);

    //定义一个队列
    channel.QueueDeclare(queueName, false, false, false, null);

    //将队列绑定到交换机
    channel.QueueBind(queueName, exchangeName, routeKey, null);

    Console.WriteLine($"\nRabbitMQ链接成功,Exchange:{exchangeName},Queue:{queueName},Route:{routeKey},\n\n请输入消息,输入exit退出!");

    string input;
    do
    {
        input = Console.ReadLine();

        var sendBytes = Encoding.UTF8.GetBytes(input);
        //发布消息
        channel.BasicPublish(exchangeName, routeKey, null, sendBytes);

    } while (input.Trim().ToLower() != "exit");
    channel.Close();
    connection.Close();
复制代码

运行效果以下:

三、Topic

  这种类型的Exchange的路由规则支持 binding key 和 routing key 的模糊匹配,会把消息路由到知足条件的Queue。 binding key 中能够存在两种特殊字符 *与 #,用于作模糊匹配,其中 * 用于匹配一个单词,# 用于匹配0个或多个单词,单词以符号“.”为分隔符。

  以上图中的配置为例,routingKey=”quick.orange.rabbit”的消息会同时路由到Q1与Q2,routingKey=”lazy.orange.fox”的消息会路由到Q1与Q2,routingKey=”lazy.brown.fox”的消息会路由到Q2,routingKey=”lazy.pink.rabbit”的消息会路由到Q2(只会投递给Q2一次,虽然这个routingKey与Q2的两个bindingKey都匹配);routingKey=”quick.brown.fox”、routingKey=”orange”、routingKey=”quick.orange.male.rabbit”的消息将会被丢弃,由于它们没有匹配任何bindingKey。

  因此,Topic Exchange使用很是灵活。
复制代码
static void Main(string[] args)
{
    string exchangeName = "TestTopicChange";
    string queueName = "hello";
    string routeKey = "TestRouteKey.*";

    //建立链接工厂
    ConnectionFactory factory = new ConnectionFactory
    {
        UserName = "guest",//用户名
        Password = "guest",//密码
        HostName = "127.0.0.1"//rabbitmq ip
    };

    //建立链接
    var connection = factory.CreateConnection();
    //建立通道
    var channel = connection.CreateModel();

    //定义一个Direct类型交换机
    channel.ExchangeDeclare(exchangeName, ExchangeType.Topic, false, false, null);

    //定义队列1
    channel.QueueDeclare(queueName, false, false, false, null);

    //将队列绑定到交换机
    channel.QueueBind(queueName, exchangeName, routeKey, null);



    Console.WriteLine($"\nRabbitMQ链接成功,\n\n请输入消息,输入exit退出!");

    string input;
    do
    {
        input = Console.ReadLine();

        var sendBytes = Encoding.UTF8.GetBytes(input);
        //发布消息
        channel.BasicPublish(exchangeName, "TestRouteKey.one", null, sendBytes);

    } while (input.Trim().ToLower() != "exit");
    channel.Close();
    connection.Close();
}
复制代码

运行效果以下:

 四、Headers

  这种类型的Exchange不依赖于 routing key 与 binding key 的匹配规则来路由消息,而是根据发送的消息内容中的 headers 属性进行匹配。

参考:

  官网:https://www.rabbitmq.com/tutorials/tutorial-one-dotnet.html

    http://www.javashuo.com/article/p-yplqyfft-o.html

    https://www.jianshu.com/p/e55e971aebd8

相关文章
相关标签/搜索