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

相关文章

Python常用命令提示符使用方法详解

《Python常用命令提示符使用方法详解》在学习python的过程中,我们需要用到命令提示符(CMD)进行环境的配置,:本文主要介绍Python常用命令提示符使用方法的相关资料,文中通过代码介绍的... 目录一、python环境基础命令【Windows】1、检查Python是否安装2、 查看Python的安

Python并行处理实战之如何使用ProcessPoolExecutor加速计算

《Python并行处理实战之如何使用ProcessPoolExecutor加速计算》Python提供了多种并行处理的方式,其中concurrent.futures模块的ProcessPoolExecu... 目录简介完整代码示例代码解释1. 导入必要的模块2. 定义处理函数3. 主函数4. 生成数字列表5.

Python中help()和dir()函数的使用

《Python中help()和dir()函数的使用》我们经常需要查看某个对象(如模块、类、函数等)的属性和方法,Python提供了两个内置函数help()和dir(),它们可以帮助我们快速了解代... 目录1. 引言2. help() 函数2.1 作用2.2 使用方法2.3 示例(1) 查看内置函数的帮助(

Linux脚本(shell)的使用方式

《Linux脚本(shell)的使用方式》:本文主要介绍Linux脚本(shell)的使用方式,具有很好的参考价值,希望对大家有所帮助,如有错误或未考虑完全的地方,望不吝赐教... 目录概述语法详解数学运算表达式Shell变量变量分类环境变量Shell内部变量自定义变量:定义、赋值自定义变量:引用、修改、删

Java使用HttpClient实现图片下载与本地保存功能

《Java使用HttpClient实现图片下载与本地保存功能》在当今数字化时代,网络资源的获取与处理已成为软件开发中的常见需求,其中,图片作为网络上最常见的资源之一,其下载与保存功能在许多应用场景中都... 目录引言一、Apache HttpClient简介二、技术栈与环境准备三、实现图片下载与保存功能1.

Python中使用uv创建环境及原理举例详解

《Python中使用uv创建环境及原理举例详解》uv是Astral团队开发的高性能Python工具,整合包管理、虚拟环境、Python版本控制等功能,:本文主要介绍Python中使用uv创建环境及... 目录一、uv工具简介核心特点:二、安装uv1. 通过pip安装2. 通过脚本安装验证安装:配置镜像源(可

LiteFlow轻量级工作流引擎使用示例详解

《LiteFlow轻量级工作流引擎使用示例详解》:本文主要介绍LiteFlow是一个灵活、简洁且轻量的工作流引擎,适合用于中小型项目和微服务架构中的流程编排,本文给大家介绍LiteFlow轻量级工... 目录1. LiteFlow 主要特点2. 工作流定义方式3. LiteFlow 流程示例4. LiteF

使用Python开发一个现代化屏幕取色器

《使用Python开发一个现代化屏幕取色器》在UI设计、网页开发等场景中,颜色拾取是高频需求,:本文主要介绍如何使用Python开发一个现代化屏幕取色器,有需要的小伙伴可以参考一下... 目录一、项目概述二、核心功能解析2.1 实时颜色追踪2.2 智能颜色显示三、效果展示四、实现步骤详解4.1 环境配置4.

使用jenv工具管理多个JDK版本的方法步骤

《使用jenv工具管理多个JDK版本的方法步骤》jenv是一个开源的Java环境管理工具,旨在帮助开发者在同一台机器上轻松管理和切换多个Java版本,:本文主要介绍使用jenv工具管理多个JD... 目录一、jenv到底是干啥的?二、jenv的核心功能(一)管理多个Java版本(二)支持插件扩展(三)环境隔

SQL中JOIN操作的条件使用总结与实践

《SQL中JOIN操作的条件使用总结与实践》在SQL查询中,JOIN操作是多表关联的核心工具,本文将从原理,场景和最佳实践三个方面总结JOIN条件的使用规则,希望可以帮助开发者精准控制查询逻辑... 目录一、ON与WHERE的本质区别二、场景化条件使用规则三、最佳实践建议1.优先使用ON条件2.WHERE用