热门标签 | HotTags
当前位置:  开发笔记 > 编程语言 > 正文

9、rabbitmq

主要说rabbitmq,kafka简单看一下,rocketmq对.net没有官方的支持,所以暂不介绍,如果在业务中有用到,用java封装一下然后对外开放api吧,其实rocketm

主要说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; }
// 获取或设置 .
public string ExchangeType { get; set; }
// 获取或设置发布消息的属性。
public IDictionary Headers { get; set; }
// 获取或设置定义 Exchange 时的参数。
public IDictionary 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 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();
var optiOns= (IOptions)Options.Create(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);
}

  

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


///


/// 定义 RabbitMQ 订阅者的抽象类。
///

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().CreateLogger(this.GetType());
this.Initialize();
}
protected RabbitMQConsumerContext Context { get; set; }
///
/// 初始化
///

protected abstract void Initialize();
public void Subscribe()
{
try
{
using (var scope = this.ServiceProvider.CreateScope())
{
var cOnfigProvider= scope.ServiceProvider.GetRequiredService();
var optiOns= scope.ServiceProvider.GetService>().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 Headers { get; set; }
public ReadOnlyMemory 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 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();
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();

 

增加一个订阅,也就是实现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

 


推荐阅读
  • Java验证码——kaptcha的使用配置及样式
    本文介绍了如何使用kaptcha库来实现Java验证码的配置和样式设置,包括pom.xml的依赖配置和web.xml中servlet的配置。 ... [详细]
  • 基于layUI的图片上传前预览功能的2种实现方式
    本文介绍了基于layUI的图片上传前预览功能的两种实现方式:一种是使用blob+FileReader,另一种是使用layUI自带的参数。通过选择文件后点击文件名,在页面中间弹窗内预览图片。其中,layUI自带的参数实现了图片预览功能。该功能依赖于layUI的上传模块,并使用了blob和FileReader来读取本地文件并获取图像的base64编码。点击文件名时会执行See()函数。摘要长度为169字。 ... [详细]
  • HDU 2372 El Dorado(DP)的最长上升子序列长度求解方法
    本文介绍了解决HDU 2372 El Dorado问题的一种动态规划方法,通过循环k的方式求解最长上升子序列的长度。具体实现过程包括初始化dp数组、读取数列、计算最长上升子序列长度等步骤。 ... [详细]
  • 本文讨论了Alink回归预测的不完善问题,指出目前主要针对Python做案例,对其他语言支持不足。同时介绍了pom.xml文件的基本结构和使用方法,以及Maven的相关知识。最后,对Alink回归预测的未来发展提出了期待。 ... [详细]
  • 本文讨论了如何优化解决hdu 1003 java题目的动态规划方法,通过分析加法规则和最大和的性质,提出了一种优化的思路。具体方法是,当从1加到n为负时,即sum(1,n)sum(n,s),可以继续加法计算。同时,还考虑了两种特殊情况:都是负数的情况和有0的情况。最后,通过使用Scanner类来获取输入数据。 ... [详细]
  • 本文介绍了OC学习笔记中的@property和@synthesize,包括属性的定义和合成的使用方法。通过示例代码详细讲解了@property和@synthesize的作用和用法。 ... [详细]
  • Mac OS 升级到11.2.2 Eclipse打不开了,报错Failed to create the Java Virtual Machine
    本文介绍了在Mac OS升级到11.2.2版本后,使用Eclipse打开时出现报错Failed to create the Java Virtual Machine的问题,并提供了解决方法。 ... [详细]
  • 本文讲述了作者通过点火测试男友的性格和承受能力,以考验婚姻问题。作者故意不安慰男友并再次点火,观察他的反应。这个行为是善意的玩人,旨在了解男友的性格和避免婚姻问题。 ... [详细]
  • 1,关于死锁的理解死锁,我们可以简单的理解为是两个线程同时使用同一资源,两个线程又得不到相应的资源而造成永无相互等待的情况。 2,模拟死锁背景介绍:我们创建一个朋友 ... [详细]
  • 《数据结构》学习笔记3——串匹配算法性能评估
    本文主要讨论串匹配算法的性能评估,包括模式匹配、字符种类数量、算法复杂度等内容。通过借助C++中的头文件和库,可以实现对串的匹配操作。其中蛮力算法的复杂度为O(m*n),通过随机取出长度为m的子串作为模式P,在文本T中进行匹配,统计平均复杂度。对于成功和失败的匹配分别进行测试,分析其平均复杂度。详情请参考相关学习资源。 ... [详细]
  • 动态规划算法的基本步骤及最长递增子序列问题详解
    本文详细介绍了动态规划算法的基本步骤,包括划分阶段、选择状态、决策和状态转移方程,并以最长递增子序列问题为例进行了详细解析。动态规划算法的有效性依赖于问题本身所具有的最优子结构性质和子问题重叠性质。通过将子问题的解保存在一个表中,在以后尽可能多地利用这些子问题的解,从而提高算法的效率。 ... [详细]
  • 高质量SQL书写的30条建议
    本文提供了30条关于优化SQL的建议,包括避免使用select *,使用具体字段,以及使用limit 1等。这些建议是基于实际开发经验总结出来的,旨在帮助读者优化SQL查询。 ... [详细]
  • 在project.properties添加#Projecttarget.targetandroid-19android.library.reference.1..Sliding ... [详细]
  • 本文介绍了lua语言中闭包的特性及其在模式匹配、日期处理、编译和模块化等方面的应用。lua中的闭包是严格遵循词法定界的第一类值,函数可以作为变量自由传递,也可以作为参数传递给其他函数。这些特性使得lua语言具有极大的灵活性,为程序开发带来了便利。 ... [详细]
  • CentOS 7部署KVM虚拟化环境之一架构介绍
    本文介绍了CentOS 7部署KVM虚拟化环境的架构,详细解释了虚拟化技术的概念和原理,包括全虚拟化和半虚拟化。同时介绍了虚拟机的概念和虚拟化软件的作用。 ... [详细]
author-avatar
优美rosner_704
这个家伙很懒,什么也没留下!
PHP1.CN | 中国最专业的PHP中文社区 | DevBox开发工具箱 | json解析格式化 |PHP资讯 | PHP教程 | 数据库技术 | 服务器技术 | 前端开发技术 | PHP框架 | 开发工具 | 在线工具
Copyright © 1998 - 2020 PHP1.CN. All Rights Reserved | 京公网安备 11010802041100号 | 京ICP备19059560号-4 | PHP1.CN 第一PHP社区 版权所有