zoukankan      html  css  js  c++  java
  • 9、rabbitmq

    主要说rabbitmq,kafka简单看一下,rocketmq对.net 没有官方的支持,所以暂不介绍,如果在业务中有用到,用java封装一下然后对外开放api吧,其实rocketmq还是更好一些,因为没有依赖项,而rocketmq需要erlang和socat

    队列概念,转自简书

    https://www.jianshu.com/p/9a0e9ffa17dd

    1、rabbitmq

    官网

    https://www.rabbitmq.com/

    .net 支持

    https://www.rabbitmq.com/dotnet.html

    2、安装

    安装erlang

    yum install erlang

    安装 socat

    yun install socat

    安装rabbitmq

    rpm -Uvh https://github.com/rabbitmq/rabbitmq-server/releases/download/v3.8.4/rabbitmq-server-3.8.4-1.el8.noarch.rpm

    3、启动

    service rabbitmq-server start #启动

    service rabbitmq-server stop #停止

    service rabbitmq-server restart #重启

    chkconfig rabbitmq-server on #开机自启

    rabbitmq-plugins enable rabbitmq_management #开启web管理界面 需要重启

    rabbitmqctl  change_password  admin admin #创建用户admin,密码admin

    rabbitmqctl set_user_tags admin administrator # 给admin用户设置权限

    这里注意需要处理一下端口权限,web管理界面的端口是15672

    安装启动部分不多说了

    4、做什么,怎么做

    队列可以做很多事情,比如用来处理并发,电商的秒杀,简单的理解就是减库存操作,比如库存是10,秒杀活动在1ms产生了10000个请求,如果用往常的减库存的操作,那么肯定会出各种各样的错误,当然可以用redis锁或者其他办法来解决,但是队列确实是一种方案。还有就是做解耦,用队列的发布订阅模式,一个消息,谁想用谁就订阅。

    这个模式有两个角色,发布者和消费者

    发布者:创建一条消息

    订阅者:接收并执行发布者发布的消息,也就是消费者(杠精滚)

    发布者发布了一条命令【去给我买瓶水】,如果有人订阅了,那么订阅的人就叫做消费者(手下办事儿的),消费者收到命令,就去买水了,买完了会给发布者一个眼神,告诉发布者水买回来了

    有多个订阅者怎么办?发布者也是恃强凌弱的,虽然多个订阅者,但是他也只发送给其中一个订阅者

    没有订阅者怎么办?那买水的指令就一直挂在那了,直到有订阅者订阅,才会被消费掉

    下面开始实现发布订阅模式

    apollo里把需要用到的参数定义好

    rabbmitmq中把topic加上

    先定义一些基础类

    mq的上下文实体,说白了就是发布、消费消息时候用到的参数

    public class RabbitMQContext
        {
            // 获取或设置交换机名称。
            public string ExchangeName { get; set; }
            // 获取或设置 RoutingKey 。
            public string RoutingKey { get; set; }
            // 获取或设置队列名称。
            public string QueueName { get; set; }
            // 获取或设置 <see cref="global::RabbitMQ.Client.ExchangeType"/>.
            public string ExchangeType { get; set; }
            // 获取或设置发布消息的属性。<see cref="global::RabbitMQ.Client.IBasicProperties.Headers"/>
            public IDictionary<string, object> Headers { get; set; }
            // 获取或设置定义 Exchange 时的参数。
            public IDictionary<string, object> ExchangeArgs { get; set; }
            public bool IsDelayMessage { get; set; }
        }
        public class RabbitMQConsumerContext : RabbitMQContext
        {
        }
        public class RabbitMQProducerContext : RabbitMQContext
        {
            public RabbitMQMessage Body { get; set; }
        }
        public class RabbitMQMessage
        {
            private object _message;
    
            [JsonProperty(PropertyName = "message")]
            public object Message
            {
                get
                {
                    return this._message;
                }
    
                set
                {
                    this._message = value;
                    this.MessageType = value == null
                        ? string.Empty
                        : value.GetType().FullName;
                }
            }
            [JsonProperty("messageType")]
            public string MessageType { get; set; }
        }
    

      

    实现发布者功能

    定义发布者接口

    public interface IRabbitMQPublisher
    {
    	// 发布消息到消息队列。
    	bool Publish(RabbitMQProducerContext context);
    }
    

    实现发布者接口

    public class RabbitMQPublisher : IRabbitMQPublisher
        {
            private readonly IConfiguration _configProvider;
            private readonly RabbitMQOptions _options;
            public RabbitMQPublisher(IConfiguration configProvider, IOptions<RabbitMQOptions> options)
            {
                this._configProvider = configProvider;
                this._options = options.Value;
            }
            public bool Publish(RabbitMQProducerContext context)
            {
                if (string.IsNullOrWhiteSpace(context.ExchangeName))
                {
                    context.ExchangeName = _configProvider["RabbitMQ:Default:Exchange"];
                }
                var factory = new ConnectionFactory
                {
                    HostName = this._options.Host,
                    Port = this._options.Port,
                    UserName = this._options.UserName,
                    Password = this._options.Password
                };
                using (var connection = factory.CreateConnection())
                using (var channel = connection.CreateModel())
                {
                    var body = Encoding.UTF8.GetBytes(JsonConvert.SerializeObject(context.Body));
                    channel.ExchangeDeclare(context.ExchangeName, context.ExchangeType, durable: true, autoDelete: false, arguments: context.ExchangeArgs);
                    if (!string.IsNullOrWhiteSpace(context.QueueName))
                    {
                        channel.QueueDeclare(context.QueueName, durable: true, exclusive: false, autoDelete: false, arguments: context.ExchangeArgs);
                        channel.QueueBind(context.QueueName, context.ExchangeName, context.RoutingKey);
                    }
                    IBasicProperties basicProperties = null;
                    if (context.Headers != null && context.Headers.Any())
                    {
                        basicProperties = channel.CreateBasicProperties();
                        basicProperties.Headers = context.Headers;
                    }
                    channel.BasicPublish(context.ExchangeName, context.RoutingKey, basicProperties: basicProperties, body: body);
                }
                return true;
            }
        }
    

    发布一个消息

    [HttpGet]
    public JsonResult RabbitMqPublish([FromQuery] string paramStr)
    {
    	var context = new RabbitMQProducerContext
    	{
    		Body = new RabbitMQMessage
    		{
    			Message = JsonConvert.SerializeObject(new { paramStr })
    		},
    		QueueName = "new_born",
    		RoutingKey = "new_born",
    		ExchangeType = ExchangeType.Topic
    	};
    	var optionValue = this._configuration.GetSection("RabbitMQ").Get<RabbitMQOptions>();
    	var options = (IOptions<RabbitMQOptions>)Options.Create<RabbitMQOptions>(optionValue);
    	var publisher = new RabbitMQPublisher(this._configuration, options);
    	publisher.Publish(context);
    	return Json(new { code = 1, message = "发布成功" });
    }
    

    执行结束后,可以看到队列里面已经有待消费的数据了

    下面实现订阅、消费功能

    订阅

    需要在系统启动的时候,就去订阅,所以依然实现 IHostedService 接口

    定义消费者/订阅者接口

    public interface IMessageQueueConsumer
    {
    	// 订阅消息。
    	void Subscribe();
    	// 取消订阅消息。
    	void Unsubscribe();
    	// 通知生产者此消息已被消费
    	void BasicAck(ulong deliveryTag, bool multiple);
    }
    

      

     定义消费者/订阅者抽象类

    /// <summary>
        /// 定义 RabbitMQ 订阅者的抽象类。
        /// </summary>
        public abstract class RabbitMQConsumer : IMessageQueueConsumer
        {
            private IConnection _connection;
            private IModel _channel;
            private string _consumerTag;
            protected readonly IServiceProvider ServiceProvider;
            protected readonly ILogger Logger;
    
            protected RabbitMQConsumer(IServiceProvider serviceProvider)
            {
                this.ServiceProvider = serviceProvider ?? throw new ArgumentNullException(nameof(serviceProvider));
                this.Logger = this.ServiceProvider.GetRequiredService<ILoggerFactory>().CreateLogger(this.GetType());
                this.Initialize();
            }
            protected RabbitMQConsumerContext Context { get; set; }
            /// <summary>
            /// 初始化 <see cref="RabbitMQContext"/> 。
            /// </summary>
            protected abstract void Initialize();
            public void Subscribe()
            {
                try
                {
                    using (var scope = this.ServiceProvider.CreateScope())
                    {
                        var configProvider = scope.ServiceProvider.GetRequiredService<IConfiguration>();
                        var options = scope.ServiceProvider.GetService<IOptionsSnapshot<RabbitMQOptions>>().Value;
                        var factory = new ConnectionFactory
                        {
                            HostName = options.Host,
                            Port = options.Port,
                            UserName = options.UserName,
                            Password = options.Password
                        };
                        this._connection = factory.CreateConnection();
                        this._channel = this._connection.CreateModel();
                        // 定义 Exchange ,持久化,不自动删除此 Exchange
                        this._channel.ExchangeDeclare(this.Context.ExchangeName, this.Context.ExchangeType, durable: true, autoDelete: false, arguments: this.Context.ExchangeArgs);
                        // 如果队列名为空,则队列默认不持久化,非独占,自动删除
                        if (string.IsNullOrWhiteSpace( this.Context.QueueName))
                        {
                            this.Context.QueueName = this._channel.QueueDeclare(durable: false, exclusive: false, autoDelete: true).QueueName;
                        }
                        else
                        {
                            this._channel.QueueDeclare(this.Context.QueueName, durable: true, exclusive: false, autoDelete: false);
                        }
                        this._channel.QueueBind(this.Context.QueueName, this.Context.ExchangeName, this.Context.RoutingKey);
                        if (options.Qos > 0)
                        {
                            this._channel.BasicQos(0, (ushort)options.Qos, true);
                        }
                        var consumer = new EventingBasicConsumer(this._channel);
                        consumer.Received += this.OnConsumerReceived;
    
                        this._consumerTag = this._channel.BasicConsume(this.Context.QueueName, false, consumer);
                    }
                }
                catch (Exception e)
                {
                    this.Logger.LogError(e, $"订阅消息异常:{ JsonConvert.SerializeObject( this.Context)}");
                }
            }
            public void Unsubscribe()
            {
                try
                {
                    this._channel.BasicCancel(this._consumerTag);
                    this._channel.Close();
                    this._channel.Dispose();
    
                    this._connection.Close(TimeSpan.FromSeconds(5));
                    this._connection.Dispose();
                }
                catch (Exception e)
                {
                    this.Logger.LogError(e, $"取消订阅异常:{this.GetType().FullName}");
                }
            }
            public void BasicAck(ulong deliveryTag, bool multiple)
            {
                this._channel.BasicAck(deliveryTag, multiple);
            }
            protected void OnConsumerReceived(object sender, BasicDeliverEventArgs e)
            {
                if (!(sender is EventingBasicConsumer consumer)) return;
                try
                {
                    var deliverEventArgs = new DeliverEventArgs
                    {
                        Body = e.Body,
                        ConsumerTag = e.ConsumerTag,
                        DeliveryTag = e.DeliveryTag,
                        Exchange = e.Exchange,
                        Redelivered = e.Redelivered,
                        RoutingKey = e.RoutingKey
                    };
                    if (e.BasicProperties?.IsHeadersPresent() ?? false)
                    {
                        deliverEventArgs.Headers = e.BasicProperties.Headers;
                    }
                    this.OnConsumerReceivedAsync(sender, deliverEventArgs).ConfigureAwait(false).GetAwaiter().GetResult();
                }
                catch (Exception ex)
                {
                    this.Logger.LogError(ex, "处理 RabbitMQ 消息异常。");
                }
            }
            protected abstract Task OnConsumerReceivedAsync(object sender, DeliverEventArgs e);
        }
        public class DeliverEventArgs : EventArgs
        {
            public IDictionary<string, object> Headers { get; set; }
            public ReadOnlyMemory<byte> Body { get; set; }
            public string ConsumerTag { get; set; }
            public ulong DeliveryTag { get; set; }
            public string Exchange { get; set; }
            public bool Redelivered { get; set; }
            public string RoutingKey { get; set; }
        }
        public static class ExchangeType
        {
            public const string Direct = "direct";
            public const string Fanout = "fanout";
            public const string Headers = "headers";
            public const string Topic = "topic";
        }
    

      

    实现IHostedService 实现系统启动时,自动执行系统内实现了IMessageQueueConsumer 接口服务的订阅方法

    public class MessageQueueHostedService : IHostedService
        {
            private readonly IServiceProvider _serviceProvider;
            private readonly RabbitMQOptions _options;
            public MessageQueueHostedService(IServiceProvider serviceProvider, IOptions<RabbitMQOptions> options)
            {
                this._serviceProvider = serviceProvider;
                this._options = options.Value;
            }
            public Task StartAsync(CancellationToken cancellationToken)
            {
                if (!this._options.Enable) return Task.CompletedTask;
                var consumers = this._serviceProvider.GetServices<IMessageQueueConsumer>();
                if (consumers.Any())
                {
                    foreach (var consumer in consumers)
                    {
                        consumer.Subscribe();
                    } 
                }
                return Task.CompletedTask;
            }
            public Task StopAsync(CancellationToken cancellationToken)
            {
                if (!this._options.Enable) return Task.CompletedTask;
                return Task.CompletedTask;
            }
        }
    

     

    注入到ServiceCollection

    services.AddSingleton<IHostedService, MessageQueueHostedService>();

    增加一个订阅,也就是实现RabbitMQConsumer抽象类

    public class NewBornConsumer : RabbitMQConsumer
    {
    	private IServiceProvider _serviceProvider;
    	public NewBornConsumer(IServiceProvider serviceProvider) : base(serviceProvider)
    	{
    		_serviceProvider = serviceProvider;
    	}
    	protected override void Initialize()
    	{
    		this.Context = new RabbitMQConsumerContext
    		{
    			ExchangeType = ExchangeType.Topic,
    			ExchangeName = "new_born",
    			QueueName = "new_born",
    			RoutingKey = "new_born",
    		};
    	}
    	protected async override Task OnConsumerReceivedAsync(object sender, DeliverEventArgs e)
    	{
    		await Task.CompletedTask;             
    		var message = Encoding.UTF8.GetString(e.Body.ToArray());
    		var data = JToken.Parse(message)["message"].ToString();
    		var paramStr = JToken.Parse(data)["paramStr"];
    		Logger.LogInformation($"我被消费拉 ~~~~参数是:{paramStr}");
    		this.BasicAck(e.DeliveryTag, false);
    	}
    }
    

     

    此时我们F5启动系统,可以看到已经有一个消费者了

    通过api增加一条消息,可以看到消费信息

    其他的队列,像kafka,ActiveMQ,rocketmq,大同小异,应用的场景也不同,多翻翻博客吧

    队列在实际的生产环境中用到的地方还挺多,比如我们做的预约挂号、到诊、进销存出库这些操作,理论上都应该用队列来处理,或者日志收集这种的,具体问题具体分析吧

    ab.exe 并发测试工具

    使用方法:.ab.exe -n 50 -c 50 http://localhost:8862/GameApi/RabbitMqPublish?paramStr=33333

    参数说明: -n 请求多少次  -c 每次并发量多少

    下载地址:

    http://files.cnblogs.com/files/gossip/ab.zip

  • 相关阅读:
    Python自然语言处理资料库
    Solr 中 Schema 结构说明
    solr 高亮显示
    HTML URL 编码
    IDEA java开发 Restful 风格的WebService
    Intellij IDEA中使用log4j日志
    IntelliJ IDEA java开发 WebService
    java 实现poi方式读取word文件内容
    Ubuntu安装nodeJS
    Ubuntu 系统下 mongodb 安装和配置
  • 原文地址:https://www.cnblogs.com/ares-core/p/13026030.html
Copyright © 2011-2022 走看看