reactor模式:多线程的reactor模式

上文说到单线程的reactor模式 reactor模式:单线程的reactor模式

单线程的reactor模式并没有解决IO和CPU处理速度不匹配问题,所以多线程的reactor模式引入线程池的概念,把耗时的IO操作交由线程池处理,处理完了之后再同步到selectionkey中,服务器架构图如下

 

 

上文(reactor模式:单线程的reactor模式)提到,以read和send阶段IO最为频繁,所以多线程的reactor版本里,把这2个阶段单独拎出来。

下面看看代码实现

 

 1 // Reactor線程 (该类与单线程的处理基本无变动) 
 2     package server;  
 3       
 4     import java.io.IOException;  
 5     import java.net.InetSocketAddress;  
 6     import java.nio.channels.SelectionKey;  
 7     import java.nio.channels.Selector;  
 8     import java.nio.channels.ServerSocketChannel;  
 9     import java.util.Iterator;  
10     import java.util.Set;  
11       
12     public class TCPReactor implements Runnable {  
13       
14         private final ServerSocketChannel ssc;  
15         private final Selector selector;  
16       
17         public TCPReactor(int port) throws IOException {  
18             selector = Selector.open();  
19             ssc = ServerSocketChannel.open();  
20             InetSocketAddress addr = new InetSocketAddress(port);  
21             ssc.socket().bind(addr); // 在ServerSocketChannel綁定監聽端口  
22             ssc.configureBlocking(false); // 設置ServerSocketChannel為非阻塞  
23             SelectionKey sk = ssc.register(selector, SelectionKey.OP_ACCEPT); // ServerSocketChannel向selector註冊一個OP_ACCEPT事件,然後返回該通道的key  
24             sk.attach(new Acceptor(selector, ssc)); // 給定key一個附加的Acceptor對象  
25         }  
26       
27         @Override  
28         public void run() {  
29             while (!Thread.interrupted()) { // 在線程被中斷前持續運行  
30                 System.out.println("Waiting for new event on port: " + ssc.socket().getLocalPort() + "...");  
31                 try {  
32                     if (selector.select() == 0) // 若沒有事件就緒則不往下執行  
33                         continue;  
34                 } catch (IOException e) {  
35                     // TODO Auto-generated catch block  
36                     e.printStackTrace();  
37                 }  
38                 Set<SelectionKey> selectedKeys = selector.selectedKeys(); // 取得所有已就緒事件的key集合  
39                 Iterator<SelectionKey> it = selectedKeys.iterator();  
40                 while (it.hasNext()) {  
41                     dispatch((SelectionKey) (it.next())); // 根據事件的key進行調度  
42                     it.remove();  
43                 }  
44             }  
45         }  
46       
47         /* 
48          * name: dispatch(SelectionKey key) 
49          * description: 調度方法,根據事件綁定的對象開新線程 
50          */  
51         private void dispatch(SelectionKey key) {  
52             Runnable r = (Runnable) (key.attachment()); // 根據事件之key綁定的對象開新線程  
53             if (r != null)  
54                 r.run();  
55         }  
56       
57     }  

 

 

 1  // 接受連線請求線程  
 2     package server;  
 3       
 4     import java.io.IOException;  
 5     import java.nio.channels.SelectionKey;  
 6     import java.nio.channels.Selector;  
 7     import java.nio.channels.ServerSocketChannel;  
 8     import java.nio.channels.SocketChannel;  
 9       
10     public class Acceptor implements Runnable {  
11       
12         private final ServerSocketChannel ssc;  
13         private final Selector selector;  
14           
15         public Acceptor(Selector selector, ServerSocketChannel ssc) {  
16             this.ssc=ssc;  
17             this.selector=selector;  
18         }  
19           
20         @Override  
21         public void run() {  
22             try {  
23                 SocketChannel sc= ssc.accept(); // 接受client連線請求  
24                 System.out.println(sc.socket().getRemoteSocketAddress().toString() + " is connected.");  
25                   
26                 if(sc!=null) {  
27                     sc.configureBlocking(false); // 設置為非阻塞  
28                     SelectionKey sk = sc.register(selector, SelectionKey.OP_READ); // SocketChannel向selector註冊一個OP_READ事件,然後返回該通道的key  
29                     selector.wakeup(); // 使一個阻塞住的selector操作立即返回  
30                     sk.attach(new TCPHandler(sk, sc)); // 給定key一個附加的TCPHandler對象  
31                 }  
32                   
33             } catch (IOException e) {  
34                 // TODO Auto-generated catch block  
35                 e.printStackTrace();  
36             }  
37         }  
38       
39           
40     }  

 

 

 1     // Handler線程  
 2     package server;  
 3       
 4     import java.io.IOException;  
 5     import java.nio.channels.SelectionKey;  
 6     import java.nio.channels.SocketChannel;  
 7     import java.util.concurrent.LinkedBlockingQueue;  
 8     import java.util.concurrent.ThreadPoolExecutor;  
 9     import java.util.concurrent.TimeUnit;  
10       
11     public class TCPHandler implements Runnable {  
12       
13         private final SelectionKey sk;  
14         private final SocketChannel sc;  
15         private static final int THREAD_COUNTING = 10;  
16         private static ThreadPoolExecutor pool = new ThreadPoolExecutor(  
17                 THREAD_COUNTING, THREAD_COUNTING, 10, TimeUnit.SECONDS,  
18                 new LinkedBlockingQueue<Runnable>()); // 線程池  
19       
20         HandlerState state; // 以狀態模式實現Handler  
21       
22         public TCPHandler(SelectionKey sk, SocketChannel sc) {  
23             this.sk = sk;  
24             this.sc = sc;  
25             state = new ReadState(); // 初始狀態設定為READING  
26             pool.setMaximumPoolSize(32); // 設置線程池最大線程數  
27         }  
28       
29         @Override  
30         public void run() {  
31             try {  
32                 state.handle(this, sk, sc, pool);  
33                   
34             } catch (IOException e) {  
35                 System.out.println("[Warning!] A client has been closed.");  
36                 closeChannel();  
37             }  
38         }  
39           
40         public void closeChannel() {  
41             try {  
42                 sk.cancel();  
43                 sc.close();  
44             } catch (IOException e1) {  
45                 e1.printStackTrace();  
46             }  
47         }  
48       
49         public void setState(HandlerState state) {  
50             this.state = state;  
51         }  
52     }  
53 
54  

 

 1     package server;  
 2       
 3     import java.io.IOException;  
 4     import java.nio.channels.SelectionKey;  
 5     import java.nio.channels.SocketChannel;  
 6     import java.util.concurrent.ThreadPoolExecutor;  
 7       
 8     public interface HandlerState {  
 9       
10         public void changeState(TCPHandler h);  
11       
12         public void handle(TCPHandler h, SelectionKey sk, SocketChannel sc,  
13                 ThreadPoolExecutor pool) throws IOException ;  
14     }  

 

 

 1     package server;  
 2       
 3     import java.io.IOException;  
 4     import java.nio.ByteBuffer;  
 5     import java.nio.channels.SelectionKey;  
 6     import java.nio.channels.SocketChannel;  
 7     import java.util.concurrent.ThreadPoolExecutor;  
 8       
 9     public class ReadState implements HandlerState{  
10       
11         private SelectionKey sk;  
12           
13         public ReadState() {  
14         }  
15           
16         @Override  
17         public void changeState(TCPHandler h) {  
18             // TODO Auto-generated method stub  
19             h.setState(new WorkState());  
20         }  
21       
22         @Override  
23         public void handle(TCPHandler h, SelectionKey sk, SocketChannel sc,  
24                 ThreadPoolExecutor pool) throws IOException { // read()  
25             this.sk = sk;  
26             // non-blocking下不可用Readers,因為Readers不支援non-blocking  
27             byte[] arr = new byte[1024];  
28             ByteBuffer buf = ByteBuffer.wrap(arr);  
29               
30             int numBytes = sc.read(buf); // 讀取字符串  
31             if(numBytes == -1)  
32             {  
33                 System.out.println("[Warning!] A client has been closed.");  
34                 h.closeChannel();  
35                 return;  
36             }  
37             String str = new String(arr); // 將讀取到的byte內容轉為字符串型態  
38             if ((str != null) && !str.equals(" ")) {  
39                 h.setState(new WorkState()); // 改變狀態(READING->WORKING)  
40                 pool.execute(new WorkerThread(h, str)); // do process in worker thread  
41                 System.out.println(sc.socket().getRemoteSocketAddress().toString()  
42                         + " > " + str);  
43             }  
44               
45         }  
46           
47         /* 
48          * 執行邏輯處理之函數 
49          */  
50         synchronized void process(TCPHandler h, String str) {  
51             // do process(decode, logically process, encode)..  
52             // ..  
53             h.setState(new WriteState()); // 改變狀態(WORKING->SENDING)  
54             this.sk.interestOps(SelectionKey.OP_WRITE); // 通過key改變通道註冊的事件  
55             this.sk.selector().wakeup(); // 使一個阻塞住的selector操作立即返回  
56         }  
57       
58         /* 
59          * 工作者線程 
60          */  
61         class WorkerThread implements Runnable {  
62       
63             TCPHandler h;  
64             String str;  
65       
66             public WorkerThread(TCPHandler h, String str) {  
67                 this.h = h;  
68                 this.str=str;  
69             }  
70       
71             @Override  
72             public void run() {  
73                 process(h, str);  
74             }  
75       
76         }  
77     }  

 

 1  package server;  
 2       
 3     import java.io.IOException;  
 4     import java.nio.channels.SelectionKey;  
 5     import java.nio.channels.SocketChannel;  
 6     import java.util.concurrent.ThreadPoolExecutor;  
 7       
 8     public class WorkState implements HandlerState {  
 9       
10         public WorkState() {  
11         }  
12           
13         @Override  
14         public void changeState(TCPHandler h) {  
15             // TODO Auto-generated method stub  
16             h.setState(new WriteState());  
17         }  
18       
19         @Override  
20         public void handle(TCPHandler h, SelectionKey sk, SocketChannel sc,  
21                 ThreadPoolExecutor pool) throws IOException {  
22             // TODO Auto-generated method stub  
23               
24         }  
25       
26     }  
 1     package server;  
 2       
 3     import java.io.IOException;  
 4     import java.nio.ByteBuffer;  
 5     import java.nio.channels.SelectionKey;  
 6     import java.nio.channels.SocketChannel;  
 7     import java.util.concurrent.ThreadPoolExecutor;  
 8       
 9     public class WriteState implements HandlerState{  
10       
11         public WriteState() {  
12         }  
13           
14         @Override  
15         public void changeState(TCPHandler h) {  
16             // TODO Auto-generated method stub  
17             h.setState(new ReadState());  
18         }  
19       
20         @Override  
21         public void handle(TCPHandler h, SelectionKey sk, SocketChannel sc,  
22                 ThreadPoolExecutor pool) throws IOException { // send()  
23             // get message from message queue  
24               
25             String str = "Your message has sent to "  
26                     + sc.socket().getLocalSocketAddress().toString() + "\r\n";  
27             ByteBuffer buf = ByteBuffer.wrap(str.getBytes()); // wrap自動把buf的position設為0,所以不需要再flip()  
28       
29             while (buf.hasRemaining()) {  
30                 sc.write(buf); // 回傳給client回應字符串,發送buf的position位置 到limit位置為止之間的內容  
31             }  
32               
33             h.setState(new ReadState()); // 改變狀態(SENDING->READING)  
34             sk.interestOps(SelectionKey.OP_READ); // 通過key改變通道註冊的事件  
35             sk.selector().wakeup(); // 使一個阻塞住的selector操作立即返回  
36         }  
37     }  
 1     package server;  
 2       
 3     import java.io.IOException;  
 4       
 5     public class Main {  
 6       
 7           
 8         public static void main(String[] args) {  
 9             // TODO Auto-generated method stub  
10             try {  
11                 TCPReactor reactor = new TCPReactor(1333);  
12                 reactor.run();  
13             } catch (IOException e) {  
14                 // TODO Auto-generated catch block  
15                 e.printStackTrace();  
16             }  
17         }  
18       
19     }  

 

总的来说,多线程版本的reactor是为了解决单线程reactor版本的IO和CPU处理速度不匹配问题,从而达到高效处理的目的

 

参考文章:

https://blog.csdn.net/yehjordan/article/details/51017025

转载于:https://www.cnblogs.com/billmiao/p/9872221.html

本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若转载,请注明出处:http://www.mzph.cn/news/252080.shtml

如若内容造成侵权/违法违规/事实不符,请联系多彩编程网进行投诉反馈email:809451989@qq.com,一经查实,立即删除!

相关文章

Elasticsearch实战篇——Spring Boot整合ElasticSearch

2019独角兽企业重金招聘Python工程师标准>>> 当前Spring Boot很是流行&#xff0c;包括我自己&#xff0c;也是在用Spring Boot集成其他框架进行项目开发&#xff0c;所以这一节&#xff0c;我们一起来探讨Spring Boot整合ElasticSearch的问题。 本文主要讲以下内容…

Python: pip升级报错了:You are using pip version 10.0.1, however version 20.3.3 is available.

1,Python使用命令&#xff1a;python -m pip install --upgrade pip升级pip的时候报了下面这个错 2,换了个命令&#xff1a; python -m pip install --upgrade pip -i https://pypi.douban.com/simple 更新成功了&#xff0c;但又报了一个新的错误&#xff1a; AttributeError:…

新手上路之Hibernate:第一个Hibernate例子

一、Hibernate概述 &#xff08;一&#xff09;什么是Hibernate&#xff1f; Hibernate核心内容是ORM&#xff08;关系对象模型&#xff09;。可以将对象自动的生成数据库中的信息&#xff0c;使得开发更加的面向对象。这样作为程序员就可以使用面向对象的思想来操作数据库&…

模板标签及模板的继承与引用

1.常用的模板标签 - 作用是什么:提供各种逻辑 view.py: def index(request):#模板标签 --常用标签 总结&#xff1a;语法 {% tag %} {% endtag %} {% tag 参数 参数 %} 示例 展示页index.html&#xff0c;包含for标签&#xff0c;if标签&#xff0c;url标签 {% extends teacher…

Golang实现一个密码生成器

小地鼠防止有人偷他的果实&#xff0c;在家里上了一把锁。这个锁怎么来的呢&#xff1f;请往下看。。 package mainimport ("flag""fmt""math/rand""time" )var (length intcharset string )const (NUmStr "0123456789"C…

C# WPF:初识布局容器

StackPanel堆叠布局 StackPanel是简单布局方式之一&#xff0c;可以很方便的进行纵向布局和横向布局 StackPanel默认是纵向布局的 <Window x:Class"WpfApplication1.MainWindow" xmlns"http://schemas.microsoft.com/winfx/2006/xaml/presentation" …

Kibana源码分析--Hapijs路由设置理解笔记

【ES6解构赋值】&#xff1a;https://developer.mozilla.org/zh-CN/docs/Web/JavaScript/Reference/Operators/Destructuring_assignment 【Joi APi】&#xff1a;https://github.com/hapijs/joi/blob/v13.1.2/API.md 转载于:https://www.cnblogs.com/lishidefengchen/p/866874…

Python打包EXE神器 pyinstaller

最近由于项目需要&#xff0c;以前的python文件需要编辑为EXE供前端客户使用。 由于最早接触的是distutils&#xff0c;所以一开始准备使用distutils和py2exe搭配来进行python的exe化&#xff0c;也就是传统的使用setup.py的方式来进行exe安装。但是结果都不是很好&#xff0c;…

20种PLC元件编号和Modbus编号地址对应表

1、三菱&#xff1a; X元件支持Modbus之02功能码&#xff1b; Y元件支持Modbus之01、05、15功能码&#xff1b; D元件支持Modbus之03、06、16功能码。 2、西门子&#xff1a; I元件支持Modbus之02功能码&#xff1b; Q元件支持Modbus之01、05、15功能码&#xff1b; V元件…

暑期学习

由于最后大作业的呈现情况与短学期所完成的还相差甚远&#xff0c;所以在暑期的时候开始进一步的细化。 在这个过程之中产生了如下的问题&#xff1a; 已解决的有&#xff1a; 1.用a标签在同一页面实现跳转。 要点&#xff1a;标记<a href"../home#pre">的时候…

五、RabbitMQ的消息属性(读书笔记)

2019独角兽企业重金招聘Python工程师标准>>> 简介 当使用RabbitMQ发布消息时&#xff0c;消息又AMQP规范中的三个低层帧类型组成&#xff1a; Basic.publish方法帧&#xff1b;内容头帧&#xff1b;消息体帧&#xff1b;这三种帧类型按顺序一起工作&#xff0c;以便…

异步和单线程

转载于:https://www.cnblogs.com/sunmarvell/p/8674748.html

C#:把dll封入exe中方法

在这个事件中,可以重新为加载失败的程序集手动加载 如果你将dll作为资源文件打包的你的应用程序中(或者类库中) 就可以在硬盘加载失败的时候 从资源文件中加载对应的dll 就像这样: class Program {static Program(){ //这个绑定事件必须要在引用到TestLibrary1这个程序…

C#结构类型图

转载于:https://www.cnblogs.com/kangao/p/8674838.html

使用gradle多渠道打包

以友盟的多渠道打包为例&#xff0c;如果我们须要打包出例如以下渠道&#xff1a;UMENG, WANDOUJIA, YINGYONGBAO。 第一种方法。是须要创建文件的。我们在写完我们的代码之后&#xff0c;在app/src以下。分别创建和main同级目录的目录umeng, wandoujia, yingyongbao,这三个目录…

四大步骤,彻底关闭Win10自动更新

尽管Win11已经发布了一段时间&#xff0c;但目前互联网上大部分电脑用户所使用的的操作系统仍是Win10&#xff0c;对于Win10&#xff0c;笔者相信大部分人应该都不陌生&#xff0c;作为目前市面上占比最高的电脑系统&#xff0c;Win10的许多功能和操作逻辑都十分优秀&#xff0…

虚拟机windows7安装启动MYSQL5.7

一.环境 环境&#xff1a;虚拟机VMVare 系统&#xff1a;windows7旗舰版 MYSQL版本&#xff1a;mysql5.7.25 二.具体步骤 1.首先下载安装mysql5.7.25&#xff0c;这里用的是安装版的mysql&#xff0c;网上大多数都是推荐去官网下载&#xff0c;这里推荐的是清华大学开源镜像站…

故障转移架构的本质:数据中心的基础设施过剩

数据中心构成了全球互联基础设施的核心&#xff0c;我们称之为“云”。从根本上讲&#xff0c;云计算指的是基础设施从桌面计算&#xff08;文件和应用程序存储在计算机的本地硬盘上&#xff09;到在线计算&#xff08;文件和应用程序存储在可通过互联网远程访问的数据中心中&a…

常用模块之hashlib,configparser,logging模块

常用模块二 hashlib模块 hashlib提供了常见的摘要算法&#xff0c;如md5和sha1等等。 那么什么是摘要算法呢?摘要算法又称为哈希算法、散列算法。它通过一个函数&#xff0c;把任意长度的数据转换为一个长度固定的数据串&#xff08;通常用16进制的字符串表示&#xff09;。 注…

浙江嘉兴徒步游

最近参加了一个徒步团&#xff0c;趁着周末时光&#xff0c;来了一场徒步旅游&#xff0c;不一样的体验图片发自简书App一开始进山探秘外蒲岛的路程&#xff0c;荒草丛生图片发自简书App树木郁郁葱葱&#xff0c;蓝天白云&#xff0c;一切都很没好图片发自简书App漫山遍野都开满…