RabbitMQ简单使用


一、windows下安装RabbitMQ

参考https://www.cnblogs.com/ericli-ericli/p/5902270.html

RabbitMQ是一个在AMQP协议标准基础上完整的,可服用的企业消息系统。它遵循Mozilla Public License开源协议,采用 Erlang 实现的工业级的消息队列(MQ)服务器,Rabbit MQ 是建立在Erlang OTP平台上。

1、安装Erlang

官网下载,一步步安装即可,添加系统环境变量(我是修改过安装目录);

2、安装RabbitMQ

官网下载rabbitmq-server-3.7.15(安装时版本),安装即可;

  • 默认安装的RabbitMQ监听的端口号:5672
  • 错误提示:Elang could not be detected...步骤1、安装Erlang是必须步骤;

3、常用命令

windows下使用cmd(以管理员身份运行)进入安装目录下,本人使用的安装路径,D:\Program Files\RabbitMQ Server\rabbitmq_server-3.7.15\sbin

  • 启用web控制台,rabbitmq-plugins enable rabbitmq_management
  • 启动服务,net start rabbitmq
  • 停止服务,net stop rabbitmq
  • 查看已有用户及其角色,rabbitmqctl list_users
  • 新增用户,rabbitmqctl add_user username password
  • 给用户设置超级管理员角色,rabbitmqctl set_user_tags username administrator
    • 超级管理员(administrator)
    • 监控者(monitoring)
    • 策略制定者(policymaker)
    • 普通管理者(management)
    • 其他的

访问web控制台,http://服务器ip:15672,默认用户名密码都是guest

二、.net core中使用RabbitMQ

参考:https://www.cnblogs.com/yan7/p/9498685.html

直接使用nuget安装RabbitMQ.Client

生产者producter

using System;
using System.Text;
using RabbitMQ.Client;

namespace RabbitMQ_Send
{
    class Program
    {
        // 生成者producter
        static void Main(string[] args)
        {
            Console.WriteLine("--->rabbitmq-producter start");
            // 创建连接工厂
            IConnectionFactory factory = new ConnectionFactory
            {
                HostName = "127.0.0.1",// ip
                Port = 5672,// 端口号
                UserName = "zsan",
                Password = "123"
            };
            // 创建连接
            var connection = factory.CreateConnection();
            // 创建通道
            var channel = connection.CreateModel();
            string queueName = "hellorabbitmq";
            // 声明一个队列
            channel.QueueDeclare(
                queue: queueName,// 消息队列名称
                durable: false,// 是否缓存
                exclusive: false,
                autoDelete: false,
                arguments: null
            );

            Console.WriteLine("\n--->rabbitmq连接成功,请输入消息,输入exit退出");
            string input = string.Empty;
            do
            {
                input = Console.ReadLine();
                var sendBytes = Encoding.UTF8.GetBytes(input);
                channel.BasicPublish(exchange: "", routingKey: queueName, basicProperties: null, body: sendBytes);// 在一对一中routingKey必须和queue一致
            } while (input.Trim().ToLower() != "exit");
            channel.Close();
            connection.Close();
        }
    }
}

消费者consumer

using System;
using System.Text;
using System.Threading;
using RabbitMQ.Client;
using RabbitMQ.Client.Events;

namespace RabbitMQ_Receive
{
    class Program
    {
        // 消费者consumer
        static void Main(string[] args)
        {
            Console.WriteLine("--->rabbitmq-consumer start");
            // 创建连接工厂
            IConnectionFactory factory = new ConnectionFactory
            {
                HostName = "127.0.0.1",
                Port = 5672,
                UserName = "zsan",
                Password = "123"
            };
            // 创建连接
            var connection = factory.CreateConnection();
            // 创建通道
            var channel = connection.CreateModel();
            string queueName = "hellorabbitmq";
            // 声明队列
            channel.QueueDeclare(
                queue: queueName,
                durable: false,
                exclusive: false,
                autoDelete: false,
                arguments: null
            );
            var consumer = new EventingBasicConsumer(channel);

            consumer.Received += (ch, ea) =>
            {
                byte[] message = ea.Body;
                Console.WriteLine("\n--->接收到信息:" + Encoding.UTF8.GetString(message));
                int random = new Random().Next(1, 6) * 1000;
                Console.WriteLine(string.Format("收到该消息[{0}] 延迟{1}s发送回执", ea.DeliveryTag, random+""));
                Thread.Sleep(random);
                Console.WriteLine(string.Format("已发送回执{0}", ea.DeliveryTag));
                //确认该消息已被消费
                channel.BasicAck(ea.DeliveryTag, false);
            };
            // 启动消费者
            channel.BasicConsume(queue: queueName, autoAck: false, consumer: consumer);
            Console.ReadKey();
        }
    }
}

注:

  • 问题1:错误提示OperationInterruptedException: The AMQP operation was interrupted: AMQP close-reason, initiated by Peer, code=530, text="NOT_ALLOWED - access to vhost '/' refused for user 'zsan'", classId=10, methodId=40, cause=

    用户角色权限问题,登录web控制台,给用户set permission