第一个例子-消费者接收消息

创建一个新控制台项目,同样使用nuget安装rabbitmq的驱动

static void Main(string[] args)
{
    var factory = new ConnectionFactory();
    factory.HostName = "localhost";
    factory.UserName = "guest";
    factory.Password = "guest";

    using (var connection = factory.CreateConnection())
    {
        using (var channel = connection.CreateModel())
        {
            channel.QueueDeclare("hello", false, false, false, null);

            var consumer = new EventingBasicConsumer(channel);
            channel.BasicConsume("hello", false, consumer);
            consumer.Received += (model, ea) =>
            {
                var body = ea.Body;
                var message = Encoding.UTF8.GetString(body);
                Console.WriteLine("已接收: {0}", message);
            };
            Console.ReadLine();
        }
    }
}Code language: JavaScript (javascript)

和发送一样,首先需要定义连接,然后声明消息队列。要接收消息,需要定义一个Consumer,然后在接收消息的事件中处理数据。

既然消息已经被接收了,那我们再来看queue的内容:

可见,消息中的内容在接收之后已被删除。

我们对比之前的

已处理的1

未消费的0

看另外一个例子

//第一步:创建连接connection
using (var connection = factory.CreateConnection())
{
    //第二步:创建通道channel
    using (var channel = connection.CreateModel())
    {
        //第三步:声明队列queue
        channel.QueueDeclare(queue: "myqueue",
                             durable: true,
                             exclusive: false,
                             autoDelete: false,
                             arguments: null);
        //第四步:定义消费者
        var consumer = new EventingBasicConsumer(channel);
        consumer.Received += (model, ea) =>
        {
            var body = ea.Body;
            var message = Encoding.UTF8.GetString(body);
            Console.WriteLine($"接受到消息【{message}】");
        };
        Console.WriteLine("消费者准备就绪....");
        //第五步:处理消息
        channel.BasicConsume(queue: "myqueue",
                             autoAck: true,
                             consumer: consumer);
        Console.ReadLine();
    }
}Code language: JavaScript (javascript)

注意:上边的代码在生产者和消费者的代码中都声明了exchange和queue,这主要是为了让这两个程序可以按任意顺序启动。比如:我们只在生产者代码中定义了exchange和queue,却先启动消费者,这会造成消费者找不到自己需要的exchange和queue(出现404错误)。实际开发中创建exchange/queue、绑定队列以及设置routingKey这些工作,都可以通过WebUI管理界面或者使用Rabbitmq Control工具完成。

QueueDeclare方法用于声明队列,ExchangeDeclare用于声明交换机。我们在使用这两个方法声明时,可以设置队列和交换机的属性,如queue的名字、长度限制、exchange是否持久化、交换机类型等。

QueueDeclare方法的参数如下:

queue:队列名字;

durable:是否持久化。设置为true时,队列信息保存在rabbitmq的内置数据库中,服务器重启时队列也会恢复(注意:重启后队列内部的消息不会恢复,怎么实现消息持久化以后会详细介绍);

exclusive:是否排外。设置为true时只有首次声明该队列的Connection可以访问,其他Connection不能访问该队列;且在Connection断开时,队列会被删除(即使durable设置为true也会被删除);

autoDelete:是否自动删除。设置为true时,表示在最后一条使用该队列的连接(Connection)断开时,将自动删除这个队列;

arguments:设置队列的一些其它属性,为Dictionary<string,object>类型,下表总结了arguments中可以设置的常用属性。

ExchangeDeclare方法详解

exchange:交换机名字。

type:交换机类型。exchange有direct、fanout、topic、header四种类型,在下一篇会详细介绍;

durable:是否持久化。设置为true时,交换机信息保存在rabbitmq的内置数据库中,服务器重启时交换机信息也会恢复;

autoDelete:是否自动删除。设置为true时,表示在最后一条使用该交换机的连接(Connection)断开时,自动删除这个exchange;

arguments:其他的一些参数,类型为Dictionary<string,object>。

发表回复

您的邮箱地址不会被公开。 必填项已用 * 标注