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/1068793

相关文章

HarmonyOS学习(七)——UI(五)常用布局总结

自适应布局 1.1、线性布局(LinearLayout) 通过线性容器Row和Column实现线性布局。Column容器内的子组件按照垂直方向排列,Row组件中的子组件按照水平方向排列。 属性说明space通过space参数设置主轴上子组件的间距,达到各子组件在排列上的等间距效果alignItems设置子组件在交叉轴上的对齐方式,且在各类尺寸屏幕上表现一致,其中交叉轴为垂直时,取值为Vert

Ilya-AI分享的他在OpenAI学习到的15个提示工程技巧

Ilya(不是本人,claude AI)在社交媒体上分享了他在OpenAI学习到的15个Prompt撰写技巧。 以下是详细的内容: 提示精确化:在编写提示时,力求表达清晰准确。清楚地阐述任务需求和概念定义至关重要。例:不用"分析文本",而用"判断这段话的情感倾向:积极、消极还是中性"。 快速迭代:善于快速连续调整提示。熟练的提示工程师能够灵活地进行多轮优化。例:从"总结文章"到"用

中文分词jieba库的使用与实景应用(一)

知识星球:https://articles.zsxq.com/id_fxvgc803qmr2.html 目录 一.定义: 精确模式(默认模式): 全模式: 搜索引擎模式: paddle 模式(基于深度学习的分词模式): 二 自定义词典 三.文本解析   调整词出现的频率 四. 关键词提取 A. 基于TF-IDF算法的关键词提取 B. 基于TextRank算法的关键词提取

使用SecondaryNameNode恢复NameNode的数据

1)需求: NameNode进程挂了并且存储的数据也丢失了,如何恢复NameNode 此种方式恢复的数据可能存在小部分数据的丢失。 2)故障模拟 (1)kill -9 NameNode进程 [lytfly@hadoop102 current]$ kill -9 19886 (2)删除NameNode存储的数据(/opt/module/hadoop-3.1.4/data/tmp/dfs/na

Hadoop数据压缩使用介绍

一、压缩原则 (1)运算密集型的Job,少用压缩 (2)IO密集型的Job,多用压缩 二、压缩算法比较 三、压缩位置选择 四、压缩参数配置 1)为了支持多种压缩/解压缩算法,Hadoop引入了编码/解码器 2)要在Hadoop中启用压缩,可以配置如下参数

【前端学习】AntV G6-08 深入图形与图形分组、自定义节点、节点动画(下)

【课程链接】 AntV G6:深入图形与图形分组、自定义节点、节点动画(下)_哔哩哔哩_bilibili 本章十吾老师讲解了一个复杂的自定义节点中,应该怎样去计算和绘制图形,如何给一个图形制作不间断的动画,以及在鼠标事件之后产生动画。(有点难,需要好好理解) <!DOCTYPE html><html><head><meta charset="UTF-8"><title>06

Makefile简明使用教程

文章目录 规则makefile文件的基本语法:加在命令前的特殊符号:.PHONY伪目标: Makefilev1 直观写法v2 加上中间过程v3 伪目标v4 变量 make 选项-f-n-C Make 是一种流行的构建工具,常用于将源代码转换成可执行文件或者其他形式的输出文件(如库文件、文档等)。Make 可以自动化地执行编译、链接等一系列操作。 规则 makefile文件

学习hash总结

2014/1/29/   最近刚开始学hash,名字很陌生,但是hash的思想却很熟悉,以前早就做过此类的题,但是不知道这就是hash思想而已,说白了hash就是一个映射,往往灵活利用数组的下标来实现算法,hash的作用:1、判重;2、统计次数;

hdu1180(广搜+优先队列)

此题要求最少到达目标点T的最短时间,所以我选择了广度优先搜索,并且要用到优先队列。 另外此题注意点较多,比如说可以在某个点停留,我wa了好多两次,就是因为忽略了这一点,然后参考了大神的思想,然后经过反复修改才AC的 这是我的代码 #include<iostream>#include<algorithm>#include<string>#include<stack>#include<

使用opencv优化图片(画面变清晰)

文章目录 需求影响照片清晰度的因素 实现降噪测试代码 锐化空间锐化Unsharp Masking频率域锐化对比测试 对比度增强常用算法对比测试 需求 对图像进行优化,使其看起来更清晰,同时保持尺寸不变,通常涉及到图像处理技术如锐化、降噪、对比度增强等 影响照片清晰度的因素 影响照片清晰度的因素有很多,主要可以从以下几个方面来分析 1. 拍摄设备 相机传感器:相机传