C# 学习笔记:RabbitMQ队列使用多线程提高消费吞吐率

本文主要是介绍C# 学习笔记:RabbitMQ队列使用多线程提高消费吞吐率,希望对大家解决编程问题提供一定的参考价值,需要的开发者们随着小编来一起学习吧!

一、引言
使用工作队列的一个好处就是它能够并行的处理队列。如果堆积了很多任务,我们只需要添加更多的工作者(workers)就可以了,扩展很简单。本例使用多线程来创建多信道并绑定队列,达到多workers的目的。

二、示例
2.1、环境准备
在NuGet上安装RabbitMQ.Client。

在这里插入图片描述

2.2、工厂类
添加一个工厂类RabbitMQFactory:
 /// <summary>/// 多路复用技术(Multiplexing)目的:为了避免创建多个TCP而造成系统资源的浪费和超载,从而有效地利用TCP连接。/// </summary>public static class RabbitMQFactory{private static IConnection sharedConnection;private static int ChannelCount { get; set; }private static readonly object _locker = new object();public static IConnection SharedConnection{get{if (ChannelCount >= 1000){if (sharedConnection != null && sharedConnection.IsOpen){sharedConnection.Close();}sharedConnection = null;ChannelCount = 0;}if (sharedConnection == null){lock (_locker){if (sharedConnection == null){sharedConnection = GetConnection();ChannelCount++;}}}return sharedConnection;}}private static IConnection GetConnection(){var factory = new ConnectionFactory{HostName = "192.168.2.242",UserName = "hello",Password = "world",Port = AmqpTcpEndpoint.UseDefaultPort,//5672VirtualHost = ConnectionFactory.DefaultVHost,//使用默认值:"/"Protocol = Protocols.DefaultProtocol,AutomaticRecoveryEnabled = true};return factory.CreateConnection();}}
2.3、主窗体

在这里插入图片描述

代码如下:
public partial class RabbitMQMultithreading : Form
{public delegate void ListViewDelegate<T>(T obj);public RabbitMQMultithreading(){InitializeComponent();}/// <summary>/// ShowMessage重载/// </summary>/// <param name="msg"></param>private void ShowMessage(string msg){if (InvokeRequired){BeginInvoke(new ListViewDelegate<string>(ShowMessage), msg);}else{ListViewItem item = new ListViewItem(new string[] { DateTime.Now.ToString("yyyy/MM/dd HH:mm:ss ffffff"), msg });lvwMsg.Items.Insert(0, item);}}/// <summary>/// ShowMessage重载/// </summary>/// <param name="format"></param>/// <param name="args"></param>private void ShowMessage(string format, params object[] args){if (InvokeRequired){BeginInvoke(new MethodInvoker(delegate (){ListViewItem item = new ListViewItem(new string[] { DateTime.Now.ToString("yyyy/MM/dd HH:mm:ss ffffff"), string.Format(format, args) });lvwMsg.Items.Insert(0, item);}));}else{ListViewItem item = new ListViewItem(new string[] { DateTime.Now.ToString("yyyy/MM/dd HH:mm:ss ffffff"), string.Format(format, args) });lvwMsg.Items.Insert(0, item);}}/// <summary>/// 生产者/// </summary>/// <param name="sender"></param>/// <param name="e"></param>private void btnSend_Click(object sender, EventArgs e){int messageCount = 100;var factory = new ConnectionFactory{HostName = "192.168.2.242",UserName = "hello",Password = "world",Port = AmqpTcpEndpoint.UseDefaultPort,//5672VirtualHost = ConnectionFactory.DefaultVHost,//使用默认值:"/"Protocol = Protocols.DefaultProtocol,AutomaticRecoveryEnabled = true};using (var connection = factory.CreateConnection()){using (var channel = connection.CreateModel()){channel.QueueDeclare(queue: "hello", durable: true, exclusive: false, autoDelete: false, arguments: null);string message = "Hello World";var body = Encoding.UTF8.GetBytes(message);for (int i = 1; i <= messageCount; i++){channel.BasicPublish(exchange: "", routingKey: "hello", basicProperties: null, body: body);ShowMessage($"Send {message}");}}}}/// <summary>/// 消费者/// </summary>/// <param name="sender"></param>/// <param name="e"></param>private async void btnReceive_Click(object sender, EventArgs e){Random random = new Random();int rallyNumber = random.Next(1, 1000);int channelCount = 0;await Task.Run(() =>{try{int asyncCount = 10;List<Task<bool>> tasks = new List<Task<bool>>();var connection = RabbitMQFactory.SharedConnection;for (int i = 1; i <= asyncCount; i++){tasks.Add(Task.Factory.StartNew(() => MessageWorkItemCallback(connection, rallyNumber)));}Task.WaitAll(tasks.ToArray());string syncResultMsg = $"集结号 {rallyNumber} 已吹起号角--" +$"本次开启信道成功数:{tasks.Count(s => s.Result == true)}," +$"本次开启信道失败数:{tasks.Count() - tasks.Count(s => s.Result == true)}" +$"累计开启信道成功数:{channelCount + tasks.Count(s => s.Result == true)}";ShowMessage(syncResultMsg);}catch (Exception ex){ShowMessage($"集结号 {rallyNumber} 消费异常:{ex.Message}");}});}/// <summary>/// 异步方法/// </summary>/// <param name="state"></param>/// <param name="rallyNumber"></param>/// <returns></returns>private bool MessageWorkItemCallback(object state, int rallyNumber){bool syncResult = false;IModel channel = null;try{IConnection connection = state as IConnection;//不能使用using (channel = connection.CreateModel())来创建信道,让RabbitMQ自动回收channel。channel = connection.CreateModel();channel.QueueDeclare(queue: "hello", durable: true, exclusive: false, autoDelete: false, arguments: null);channel.BasicQos(prefetchSize: 0, prefetchCount: 1, global: false);var consumer = new EventingBasicConsumer(channel);consumer.Received += (model, ea) =>{var message = Encoding.UTF8.GetString(ea.Body);Thread.Sleep(1000);ShowMessage($"集结号 {rallyNumber} Received {message}");channel.BasicAck(deliveryTag: ea.DeliveryTag, multiple: false);};channel.BasicConsume(queue: "hello", autoAck: false, consumer: consumer);syncResult = true;}catch (Exception ex){syncResult = false;ShowMessage(ex.Message);}return syncResult;}
}
2.4、运行结果

在这里插入图片描述

多点几次消费者即可增加信道,提升消费能力。

这篇关于C# 学习笔记:RabbitMQ队列使用多线程提高消费吞吐率的文章就介绍到这儿,希望我们推荐的文章对编程师们有所帮助!



http://www.chinasem.cn/article/1068794

相关文章

Java中流式并行操作parallelStream的原理和使用方法

《Java中流式并行操作parallelStream的原理和使用方法》本文详细介绍了Java中的并行流(parallelStream)的原理、正确使用方法以及在实际业务中的应用案例,并指出在使用并行流... 目录Java中流式并行操作parallelStream0. 问题的产生1. 什么是parallelS

Linux join命令的使用及说明

《Linuxjoin命令的使用及说明》`join`命令用于在Linux中按字段将两个文件进行连接,类似于SQL的JOIN,它需要两个文件按用于匹配的字段排序,并且第一个文件的换行符必须是LF,`jo... 目录一. 基本语法二. 数据准备三. 指定文件的连接key四.-a输出指定文件的所有行五.-o指定输出

Linux jq命令的使用解读

《Linuxjq命令的使用解读》jq是一个强大的命令行工具,用于处理JSON数据,它可以用来查看、过滤、修改、格式化JSON数据,通过使用各种选项和过滤器,可以实现复杂的JSON处理任务... 目录一. 简介二. 选项2.1.2.2-c2.3-r2.4-R三. 字段提取3.1 普通字段3.2 数组字段四.

Linux kill正在执行的后台任务 kill进程组使用详解

《Linuxkill正在执行的后台任务kill进程组使用详解》文章介绍了两个脚本的功能和区别,以及执行这些脚本时遇到的进程管理问题,通过查看进程树、使用`kill`命令和`lsof`命令,分析了子... 目录零. 用到的命令一. 待执行的脚本二. 执行含子进程的脚本,并kill2.1 进程查看2.2 遇到的

详解SpringBoot+Ehcache使用示例

《详解SpringBoot+Ehcache使用示例》本文介绍了SpringBoot中配置Ehcache、自定义get/set方式,并实际使用缓存的过程,文中通过示例代码介绍的非常详细,对大家的学习或者... 目录摘要概念内存与磁盘持久化存储:配置灵活性:编码示例引入依赖:配置ehcache.XML文件:配置

Java 虚拟线程的创建与使用深度解析

《Java虚拟线程的创建与使用深度解析》虚拟线程是Java19中以预览特性形式引入,Java21起正式发布的轻量级线程,本文给大家介绍Java虚拟线程的创建与使用,感兴趣的朋友一起看看吧... 目录一、虚拟线程简介1.1 什么是虚拟线程?1.2 为什么需要虚拟线程?二、虚拟线程与平台线程对比代码对比示例:三

k8s按需创建PV和使用PVC详解

《k8s按需创建PV和使用PVC详解》Kubernetes中,PV和PVC用于管理持久存储,StorageClass实现动态PV分配,PVC声明存储需求并绑定PV,通过kubectl验证状态,注意回收... 目录1.按需创建 PV(使用 StorageClass)创建 StorageClass2.创建 PV

一文解析C#中的StringSplitOptions枚举

《一文解析C#中的StringSplitOptions枚举》StringSplitOptions是C#中的一个枚举类型,用于控制string.Split()方法分割字符串时的行为,核心作用是处理分割后... 目录C#的StringSplitOptions枚举1.StringSplitOptions枚举的常用

Redis 基本数据类型和使用详解

《Redis基本数据类型和使用详解》String是Redis最基本的数据类型,一个键对应一个值,它的功能十分强大,可以存储字符串、整数、浮点数等多种数据格式,本文给大家介绍Redis基本数据类型和... 目录一、Redis 入门介绍二、Redis 的五大基本数据类型2.1 String 类型2.2 Hash

Redis中Hash从使用过程到原理说明

《Redis中Hash从使用过程到原理说明》RedisHash结构用于存储字段-值对,适合对象数据,支持HSET、HGET等命令,采用ziplist或hashtable编码,通过渐进式rehash优化... 目录一、开篇:Hash就像超市的货架二、Hash的基本使用1. 常用命令示例2. Java操作示例三