Java I/O模型

event multiplexing

什么是 I/O multiplexing:当指定的I/O条件被满足时,操作系统会进行通知。

UNIX Network Programming中给出的一个描述:

1
What we need is the capability to tell the kernel that we want to be notified if one or more I/O conditions are ready (i.e., input is ready to be read, or the descriptor is capable of taking more output). This capability is called I/O multiplexing and is provided by the select and poll functions.

为了区分开下面三种I/O模型,我们需要了解下:

read操作分为两个步骤:

  • 等待数据从网上传输过来,写入内核的buffer
  • 把数据从内核的buffer复制到用户的buffer

阻塞I/O (Blocking I/O model):

调用read时,线程会被阻塞住,直到数据写入到用户的buffer。

​ 问题:需要大量的线程。费内存,线程切换对CPU的开销

image-20180803034421623

非阻塞I/O(Nonblocking I/O model)

read时,如果有数据,就和阻塞I/O一样;如果没有数据,就立即返回失败。

​ 问题:需要轮询,也很费CPU

image-20180803034547878

IO多路复用(I/O Multiplexing model):

使用一个线程来检查多个文件描述符的状态,看是否有数据到达,阻塞直到有一个文件描述符有数据到达。后续的操作,可以在同一个线程里面执行,也可以另起线程。

image-20180803034630907

I/O多路复用的多种实现

select

1
select(int nfds, fd_set *r, fd_set *w, fd_set *e, struct timeval *timeout)

poll

1
2
3
4
5
6
7
poll(struct pollfd *fds, int nfds, int timeout)

struct pollfd {
int fd;
short events;
short revents;
}

kqueue

1
2
3
4
5
6
7
8
9
10
11
12
int kqueue(void);
int kevent(int kq, const struct kevent *changelist, int nchanges,
struct kevent *eventlist, int nevents, const struct timespec *timeout);

struct kevent {
uintptr_t ident; /* identifier for this event */
int16_t filter; /* filter for event */
uint16_t flags; /* general flags */
uint32_t fflags; /* filter-specific flags */
intptr_t data; /* filter-specific data */
void *udata; /* opaque user data identifier */
};

NIO中的selector

Java doc中对它的描述:A mutiplexor of SelectableChannel objects

所以很自然的想到:他的底层实现应该就是基于 event multiplexing 的各个操作系统的实现: epoll, kqueue。

在macOS上,就是这个 KQueueSelectorImpl

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
class KQueueSelectorImpl extends SelectorImp{

protected int doSelect(long var1) throws IOException {
boolean var3 = false;
if(this.closed) {
throw new ClosedSelectorException();
} else {
this.processDeregisterQueue();

int var7;
try {
this.begin();
var7 = this.kqueueWrapper.poll(var1);
} finally {
this.end();
}

this.processDeregisterQueue();
return this.updateSelectedKeys(var7);
}
}
}

class KQueueArrayWrapper {
int poll(long var1) {
this.updateRegistrations();
int var3 = this.kevent0(this.kq, this.keventArrayAddress, 128, var1);
return var3;
}
private native int kevent0(int var1, long var2, int var4, long var5);
}

java中的selector最终会调用到native方法 kevent0.

NIO中的Channel

一个Channel(通道)代表和某一实体的连接,这个实体可以是文件、网络套接字等。

1
2
3
4
5
6
class ServerSocketChannelImpl extends ServerSocketChannel implements SelChImp{
private static NativeDispatcher nd;
private final FileDescriptor fd;
ServerSocket socket;
//......
}

ServerSocketChannel 通过register方法与一个Selector绑定。

NIO中的Buffer

Buffer中有3个很重要的变量,它们是理解Buffer工作机制的关键,分别是

  • capacity (总容量)
  • position (指针当前位置)
  • limit (读/写边界位置)

在对Buffer进行读/写操作前,我们可以调用Buffer类提供的一些辅助方法来正确设置 position 和 limit 的值,主要有如下几个

  • flip(): 设置 limit 为 position 的值,然后 position 置为0。对Buffer进行读取操作前调用。
  • rewind(): 仅仅将 position 置0。一般是在重新读取Buffer数据前调用,比如要读取同一个Buffer的数据写入多个通道时会用到。
  • clear(): 回到初始状态,即 limit 等于 capacity,position 置0。重新对Buffer进行写入操作前调用。
  • compact(): 将未读取完的数据(position 与 limit 之间的数据)移动到缓冲区开头,并将 position 设置为这段数据末尾的下一个位置。其实就等价于重新向缓冲区中写入了这么一段数据。

NIO server完整实例

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
public class NioServer {

public static void main(String[] args) throws IOException {
// 创建一个selector
Selector selector = Selector.open();

// 初始化TCP连接监听通道
ServerSocketChannel listenChannel = ServerSocketChannel.open();
listenChannel.bind(new InetSocketAddress(9999));
listenChannel.configureBlocking(false);
// 注册到selector(监听其ACCEPT事件)
listenChannel.register(selector, SelectionKey.OP_ACCEPT);

// 创建一个缓冲区
ByteBuffer buffer = ByteBuffer.allocate(100);

while (true) {
selector.select(); //阻塞,直到有监听的事件发生
Iterator<SelectionKey> keyIter = selector.selectedKeys().iterator();

// 通过迭代器依次访问select出来的Channel事件
while (keyIter.hasNext()) {
SelectionKey key = keyIter.next();

if (key.isAcceptable()) { // 有连接可以接受
SocketChannel channel = ((ServerSocketChannel) key.channel()).accept();
channel.configureBlocking(false);
channel.register(selector, SelectionKey.OP_READ);

System.out.println("与【" + channel.getRemoteAddress() + "】建立了连接!");

} else if (key.isReadable()) { // 有数据可以读取
buffer.clear();

// 读取到流末尾说明TCP连接已断开,
// 因此需要关闭通道或者取消监听READ事件
// 否则会无限循环
if (((SocketChannel) key.channel()).read(buffer) == -1) {
key.channel().close();
continue;
}

// 按字节遍历数据
buffer.flip();
while (buffer.hasRemaining()) {
byte b = buffer.get();

if (b == 0) { // 客户端消息末尾的\0
System.out.println();

// 响应客户端
buffer.clear();
buffer.put("Hello, Client!\0".getBytes());
buffer.flip();
while (buffer.hasRemaining()) {
((SocketChannel) key.channel()).write(buffer);
}
} else {
System.out.print((char) b);
}
}
}

// 已经处理的事件一定要手动移除
keyIter.remove();
}
}
}
}

image-20180422003741608

上面的例子可以用上图来描述,上图中的Acceptor,就是上面的代码运行的线程。

NIO与传统IO的区别

NIO和传统IO(一下简称IO)之间第一个最大的区别是,IO是面向流的,NIO是面向缓冲区的。 Java IO面向流意味着每次从流中读一个或多个字节,直至读取所有字节,它们没有被缓存在任何地方。此外,它不能前后移动流中的数据。如果需要前后移动从流中读取的数据,需要先将它缓存到一个缓冲区。NIO的缓冲导向方法略有不同。数据读取到一个它稍后处理的缓冲区,需要时可在缓冲区中前后移动。这就增加了处理过程中的灵活性。

IO的各种流是阻塞的。这意味着,当一个线程调用read() 或 write()时,该线程被阻塞,直到有一些数据被读取,或数据完全写入。

但是,由于NIO的非阻塞的特点,就导致了一些其他问题,编码复杂,而且存在半包问题。

Reactor模型

Reactor模型中的组件

  • Reactor:Reactor是IO事件的派发者。
  • Acceptor:Acceptor接受client连接,建立对应client的Handler,并向Reactor注册此Handler。
  • Handler:和一个client通讯的实体,按这样的过程实现业务的处理。

对应上面的NIO代码来看:

  • Reactor:相当于有分发功能的Selector
  • Acceptor:NIO中建立连接的那个判断分支
  • Handler:消息读写处理等操作类

Reactor单线程模型

和NIO很相似,只是对IO读到的数据的处理(有数据可以读取的分支里面的代码),放到了Handler中。

image-20180422010715580

Reactor多线程模型

将数据的处理代码,改为多线程实现。

image-20180422010833942

主从Reactor模型

因为Reactor既要处理IO连接请求,又要处理Read,所以出现了下面的方式,MainReactor处理IO连接请求,SubReactor处理Read等。

image-20180422011045446

引用

https://people.eecs.berkeley.edu/~sangjin/2012/12/21/epoll-vs-kqueue.html

http://www.ivaneye.com/2016/07/23/iomodel.html

https://www.cnblogs.com/coderjun/p/7100423.html