MongoDB

MongoDB is a document database with the scalability and flexibility that you want with the querying and indexing that you need

上面是mongodb官方网站对 mongoDB的描述。mongoDB有什么不一样的地方?

特点1:flexibility

可以很方便的扩展结构,使用JSON-like的结构来保存document。使用过mysql的人应该深有体会,添加一个字段是多么的痛苦。

在mongodb的官方文档,是这样评价Document Model的: The best way to work with Data。

主要从下面4个方面来论证:

Easy

关系数据库使用tabular data model,由于数据规范化的要求,一个现实中的实体,经常需要好多表来表示。 带来的问题是:开发效率降低。虽然可以通过ORM层(object-relational mapping layer)来解决,但是也是会带来其他问题。

对比之下,mongo使用了document data model,document可以以更加自然的方式来描述数据,就是使用一个单一的文档,把相关的数据作为sub-document嵌入其中。

Flexible

document的格式可以随便变更,确实很方便。但是会带来其他的问题么?

Fast

mongoDB的数据模型,性能会更好?关系数据库可以JOIN,JOIN的性能确实不好。Mongo因为一个实体被映射为单个文档,所以不需要JOIN,所以性能更好。。虽然也提供了类似JOIN的功能:$lookup,可以跨不同的collection之间。

Versatile (多功能)

参考:

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

Concurrency

如何实现高并发

I/O 策略的选择

首先,有哪些I/O model 可以选择?

UNIX Netwrok Programming中,列出了以下5种I/O模型:

blocking I/O

nonblocking I/O

I/O multiplexing (select and poll)

singnal driven I/O (SIGIO)

asynchronous I/O (the POSIX aio_functions)

我们先不仔细研究,先看下 The C10K problem中,列出了哪些 I/O Strategies。

1 一个线程服务多个客户端,使用非阻塞I/O和水平触发的就绪通知

2 一个线程服务多个客户端,使用非阻塞I/O和边缘触发的就绪通知

3 一个服务线程服务多个客户端,使用异步I/O

4 一个服务线程服务多个客户端,使用阻塞I/O

5 把服务代码编译进内核

1,2 都是使用单个线程来关注一些列nonblocking sockets,看它们何时ready for I/O。但是1,2所说的非阻塞I/O不是上面列出的5种I/O模型中的第二个,它指的应该是第一种的对立面,主要包含了第三种,既I/O multiplexing模式。

1 一个线程服务多个客户端,使用非阻塞I/O和水平触发的就绪通知

实现方式:

  • 传统的select:受限于句柄个数
  • 传统的poll:性能低,因为扫码大量的描述符很耗费时间
  • kqueue(level-trigger):for FreeBSD

2 一个线程服务多个客户端,使用非阻塞I/O和边缘触发的就绪通知

实现方式:

  • kqueue(edge-trigger):
  • epoll:Linux2.6中推荐使用的edge-triggered poll
  • Realtime Signal:Linux2.4中推荐使用的edge-triggered poll

####关于水平触发边缘触发

水平触发(level-triggered)——只要满足条件,就触发一个事件(只要有数据没有被获取,内核就不断通知你);边缘触发(edge-triggered)——每当状态变化时,触发一个事件。

select,poll,Epoll区别:

select poll Poll
支持最大连接数 1024(x86) 2048(x64) 无限制 无限制
I/O效率 每次调用遍历所有 每次调用遍历所有 使用“事件”通知方式,每当fd就绪,系统注册的回调函数就会被调用,将就绪fd放到rdllist里面,这样epoll_wait返回的时候我们就拿到了就绪的fd。时间发复杂度O(1)
fd拷贝 每次调用需要拷贝 每次调用拷贝 调用epoll_ctl时拷贝进内核并由内核保存,之后每次epoll_wait不拷贝

3一个服务线程服务多个客户端,使用异步I/O

该方法目前还没有在Unix上普遍的使用,可能因为很少的操作系统支持异步I/O。在标准Unix下,异步I/O是由“aio_”接口 提供的,既上面5种模型中的最后一种:asynchronous I/O 。它把一个信号和值与每一个I/O操作关联起来。信号和其值的队列被有效地分配到用户的 进程上。

关于异步I/O和同步I/O

UNIX Netwrok Programming中有这样的描述:

1
2
3
POSIX defines these two terms as follows:
A synchronous I/O operation causes the requesting process to be blocked until that I/O operation completes.
An asynchronous I/O operation does not cause the requesting process to be blocked.

所以,上面的5中模型中,1~4是同步I/O,5才是异步I/O。

关于阻塞I/O和非阻塞I/O

阻塞I/O :BIO , 非阻塞I/O:NIO ,

一个IO操作其实分成了两个步骤:发起IO请求和实际的IO操作,阻塞IO和非阻塞IO的区别在于第一步:发起IO请求是否会被阻塞,如果阻塞直到完成那么就是传统的阻塞IO;如果不阻塞,那么就是非阻塞IO

4一个服务线程服务多个客户端,使用阻塞I/O

一个线程一个socket?

5 把服务代码编译进内核

并发模型

实现并发的途径有两种,基于线程和基于事件。

并发和并行

并发concurrency属于问题域(problem domain), 并行parallelism属于( solution domain)。并行和并发的区别在于有无状态,并行计算适合无状态应用,而并发解决的是有状态的高性能; 有状态要着力解决并发计算,无状态要着力并行计算。

参考:

Mysql InnoDB

InnoDB中数据是如何存储的

system tablespace : 一个数据文件,用来存储一个或多个InnorDB 以及相关的 索引

image-20180412010618020

定义

InnoDB数据目录

重做日志

在InnoDB的数据目录下,会有2个文件:ib_logfile0,ib_logfile1,他们就是重做日志(redo log)。日志的大小会影响数据库的性能。写入的时候,是先写入redo log buffer,然后按照一定的条件写入日志文件。从缓存写入磁盘的时候,是以512字节(一个扇区)为单位写入的。因为扇区是写入的最小单元,所以写入必定成功。从缓存写入磁盘,有几种选择:

  • 0 事物提交,不将重做日志写入磁盘,等待主线程每秒的刷新
  • 1 事物提交,写入磁盘,并调用fsync
  • 2 事物提交,写入磁盘,不fsync,其实是写入的文件系统的缓存。

The Little Book of Semaphores

###Semaphores
现实世界的semaphores是一个用来沟通的信号系统。计算机世界的semaphores是一个数据结构,被用来解决synchronization问题。Semaphores的发明者是Edsger Dijkstra

###定义
Semaphores是一个Integer,再加上下面的约束:

  • 创建的时候,可以指定任意的数字。只能对它执行increment和decrement操作,不能直接访问semaphore的int value。
  • 当一个thread decrement一个semaphore:如果结果为负,则thread阻塞。但是,在执行decrement之前,无法知道当前thread是否会被阻塞。
  • 当一个thread increment 一个semaphore:选择一个阻塞的thread使其被唤起。然后,两个thread都继续执行。

###语法
increment :signal 或者 V
decrement:wait 或者P

如果说要一个完整的函数名,可能是下面的形式:

1
2
fred.increment_and_wake_a_waiting_process_if_any()
fred.decrement_and_block_if_the_result_is_negative()

increment 和decrement 描述了方法做了什么;

signal 和 wait 描述了方法被用来做什么;

V 和 P 是Dijkstra提出的原语。

基本的同步模式 (Basic synchronization patterns)

Signaling

可以用来保证一个线程中的一段代码会早于另外一个线程中的一段代码运行。

Thread A

1
2
statement a1
sem.signal()

Thread B

1
2
sem.wait()
statement b1

a1 > b1

Rendezvous

要求: a1 > b2 and b1 > a2

Thread A

1
2
3
4
statement a1
aArrived.signal()
bArrived.wait()
statement a2 // critical point

Thread B

1
2
3
4
statement b1
bArrived.signal()
aArrived.wait()
statement b2 // critical point

mutex

mutex很像一个在线程间传递的token,获得到token的线程可以处理。

Mutex 是mutual exclusion的缩写。用mutex去保护critical regions

Thread A 和 Thread B

1
2
3
4
mutex.wait()
#critical section
count = count + 1 //for example
mutex.signal()

Multiplex

和mutex和像,只是容许多个线程同时进入critical regions。为了实现这样的效果,直接将semaphores 初始化为n(n个线程同时进入critical regions)。

1
2
3
multiplex.wait()
//critical section
multiplex.signal()

Barrier

Rendezvous只容许两个线程。但是如果容许多个线程,就是Barrier,要求所有线程都不能执行critical point,除非所有的线程都已经执行了rendezvous。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
n = the number of threads
count = 0 //有多少线程已经到来了
mutex = Semaphore(1) //对count的修改,进行并发控制
barrier = Semaphore(0) //一直阻塞,直到所有的线程到达。

rendezvous
mutex.wait()
count = count + 1
mutex.signal()

if count == n: barrier.signal()

barrier.wait() // 这样的wait和signal紧接着的情况,很常见,被叫做turnstile(旋转门)
barrier.signal() // 因为它既可以控制线程们一个个的通过,也可以被锁住,不让所有的线程通过。

cirtical point

Reusable barrier

当所有的线程通过后,再次不容许任何线程通过,就可以再次使用。也被叫做two-phase barrier

Preloaded turnstile

turnstile是一个常用的组件,它有个不好的地方就是强制线程一个个的通过。可能造成大量的线程切换。

基于reusable barrier ,如果最后一个打开旋转门的线程,可以预加载足够多的signal,就可以让相同数量的线程通过旋转门。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
# rendezvous
mutex.wait()
count += 1
if count == n:
turnstile.signal(n) // unlock the first
mutex.signal()

turnstile.wait() // first turnstile
# critical point

mutex.wait()
count -= 1
if count == 0:
turnstile2.signal(n) // unlock the second
mutex.signal()

turnstile2.wait() // second turnstile

Queue

//TODO

Netty

Channel Sockets
EventLoop Control flow, multithreading, concurrency
ChannelFuture Asynchronous notification
–Netty in Action

Channel

在java中,基础I/O操作的依赖的基石就是Socket,而Channel就是提供了一组API,极大的化简了对Socket的操作。

EventLoop

EventLoop处理一个Connection生命周期内的所有事件。
几个概念间的关系:
EventLoopGroup包含一个或者多个EventLoop;
EventLoop会绑定一个Thtread。
一个Channel会注册到一个EventLoop上。
一个EventLoop可能会被分配给多个Channel。

因为I/O操作在Netty中,都是异步的,没有一个操作可以立即返回。所以需要ChannelFuture。

ChannelHandler

ChannelHandler:是应用开发者最关注的,它被网络事件触发。处理inbound,outbound数据。
ChannelPipeline:包含一连串的ChannelHandler。而且,当Channel被创建时,就会被指定一个ChannelPipeline。细节:ChannelInitializer会负责安装一系列的ChannelHandler到Pipeline中。
Event顺着ChannelPipeline流动,中间经过一个个的ChannelHandler。PipeLine是双向流动的,所以会有两种ChannelHandler:ChannelInboundHandler 和 ChannelOutboundHandler。
ChannelHandlerContext保存着ChannelHandler和ChannelPipeline的绑定关系。
Adapter:ChannelhandlerAdapter是包含了基础实现的ChannelHandler。
Encoders 和 Decoders:netty提供了很多这样的抽象类,它们都是ChannelHandler。

服务端,需要两套Channels,既两个EventLoopGroup。其中一个与ServerChannle相关联,里面的EventLoop负责为到来的l连接请求创建Channel,一旦这个channel被创建,另外一个EventLoopGroupj就会分配一个EventLoop,来和这个channel绑在一起。

参考: http://owvvyywxv.bkt.clouddn.com/nio.pdf

Java AQS

java AbstractQueuedSynchronizer

###以ReentrantLock为例

首先ReentrantLock有公平和非公平模式:NonfairSyncFairSync

FairSync

lock方法的定义:获取锁。

  • 如果锁没有被其他线程获取,则获取锁并立刻返回,把锁的count设为1.
  • 如果锁已经被当前线程持有,则立刻返回,把锁的count +1.
  • 如果锁被其他线程持有,则当前线程会失去被调度的权利,等待被唤醒。

假设有3个线程同时调用ReentrantLock的 lock方法。

image-20180824040623739

结合代码的解释可以参考:https://www.jianshu.com/p/d8eeb31bee5c

Synchronizer 同步器

同步队列,同步节点,同步状态

AbstractQueuedSynchronizer 被用来实现blocking lock或者相关的同步器(semaphores,events)。这个类的基础是 FIFO waiting queue。且用一个原子的int值来表示同步状态。

###子类该怎么定义?

子类应该被定义为一个非public的内部帮助类。

这个类为internal queue提供了各种便利的方法,同时也适用于condition objects。你可以引入这些到你的类里,做一些synchronization mechanics。

独占式同步状态的获取,释放

通过调用同步器的 acquire(int arg)方法来获取到同步状态。

线程获取锁失败后,进入同步队列。

首先,调用tryAcquire方法,尝试获取锁。如果获取失败,则构造同步节点,并通过addWaiter方法将同步节点加入到同步队列的尾部,最后调用acquireQueued方法,使该同步节点以死循环的方式获取到锁,且非中断。只有前置节点为头节点并且tryAcquire成功返回的时候,才会返回。否则就会让线程停止调度,直到unpark或者线程中断。但是,如果前置节点为头结点,但是tryAcquire返回失败,怎么办?

通过同步器的relase方法来释放同步状态。

该方法会先调用tryRelease方法,然后会唤醒其他后续节点,通过调用unparkSuccessor方法,即调用了unpark方法。

同步节点进入同步队列之后,就进入了一个自旋的过程,每个同步节点都在自己观察,当条件满足,就可以从自旋中退出。

###用法:

如果使用这个类作为自定义的synchronozier的基础。

#####1必须实现下面的方法:

tryAcquire : 查询是否可以获取,如果可以则获取。
tryRelease
tryAcquireShared
tryReleaseShared
isHeldExclusively

上面的方法的实现,要尽量简短,不能阻塞

#####2 同步状态的修改(Synchronization state)
getState
setState
compareAndSetState
来修改同步状态,可以保证状态的改变是安全的。

######3 The CORE of exclusive synchronization:

1
2
3
4
5
6
7
8
9
Acquire:
while (!tryAcquire(arg)) {
//enqueue thread if it is not already queued;
//possibly block current thread;
}

Release:
if (tryRelease(arg))
//unblock the first queued thread;

#####Barging Strategy
TODO

#####Examples

举了2个例子:MutexBooleanLatch

Mutext
什么是mutex?Mutual(彼此的) exclusion 的缩写。我们使用mutex去保护临界区,从而防止竞争。进入临界区之前调用 acquire ,离开临界区之前调用release

1
2
3
4
5
6
7
8
9
10
aquire(){
while (!available){
;// busy wait
}
available = false;
}

release() {
available = true;
}

BooleanLatch
BooleanLatch是non-exclusive,所以使用shared aquire 和 release方法。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
class BooleanLatch {

private static class Sync extends AbstractQueuedSynchronizer {
boolean isSignalled() { return getState() != 0; }

protected int tryAcquireShared(int ignore) {
return isSignalled() ? 1 : -1;
}

protected boolean tryReleaseShared(int ignore) {
setState(1);
return true;
}
}

private final Sync sync = new Sync();
public boolean isSignalled() { return sync.isSignalled(); }
public void signal() { sync.releaseShared(1); }
public void await() throws InterruptedException {
sync.acquireSharedInterruptibly(1);
}

waitStatus 讲解

Latch

作用:A synchronization aid that allows one or more threads to wait until a set of operations being performed in other threads completes.

主要方法:

await() 阻塞等待the count of latch变为0.

CountDown() 给the count of latch -1.