爆款云主机2核4G限时秒杀,88元/年起!
查看详情

活动

天翼云最新优惠活动,涵盖免费试用,产品折扣等,助您降本增效!
热门活动
  • 618智算钜惠季 爆款云主机2核4G限时秒杀,88元/年起!
  • 免费体验DeepSeek,上天翼云息壤 NEW 新老用户均可免费体验2500万Tokens,限时两周
  • 云上钜惠 HOT 爆款云主机全场特惠,更有万元锦鲤券等你来领!
  • 算力套餐 HOT 让算力触手可及
  • 天翼云脑AOne NEW 连接、保护、办公,All-in-One!
  • 中小企业应用上云专场 产品组合下单即享折上9折起,助力企业快速上云
  • 息壤高校钜惠活动 NEW 天翼云息壤杯高校AI大赛,数款产品享受线上订购超值特惠
  • 天翼云电脑专场 HOT 移动办公新选择,爆款4核8G畅享1年3.5折起,快来抢购!
  • 天翼云奖励推广计划 加入成为云推官,推荐新用户注册下单得现金奖励
免费活动
  • 免费试用中心 HOT 多款云产品免费试用,快来开启云上之旅
  • 天翼云用户体验官 NEW 您的洞察,重塑科技边界

智算服务

打造统一的产品能力,实现算网调度、训练推理、技术架构、资源管理一体化智算服务
智算云(DeepSeek专区)
科研助手
  • 算力商城
  • 应用商城
  • 开发机
  • 并行计算
算力互联调度平台
  • 应用市场
  • 算力市场
  • 算力调度推荐
一站式智算服务平台
  • 模型广场
  • 体验中心
  • 服务接入
智算一体机
  • 智算一体机
大模型
  • DeepSeek-R1-昇腾版(671B)
  • DeepSeek-R1-英伟达版(671B)
  • DeepSeek-V3-昇腾版(671B)
  • DeepSeek-R1-Distill-Llama-70B
  • DeepSeek-R1-Distill-Qwen-32B
  • Qwen2-72B-Instruct
  • StableDiffusion-V2.1
  • TeleChat-12B

应用商城

天翼云精选行业优秀合作伙伴及千余款商品,提供一站式云上应用服务
进入甄选商城进入云市场创新解决方案
办公协同
  • WPS云文档
  • 安全邮箱
  • EMM手机管家
  • 智能商业平台
财务管理
  • 工资条
  • 税务风控云
企业应用
  • 翼信息化运维服务
  • 翼视频云归档解决方案
工业能源
  • 智慧工厂_生产流程管理解决方案
  • 智慧工地
建站工具
  • SSL证书
  • 新域名服务
网络工具
  • 翼云加速
灾备迁移
  • 云管家2.0
  • 翼备份
资源管理
  • 全栈混合云敏捷版(软件)
  • 全栈混合云敏捷版(一体机)
行业应用
  • 翼电子教室
  • 翼智慧显示一体化解决方案

合作伙伴

天翼云携手合作伙伴,共创云上生态,合作共赢
天翼云生态合作中心
  • 天翼云生态合作中心
天翼云渠道合作伙伴
  • 天翼云代理渠道合作伙伴
天翼云服务合作伙伴
  • 天翼云集成商交付能力认证
天翼云应用合作伙伴
  • 天翼云云市场合作伙伴
  • 天翼云甄选商城合作伙伴
天翼云技术合作伙伴
  • 天翼云OpenAPI中心
  • 天翼云EasyCoding平台
天翼云培训认证
  • 天翼云学堂
  • 天翼云市场商学院
天翼云合作计划
  • 云汇计划
天翼云东升计划
  • 适配中心
  • 东升计划
  • 适配互认证

开发者

开发者相关功能入口汇聚
技术社区
  • 专栏文章
  • 互动问答
  • 技术视频
资源与工具
  • OpenAPI中心
开放能力
  • EasyCoding敏捷开发平台
培训与认证
  • 天翼云学堂
  • 天翼云认证
魔乐社区
  • 魔乐社区

支持与服务

为您提供全方位支持与服务,全流程技术保障,助您轻松上云,安全无忧
文档与工具
  • 文档中心
  • 新手上云
  • 自助服务
  • OpenAPI中心
定价
  • 价格计算器
  • 定价策略
基础服务
  • 售前咨询
  • 在线支持
  • 在线支持
  • 工单服务
  • 建议与反馈
  • 用户体验官
  • 服务保障
  • 客户公告
  • 会员中心
增值服务
  • 红心服务
  • 首保服务
  • 客户支持计划
  • 专家技术服务
  • 备案管家

了解天翼云

天翼云秉承央企使命,致力于成为数字经济主力军,投身科技强国伟大事业,为用户提供安全、普惠云服务
品牌介绍
  • 关于天翼云
  • 智算云
  • 天翼云4.0
  • 新闻资讯
  • 天翼云APP
基础设施
  • 全球基础设施
  • 信任中心
最佳实践
  • 精选案例
  • 超级探访
  • 云杂志
  • 分析师和白皮书
  • 天翼云·创新直播间
市场活动
  • 2025智能云生态大会
  • 2024智算云生态大会
  • 2023云生态大会
  • 2022云生态大会
  • 天翼云中国行
天翼云
  • 活动
  • 智算服务
  • 产品
  • 解决方案
  • 应用商城
  • 合作伙伴
  • 开发者
  • 支持与服务
  • 了解天翼云
      • 文档
      • 控制中心
      • 备案
      • 管理中心

      NetMQ Pull-Push 消息模式 + 多线程 + 序列化

      首页 知识中心 软件开发 文章详情页

      NetMQ Pull-Push 消息模式 + 多线程 + 序列化

      2024-07-29 08:02:08 阅读次数:31

      多线程

      近期研究了一下NetMQ,设想把他用在分布式爬虫上面,NetMQ是一个封装了Socket队列的开源库,他是ZeroMQ的.net移植版,而ZeroMQ是用C写成的,有人测试过他的性能,几乎可以秒杀其他所有的MQ(MSMQ,RabitMQ等等,都不是他的对手),不过他也有一个弱点,消息不支持持久化!当然,这个功能可以自己实现,我这里只讲性能,不需要持久化

      下面的例子是我基于NetMQ官网的例子修改的,下面有三个对象Ventilator 消息分发者,Worker 消息处理者,Sink 接受Worker处理消息后返回的结果,耗时的计算处理工作是交给Worker的,如果开多个Worker.exe,可以提升处理速度,Worker的最终目的是分布式计算,部署到多台PC上面,把计算工作交给他们去做(在分布式爬虫上面,每个Worker相当于一个爬虫)。

      不废话,上代码(本来打算用作为序列化格式,在多线程环境下老是报一个错,暂时不知道是什么原因,所以这段注释掉了)

      首先是定义要发送到消息里的对象

      using System;
      using ProtoBuf;
      
      namespace Model
      {
          [Serializable]
          [ProtoContract]
          public class Person
          {
              [ProtoMember(1)]
              public int Id { get; set; }
              [ProtoMember(2)]
              public string Name { get; set; }
              [ProtoMember(3)]
              public DateTime BirthDay { set; get; }
              [ProtoMember(4)]
              public Address Address { get; set; }
          }
      }
      
      using System;
      using ProtoBuf;
      
      namespace Model
      {
          [Serializable]
          [ProtoContract]
          public class Address
          {
              [ProtoMember(1)]
              public string Line1 { get; set; }
              [ProtoMember(2)]
              public string Line2 { get; set; }
          }
      }
      

      然后是消息的发送者

      using System;
      using System.IO;
      using System.Runtime.Remoting.Channels;
      using System.Runtime.Serialization.Formatters.Binary;
      using System.Threading;
      using System.Threading.Tasks;
      using Model;
      using NetMQ;
      using ProtoBuf;
      using ProtoBuf.Meta;
      
      namespace Ventilator
      {
          sealed class Ventilator
          {
              public void Run()
              {
                  Task.Run(() =>
                  {
                      using (var ctx = NetMQContext.Create())
                      using (var sender = ctx.CreatePushSocket())
                      using (var sink = ctx.CreatePushSocket())
                      {
                          sender.Bind("tcp://*:5557");
                          sink.Connect("tcp://localhost:5558");
                          sink.Send("0");
      
                          Console.WriteLine("Sending tasks to workers");
                          RuntimeTypeModel.Default.MetadataTimeoutMilliseconds = 300000;
      
                          //send 100 tasks (workload for tasks, is just some random sleep time that
                          //the workers can perform, in real life each work would do more than sleep
                          for (int taskNumber = 0; taskNumber < 10000; taskNumber++)
                          {
                              Console.WriteLine("Workload : {0}", taskNumber);
                              var person = new Person
                              {
                                  Id = taskNumber,
                                  Name = "First",
                                  BirthDay = DateTime.Parse("1981-11-15"),
                                  Address = new Address { Line1 = "Line1", Line2 = "Line2" }
                              };
                              using (var sm = new MemoryStream())
                              {
                                  //Serializer.PrepareSerializer<Person>();
                                  //Serializer.Serialize(sm, person);
                                  //sender.Send(sm.ToArray());
      
                                  var binaryFormatter = new BinaryFormatter();
                                  binaryFormatter.Serialize(sm, person);
                                  sender.Send(sm.ToArray());
                              }
                          }
                      }
                  });
              }
          }
      }
      
      using System;
      using System.Collections.Generic;
      using System.Linq;
      using System.Text;
      using System.Threading;
      using System.Threading.Tasks;
      using NetMQ;
      
      namespace Ventilator
      {
          public class Program
          {
      
              public static void Main(string[] args)
              {
                  // Task Ventilator
                  // Binds PUSH socket to tcp://localhost:5557
                  // Sends batch of tasks to workers via that socket
                  Console.WriteLine("====== VENTILATOR ======");
      
                  Console.WriteLine("Press enter when worker are ready");
                  Console.ReadLine();
      
                  //the first message it "0" and signals start of batch
                  //see the Sink.csproj Program.cs file for where this is used
                  Console.WriteLine("Sending start of batch to Sink");
      
                  var ventilator = new Ventilator();
                  ventilator.Run();
      
                  Console.WriteLine("Press Enter to quit");
                  Console.ReadLine();
              }
          }
      }
      

      消息的处理者

      using System;
      using System.IO;
      using System.Runtime.Serialization.Formatters.Binary;
      using System.Threading;
      using System.Threading.Tasks;
      using Model;
      using NetMQ;
      using ProtoBuf;
      
      namespace Worker
      {
          sealed class Worker
          {
              public void Run()
              {
                  Task.Run(() =>
                  {
                      using (NetMQContext ctx = NetMQContext.Create())
                      {
                          //socket to receive messages on
                          using (var receiver = ctx.CreatePullSocket())
                          {
                              receiver.Connect("tcp://localhost:5557");
      
                              //socket to send messages on
                              using (var sender = ctx.CreatePushSocket())
                              {
                                  sender.Connect("tcp://localhost:5558");
      
                                  //process tasks forever
                                  while (true)
                                  {
                                      //workload from the vetilator is a simple delay
                                      //to simulate some work being done, see
                                      //Ventilator.csproj Proram.cs for the workload sent
                                      //In real life some more meaningful work would be done
      
                                      //string workload = receiver.ReceiveString();
      
                                      var receivedBytes = receiver.Receive();
                                      using (var sm = new MemoryStream(receivedBytes))
                                      {
                                          // 序列化在多线程方式下报错:
                                          /*
                                            Timeout while inspecting metadata; this may indicate a deadlock. 
                                            This can often be avoided by preparing necessary serializers during application initialization, 
                                            rather than allowing multiple threads to perform the initial metadata inspection; 
                                            please also see the LockContended event
                                           */
                                          //var person = Serializer.Deserialize<Person>(sm);
      
                                          //采用二进制方式
                                          var binaryFormatter = new BinaryFormatter();
                                          var person = binaryFormatter.Deserialize(sm) as Person;
                                          Console.WriteLine("Person {Id:" + person.Id + ",Name:" + person.Name + ",BirthDay:" +
                                                            person.BirthDay + ",Address:{Line1:" + person.Address.Line1 +
                                                            ",Line2:" + person.Address.Line2 + "}}");
                                          Console.WriteLine("Sending to Sink:" + person.Id);
                                          sender.Send(person.Id + "");
                                      }
      
                                      //simulate some work being done
                                      //Thread.Sleep(int.Parse(workload));
                                  }
                              }
                          }
                      }
                  });
              }
          }
      }
      
      using System;
      using System.Collections.Generic;
      using System.Linq;
      using System.Threading.Tasks;
      
      namespace Worker
      {
          public class Program
          {
              public static void Main(string[] args)
              {
                  // Task Worker
                  // Connects PULL socket to tcp://localhost:5557
                  // collects workload for socket from Ventilator via that socket
                  // Connects PUSH socket to tcp://localhost:5558
                  // Sends results to Sink via that socket
                  Console.WriteLine("====== WORKER ======");
      
                  //Task 方式多线程
                  //foreach (Worker client in Enumerable.Range(0, 1000).Select(
                  //    x => new Worker()))
                  //{
                  //    client.Run();
                  //}
      
                  //多核计算方式多线程
                  var actList =
                      Enumerable.Range(0, 50).Select(x => new Worker()).Select(client => (Action)(client.Run)).ToList();
                  var paraOption = new ParallelOptions { MaxDegreeOfParallelism = Environment.ProcessorCount };
                  Parallel.Invoke(paraOption, actList.ToArray());
      
                  Console.ReadLine();
              }
          }
      }
      

      接受消息处理的结果

      using System;
      using System.Collections.Generic;
      using System.Diagnostics;
      using System.Linq;
      using System.Text;
      using System.Threading.Tasks;
      using NetMQ;
      
      namespace Sink
      {
          public class Program
          {
              public static void Main(string[] args)
              {
                  // Task Sink
                  // Bindd PULL socket to tcp://localhost:5558
                  // Collects results from workers via that socket
                  Console.WriteLine("====== SINK ======");
      
                  using (NetMQContext ctx = NetMQContext.Create())
                  {
                      //socket to receive messages on
                      using (var receiver = ctx.CreatePullSocket())
                      {
                          receiver.Bind("tcp://localhost:5558");
      
                          //wait for start of batch (see Ventilator.csproj Program.cs)
                          var startOfBatchTrigger = receiver.ReceiveString();
                          Console.WriteLine("Seen start of batch");
      
                          //Start our clock now
                          Stopwatch watch = new Stopwatch();
                          watch.Start();
      
                          for (int taskNumber = 0; taskNumber < 10000; taskNumber++)
                          {
                          //while (true)
                          //{
                              var workerDoneTrigger = receiver.ReceiveString();
                              Console.WriteLine(workerDoneTrigger);
                          //}
                          }
                          watch.Stop();
                          //Calculate and report duration of batch
                          Console.WriteLine();
                          Console.WriteLine("Total elapsed time {0} msec", watch.ElapsedMilliseconds);
                          Console.ReadLine();
                      }
                  }
              }
          }
      }
      再次提醒,Worker.exe 可以开多个,以提高效率
      版权声明:本文内容来自第三方投稿或授权转载,原文地址:https://blog.csdn.net/lee576/article/details/47067233,作者:lee576,版权归原作者所有。本网站转在其作品的目的在于传递更多信息,不拥有版权,亦不承担相应法律责任。如因作品内容、版权等问题需要同本网站联系,请发邮件至ctyunbbs@chinatelecom.cn沟通。

      上一篇:R语言中的block Gibbs吉布斯采样贝叶斯多元线性回归|附代码数据

      下一篇:Ultrawebgrid中利用JS得到选中行的值

      相关文章

      2025-05-19 09:04:44

      spark控制台没显示其他机器

      spark控制台没显示其他机器

      2025-05-19 09:04:44
      Spark , 节点 , 集群
      2025-05-06 09:19:12

      Spring多线程事务 能否保证事务的一致性(同时提交、同时回滚)?

      Spring的事务信息是存在ThreadLocal中的Connection, 所以一个线程永远只能有一个事务

      2025-05-06 09:19:12
      Spring , 事务 , 多线程 , 线程
      2025-04-11 07:08:33

      Java线程中的run()和start()区别

      Java线程中的run()和start()区别

      2025-04-11 07:08:33
      run , start , 启动 , 多线程 , 方法 , 线程 , 运行
      2025-04-09 09:16:00

      Java线程的基础概念介绍(结合代码说明)

      线程是操作系统能够进行运算调度的最小单位,它是进程中的实际运作单位。每个线程执行的都是某一个进程的代码的某个片段。

      2025-04-09 09:16:00
      CPU , 多线程 , 方法 , 状态 , 线程 , 进程 , 阻塞
      2025-03-27 09:34:39

      阻塞与唤醒:多线程编程的神秘面纱

      在多线程编程中,线程状态切换是一个非常关键的概念。了解线程状态切换的原理,对于编写高效、稳定的多线程程序至关重要。

      2025-03-27 09:34:39
      切换 , 多线程 , 状态 , 等待 , 线程
      2025-03-26 08:57:33

      三种方法教你实现多线程交替打印ABC,干货满满!

      假设有三个线程,分别打印字母A、B、C。我们需要让这三个线程交替运行,按顺序打印出“ABCABCABC...”,直到打印一定次数或者满足某个条件。如何通过多线程的协调实现这个任务呢?这听起来简单,实际涉及到线程之间的同步和互斥,是我们学习多线程编程的一个很好的练习。

      2025-03-26 08:57:33
      Condition , wait , 信号量 , 多线程 , 线程
      2025-03-21 09:33:29

      C运行时库(C Run-Time Libraries)

      C运行时库(C Run-Time Libraries)

      2025-03-21 09:33:29
      DLL , lib , link , 多线程 , 链接 , 静态
      2025-03-21 08:23:07

      深入理解Java中的多线程编程

      在本文中,我们将探讨Java多线程编程的核心概念和实践。我们将从基本概念开始,逐步深入到线程的创建与管理、同步与锁机制以及高级并发工具的应用。通过实例代码和详细解释,帮助读者全面掌握Java多线程编程的精髓。

      2025-03-21 08:23:07
      Java , Runnable , 同步 , 多线程 , 并发 , 线程 , 编程
      2025-03-18 09:59:32

      深入学习Java语言核心技术

      深入学习Java语言核心技术

      2025-03-18 09:59:32
      Java , JVM , 多线程 , 并发 , 框架 , 线程 , 集合
      2025-03-11 09:35:24

      【高并发】java高并发核心知识

      【高并发】java高并发核心知识

      2025-03-11 09:35:24
      java , 多线程 , 并发 , 线程 , 编程
      查看更多
      推荐标签

      作者介绍

      天翼云小翼
      天翼云用户

      文章

      33561

      阅读量

      5255355

      查看更多

      最新文章

      Spring多线程事务 能否保证事务的一致性(同时提交、同时回滚)?

      2025-05-06 09:19:12

      Java线程中的run()和start()区别

      2025-04-11 07:08:33

      Java线程的基础概念介绍(结合代码说明)

      2025-04-09 09:16:00

      阻塞与唤醒:多线程编程的神秘面纱

      2025-03-27 09:34:39

      三种方法教你实现多线程交替打印ABC,干货满满!

      2025-03-26 08:57:33

      C运行时库(C Run-Time Libraries)

      2025-03-21 09:33:29

      查看更多

      热门文章

      JAVA多线程学习笔记

      2023-05-11 06:05:48

      Thrift第七课 服务器多线程发送异常

      2023-05-16 09:42:24

      synchronized实现两个线程交替运行

      2022-12-28 07:22:30

      线程池笔记(一)

      2022-12-28 07:22:30

      【多线程】synchronized 中的 锁优化的机制 (偏向锁->轻量级锁->重量级锁)

      2023-04-13 09:26:52

      kotlin创建简单多线程的3种方式

      2023-04-17 09:39:23

      查看更多

      热门标签

      java Java python 编程开发 代码 开发语言 算法 线程 Python html 数组 C++ 元素 javascript c++
      查看更多

      相关产品

      弹性云主机

      随时自助获取、弹性伸缩的云服务器资源

      天翼云电脑(公众版)

      便捷、安全、高效的云电脑服务

      对象存储

      高品质、低成本的云上存储服务

      云硬盘

      为云上计算资源提供持久性块存储

      查看更多

      随机文章

      三种方法教你实现多线程交替打印ABC,干货满满!

      JVM对synchronized做了哪些优化

      自己开发的在线视频下载工具,基于 Java 多线程

      使用Java实现高并发的数据同步

      Linux生产者消费者模型

      Python爬虫-第四章-3-多线程多进程提升任务执行效率

      • 7*24小时售后
      • 无忧退款
      • 免费备案
      • 专家服务
      售前咨询热线
      400-810-9889转1
      关注天翼云
      • 旗舰店
      • 天翼云APP
      • 天翼云微信公众号
      服务与支持
      • 备案中心
      • 售前咨询
      • 智能客服
      • 自助服务
      • 工单管理
      • 客户公告
      • 涉诈举报
      账户管理
      • 管理中心
      • 订单管理
      • 余额管理
      • 发票管理
      • 充值汇款
      • 续费管理
      快速入口
      • 天翼云旗舰店
      • 文档中心
      • 最新活动
      • 免费试用
      • 信任中心
      • 天翼云学堂
      云网生态
      • 甄选商城
      • 渠道合作
      • 云市场合作
      了解天翼云
      • 关于天翼云
      • 天翼云APP
      • 服务案例
      • 新闻资讯
      • 联系我们
      热门产品
      • 云电脑
      • 弹性云主机
      • 云电脑政企版
      • 天翼云手机
      • 云数据库
      • 对象存储
      • 云硬盘
      • Web应用防火墙
      • 服务器安全卫士
      • CDN加速
      热门推荐
      • 云服务备份
      • 边缘安全加速平台
      • 全站加速
      • 安全加速
      • 云服务器
      • 云主机
      • 智能边缘云
      • 应用编排服务
      • 微服务引擎
      • 共享流量包
      更多推荐
      • web应用防火墙
      • 密钥管理
      • 等保咨询
      • 安全专区
      • 应用运维管理
      • 云日志服务
      • 文档数据库服务
      • 云搜索服务
      • 数据湖探索
      • 数据仓库服务
      友情链接
      • 中国电信集团
      • 189邮箱
      • 天翼企业云盘
      • 天翼云盘
      ©2025 天翼云科技有限公司版权所有 增值电信业务经营许可证A2.B1.B2-20090001
      公司地址:北京市东城区青龙胡同甲1号、3号2幢2层205-32室
      • 用户协议
      • 隐私政策
      • 个人信息保护
      • 法律声明
      备案 京公网安备11010802043424号 京ICP备 2021034386号