导图社区 Java并发编程的艺术
这是一篇关于Java并发编程的艺术的思维导图
提示: 本内容由社区用户上传并分享。平台不对内容的真实性、合法性、知识产权归属及是否侵害第三方权利进行事前审核或保证。本内容可能包含受版权保护的图片、字体或其他第三方素材,使用前请自行确认授权范围。
Java并发编程的艺术
Java并发容器和框架
ConcurrentHashMap的实现原理与使用(其中的代码的版本为Java7)
ConcurrentHashMap是线程安全且高效的HashMap
为什么要使用ConcurrentHashMap?
(1)在并发编程时使用HashMap可能导致程序死循环
多个线程执行put()方法,发生扩容时可能导致Entry形成循环链表,即next节点永远不为空,导致获取Entry时产生死循环
Java7链表新节点采用的是头插法,这样在线程一扩容迁移元素时,会将元素顺序改变,导致两个线程中出现元素的相互指向而形成循环链表,Java8采用了尾插法,从根源上杜绝了这种情况的发生
参考链接
https://www.jianshu.com/p/1e9cf0ac07f4
https://www.jianshu.com/p/0df1f25139e4
(2)线程安全的Hashtable效率较低
Hashtable内部使用了synchronized保证线程安全,在多个线程进行竞争时会导致阻塞,效率降低
(3)ConcurrentHashMap的锁分段技术可以有效提升并发访问率
Hashtable在并发环境下效率较低主要是因为多个线程竞争同一把锁,而ConcurrentHashMap使用的锁分段技术把数据分为一段一段存储,每一段配一把锁,当一个线程占用一把锁访问一段数据时,其他段的数据也能被其他线程访问
ConcurrentHashMap的结构
类图
ConcurrentHashMap由Segment数组结构和HashEntry数组结构组成,一个ConcurrentHashMap包含一个Segment数组,一个Segment包含一个HashEntry数组
Segment:是一种可重入锁(ReentrantLock),结构类似于HashMap,是一种数组和链表结构,扮演锁的角色
HashEntry:是一种链表结构的元素,用于存储键值对数据
每个Segment守护着一个HashEntry数组里的元素,当对HashEntry数组的数据进行修改时,必须首先要获取与之对应的Segment锁
ConcurrentHashMap的初始化
ConcurrentHashMap初始化方法主要是使用initialCapacity、loadFactor和concurrencyLevel等参数来初始化segment数组、段偏移量segmentShift、段掩码segmentMask和每个segment里的HashEntry数组
(1)初始化segments数组
segments数组的长度ssize是通过concurrencyLevel计算出的,为了能通过按位与的散列算法来定位segments数组的索引,必须保证该数组的长度是2的N次方,所以ssize的值为大于或等于concurrencyLevel的最小的2的N次方值
注意:concurrencyLevel最大值为65535,所以segments数组的长度最大为65536,即16位
(2)初始化segmentShift和segmentMask
sshift = ssize从1向左移位的次数。concurrencyLevel默认为16,ssize需要左移4次,所以sshift等4
segmentShift用于定位要参与散列运算的位数,值为32减sshift,所以是28
为什么是32减sshift?因为ConcurrentHashMap的hash()方法输出的最大数为32位
segmentMask是散列运算的掩码,其值为ssize减1,即15(二进制为1111)
ssize的最大值为65536,所以segmentShift的最大值为16,segmentMask的最大值为65535,对应的二进制是16位,每个位都是1
(3)初始化每个segment
代码中,cap为segment中HashEntry数组的长度,其值为initialCapacity除以ssize的倍数c,若c大于1,则取大于等于c的2的N次方值,否则取1
segment的容量threshold = (int)cap * loadFactor,默认initialCapacity等于16,loadFactor等于0.75,cap等于1,threshold等于0
initialCapacity是ConcurrentHashMap的初始容量,loadFactor是segment的负载因子,在构造方法里需要使用这两个参数初始化数组中的每个segment
定位Segment
ConcurrentHashMap使用分段锁Segment保护不同段的数据,在插入和获取元素时,必须先通过散列算法定位到Segment
(1)首先使用Wang/Jenkins hash 的变种算法对元素的 hashCode 进行一次再散列。减少散列冲突,使元素均匀的分布在不同的Segment上
(2)使用此散列算法定位segment
ConcurrentHashMap的操作
(1)get()操作
首先进行一次再散列,使用该散列值通过散列运算定位到Segment,然后再通过散列算法定位到元素
为什么不用加锁?因为get()方法里将需要共享的变量声明为volatile类型,保证多个线程间的可见性,能够被多个线程同时读
(2)put()操作,首先定位到Segment,然后需要经历以下两个步骤
(1)判断Segment中的HashEntry数组是否需要扩容,插入前先判断Segment中的HashEntry数组是否超过容量(threshold),超过阈值则进行扩容
如何扩容?首先创建一个容量为原数组容量的两倍的数组,然后将原数组中的元素再散列后插入到新数组中。ConcurrentHashMap不会对整个容器进行扩容,而只对某个segment进行扩容
(2)定位添加元素的位置,将其添加到HashEntry数组中
(3)size()操作,想要统计ConcurrentHashMap的大小,就必须统计所有Segment的大小然后求和。
ConcurrentHashMap先尝试两次不锁住Segment的方式统计各个Segment的大小,如果统计过程中,count发生了变化,则采用锁住Segment的方式统计各个Segment的大小
如何知道统计过程中count是否发生了变化?使用modCount变量,put()、remove()等方法会将modCount+1,然后在统计size前后比较modCount是否发生了变化
ConcurrentLinkedQueue
是一个基于链接节点的无界线程安全队列,按先进先出的规则对节点进行排序,采用了“wait-free”算法(即CAS算法)来实现
ConcurrentLinkedQueue的结构
ConcurrentLinkedQueue是一个基于链表实现的队列,由head和tail节点组成,默认情况下head存储的元素为空,head = tail
入队列
(1)入队列的过程
主要步骤
(1)将入队节点设置为当前队尾节点的下一个节点
(2)更新tail节点,如果tail节点的next不为空,则将入队节点设置为tail节点,否则将其设置为tail节点的next节点,所以tail节点并不总是尾节点
源码分析
从源码角度看,入队过程主要是定位出尾节点,然后使用CAS算法将入队节点设置成尾节点的next节点,不成功则重试
定位尾节点
每次入队必须通过tail节点来查找尾节点,因为tail节点并不总是尾节点,尾节点有可能是tail节点或者是tail节点的next节点
需要注意p节点等于p的next节点的情况,此时表示这个队列刚初始化,所以需要返回head节点
设置入队节点为尾节点
p.casNext(null,n)设置当前入队节点为当前队列尾节点的next节点,如果p的next是null,则p是尾节点,如果不是null,则表示其他线程更新了尾节点,需要重新获取尾节点
HOPS的设计意图
使用HOPS可以减少循环CAS更新tail节点的次数,提高入队的效率,只有当tail节点和尾节点的距离大于等于常量HOPS的值时才更新tail节点,这样更新tail次数会变少,定位尾节点次数会变多,但仍然可以提高入队效率
本质上是通过增加对volatile变量的读来减少对volatile变量的写,因为对volatile变量的写的开销要大于读的开销
注意:入队方法永远返回true,不要通过返回值来判断是否成功
出队列
从队列中返回一个节点元素,并情况该节点对元素的引用
并不是每次出队都更新head节点,如果head节点不为空直接弹出head节点中的元素,否则才会更新head节点,这种操作也是通过hops变量减少CAS更新head节点的消耗,提升效率
源码分析
首先获取并判断头元素是否为空,为空则表示其他线程已经执行了一次出队,不为空则使用CAS操作将头节点引用设置为null,若CAS成功,则返回头节点的元素,否则表示其他线程进行了一次出队,头节点已经改变,需要重新获取头节点
Java中的阻塞队列
阻塞队列(BlockingQueue)是支持两个附加操作的队列,这两个附加操作支持阻塞的插入和移除方法
支持阻塞的插入方法:当队列满时,队列会阻塞插入元素的线程,直到队列不满
支持阻塞的移除方法:当队列为空时,获取元素的线程会等待队列变为非空
阻塞队列常用于生产者和消费者场景,生产者是向队列中添加元素的线程,消费者是从队列中获取元素的线程
阻塞队列不可用时,两个附加操作提供了4中处理方式
抛出异常:当队列满时,再向队列插入元素,会抛出IllegalStateException("Queue full")异常。当队列空时,从队列获取元素时会抛出NoSuchElementException
返回特殊值:当向队列中插入元素时,会返回元素是否插入成功,成功则返回true。如果是移除方法,则会从队列中取出一个元素,如果没有则返回null
一直阻塞:当队列满时,如果生产者线程向队列中put元素,队列会一直阻塞生产者线程,直到队列可用或响应中断退出。当队列空时,如果消费者线程从队列中take元素,队列会阻塞住消费者线程,直到队列可用
超时退出:当队列满时,如果生产者线程向队列中插入元素,队列会阻塞生产者线程一段时间,如果超过指定的时间,生产者线程就会退出
注意:如果是无界阻塞队列,队列不可能出现满的情况,所以如果使用put或offer方法永远不会被阻塞,而且使用offer方法时,该方法永远返回true
Java里的阻塞队列
ArrayBlockingQueue:用数组实现的有界阻塞队列。此队列按照先进先出(FIFO)的原则对元素进行排序
公平访问队列:阻塞的线程可以按照阻塞的先后顺序访问队列,即先阻塞的线程先访问队列
默认情况下不保证线程公平的访问队列,队列可用时,阻塞的线程都可以争夺访问队列 的资格,有可能先阻塞的线程最后才访问到队列
如何创建公平的阻塞队列?
LinkedBlockingQueue:用链表实现的有界阻塞队列,先进先出(FIFO),默认和最大长度为Integer.MAX_VALUE
PriorityBlockingQueue:支持优先级的无界阻塞队列,默认情况下元素采取自然顺序升序排列,可以自定义类实现compareTo()方法或初始化PriorityBlockingQueue时指定构造参数Comparator来实现自定义排序规则。注意:不能保证同优先级元素的顺序
DelayQueue:支持延时获取元素的无界阻塞队列,使用PriorityQueue实现。
队列中的元素必须实现Delayed接口,在创建元素时可以指定多久才能从队列中获取当前元素
如何实现Delayed接口?参考ScheduledTheadPoolExcecutor里ScheduledFutureTask类的实现,需要三步
(1)对象创建的时候初始化基本数据,比如记录当前对象延迟到何时可以使用和元素在队列中的先后顺序等
(2)实现getDelay()方法,返回当前元素还需要延时多长时间,单位为纳秒。注意:time小于当前时间时,getDelay()会返回负数
(3)实现compareTo()方法来指定元素的顺序
只有在延迟期满时才能从队列中提取元素
应用场景
缓存系统的设计:可以使用DelayQueue保存元素的有效期,使用一个线程循环查询DelayQueue,一旦能从DelayQueue中获取元素时,说明缓存有效期到了
定时任务调度:使用DelayQueue保存当天将要执行的任务和执行时间,一旦从DelayQueue中获取到任务就开始执行,比如TimerQueue就是使用DelayQueue实现的
如何实现延时阻塞队列?
当消费者从队列中获取元素时,如果元素没有到达延时时间,就阻塞当前线程
代码中的leader是等待获取队列头部元素的线程
SynchronousQueue:是一个不存储元素的阻塞队列,每一个put操作必须等待一个take操作,否则不能继续添加元素
支持公平访问队列,默认使用非公平策略访问队列
设置为true则线程采用先进先出的顺序访问队列
此队列本身不存储任何元素,非常适合传递数据的场景,负责把生产者线程处理的数据直接传递给消费者线程
LinkedTransferQueue:由链表结构组成的无界阻塞TransferQueue队列,相对于其他阻塞队列,此队列多了tryTransfer()方法和transfer()方法
public void transfer(E e):如果当前有消费者在等待元素(消费者使用take()或带时间的poll()方法时)此方法可以将数据立刻传输给消费者,否则将元素放在队列的tail节点,并等到该元素被消费者消费后才返回
public boolean tryTransfer(E e):试探生产者传入的元素是否能直接传给消费者,如果没有消费者等待接收,则返回false。此方法无论消费者是否消费了元素,都会直接返回,而transfer()方法会等待元素被消费后才返回
public boolean tryTransfer(E e, long timeout, TimeUnit unit):试图把生产者传入的元素直接传给消费者,如果没有消费者消费则等待指定的时间,若超时还没消费元素则返回false,否则返回true
LinkedBlockingDeque:由链表结构组成的双向阻塞队列。
双向阻塞队列:可以从队列的两端插入和移除元素
多了一个操作队列的入口,在多线程同时入队时,就减少了一半的竞争
LinkedBlockingDeque 多了addFirst、addLast、offerFirst、offerLast、peekFirst 和 peekLast 等方法,以 First 单词结尾的方法,表示插入、获取(peek)或移除双端队列的第一个元素。以 Last 单词结尾的方法,表示插入、获取或移除双端队列的最后一个元素。add等同于addLast,remove等同于removeFirst,take等同于takeFirst
初始化此队列时可以指定容量,防止其过度膨胀。
双向阻塞队列可以运用在“工作窃取”模式中
阻塞队列的实现原理
使用通知模式实现,所谓通知模式是指当生产者向满的队列里添加数据时,会阻塞生产者,当消费者消费了队列的一个元素时,会通知生产者队列可用
ArrayBlockingQueue内部使用Condition实现
await()方法内部使用了LockSupport.park()阻塞生产者线程
LockSupport.park()内部使用了unsafe.park()方法
public native void park(boolean isAbsolute, long time),此方法会阻塞当前线程,只有当以下四种情况的一种发生时,该方法才会返回
(1)与park()对应的unpark()执行或已经执行,已经执行是指先执行unpark()后执行park()
(2)线程被中断时
(3)等待完time参数指定的毫秒数时
(4)发生异常时
当线程被阻塞队列阻塞时,会进入WAITING(parking)状态
Fork/Join框架
什么是Fork/Join框架?
Fork/Join框架是Java7提供的一个并行执行任务的框架,将一个大任务切分为若干个小任务,然后汇总每个小任务的结果得到大任务的结果
Fork:将一个大任务切分为若干个小任务并行执行
Join:将每个小任务的结果汇合得到大任务的结果
工作窃取算法
工作窃取(work-stealing)算法是指某个线程从其他队列中窃取任务来执行
为何需要工作窃取算法?
将一个大任务切分为多个互不依赖的子任务并放在不同的队列中,并为每个队列单独创建一个线程来执行队列中的任务,线程和队列一一对应;当某个线程执行完自己队列中的任务时,可以从其他线程的队列中窃取任务来执行,为了减少被窃取任务的线程和窃取任务的线程之间的竞争,通常使用双端队列,被窃取任务的线程从队头取任务,窃取任务的线程从队尾取任务
优点:充分利用线程进行并行计算,减少线程间的竞争
缺点:在某些情况下还是会存在线程间的竞争,比如当双端队列中只有一个任务时。会消耗更多资源,比如创建了多个线程和双端队列
Fork/Join框架的设计
步骤一:fork类将大任务分割成子任务,可能子任务还是很大,需要不停分割,直到分割出的子任务足够小
步骤二:执行任务并合并结果:分割出的子任务分别放在双端队列中,然后启动几个线程分别从双端队列中获取任务并执行,子任务的执行结果统一放在一个队列中,启动一个线程从这个队列中拿数据并合并
ForkJoinTask:使用Fork/Join框架,需要创建ForkJoin任务,它提供在任务中执行fork()和join()操作的机制
子类
RecursiveAction:用于没有返回结果的任务
RecursiveTask:用于有返回结果的任务
ForkJoinPool:ForkJoinTask需要通过ForkJoinPool来执行
任务分割出的子任务会添加到当前工作线程维护的双端队列的队头,当工作线程的队列中暂时没有任务时,它会随机从其他工作线程的队列的队尾获取一个任务
使用Fork/Join框架
在ForkJoinTask的compute()方法中需要判断,如果任务足够小则执行,否则继续分割任务,子任务执行fork()时会再次进入compute()方法看子任务是否还需要分割,如果不需要则执行子任务并返回结果。使用join()方法会等待子任务执行完毕并获取其结果
Fork/Join框架的异常处理
ForkJoinTask在执行时可能出现异常,在主线程中无法直接捕获
public final boolean isCompletedAbnormally():用来检查ForkJoinTask中是否抛出了异常或已经被取消
public final Throwable getException():获取异常,返回Throwable对象,任务被取消返回CancellationException,如果任务没有完成或没有出现异常则返回null
Fork/Join框架的实现原理
ForkJoinPool由ForkJoinTask数组和ForkJoinWorkerThread组成
ForkJoinTask数组存储程序提交给ForkJoinPool的任务
fork():当调用fork()方法时,程序会调用ForkJoinWorkerThread的pushTask()方法异步的执行这个任务,然后立即返回结果
pushTask():此方法会把当前任务放在ForkJoinTask数组队列里,然后再调用ForkJoinPool的signalWork()唤醒或创建一个线程来执行任务
join():阻塞当前线程并等待获取结果,其中调用了doJoin()获取任务的状态判断返回什么结果
任务的四种状态
已完成(NORMAL)
被取消(CANCELLED)
信号(SIGNAL)
出现异常(EXCEPTIONAL)
doJoin()
已完成,直接返回任务状态
未执行完,则从任务数组中取出任务并执行
任务顺利完成,状态为NORMAL
出现异常,状态为EXCEPTIONAL
ForkJoinWorkerThread数组负责执行ForkJoinTask数组中的任务
Java中的13个原子操作类
JDK1.5提供的原子操作类提供了一种线程安全地更新一个变量的方式。Atomic 包里的类基本都是使用 Unsafe 实现的包装类
原子更新基本类型类
AtomicBoolean:原子更新布尔类型
AtomicInteger:原子更新整型
int addAndGet(int delta):以原子方式将输入的数值与实例中的值(AtomicInteger中的value)相加,并返回结果
boolean compareAndSet(int expect,int update):如果实例中的值等于预期值,则以原子方式设置该值为给定的更新值
int getAndIncrement():以原子方式将当前值+1,并返回自增前的值
(1)获取AtomicInteger的当前值
(2)将当前值+1
(3)调用compareAndSet(current,next)方法进行原子更新,首先会判断当前实例的值是否等于current,相等说明没有其他线程修改过实例的值,则将实例的值更新为next,否则compareAndSet()会返回false,程序进入for循环重新执行compareAndSet()操作
void lazySet(int newValue):最终会设置成newValue,使用lazySet()设置值后,可能会导致其他线程在之后的一小段时间内读取到的还是旧值
参考链接:http://ifeve.com/how-does-atomiclong-lazyset-work/
int getAndSet(int newValue):以原子方式设置实例的值为newValue,并返回旧的值
AtomicLong:原子更新长整型
这三个类提供的方法几乎一样,所以只以AtomicInteger为例
如何原子的更新其他基本类型呢?
Atomic包中的类基本上都是使用Unsafe实现的
参考AtomicBoolean,它是先把Boolean转成整型,然后使用compareAndSwapInt()进行CAS,所以原子更新char、float和double也可以用类似的思路来实现
原子更新数组
通过原子的方式更新数组里的某个元素
AtomicIntegerArray:原子更新整型数组里的元素
int addAndGet(int i, int delta):以原子方式将输入值delta与数组索引为i的元素相加
boolean compareAndSet(int i, int expect, int update):如果当前值等于预期值expect,则以原子方式将数组索引为i的元素设置为update
AtomicLongArray:原子更新长整型数组里的元素
AtomicReferenceArray:原子更新引用类型数组里的元素
这些原子类会将通过构造方法传入的数组value复制一份,所以修改其内部的数组元素不会影响到传入的数组
原子更新引用类型
如果要原子更新多个变量,就要使用原子更新引用类型提供的类
AtomicReference:原子更新引用类型
AtomicReferenceFieldUpdater:原子更新引用类型里的字段
AtomicMarkableReference:原子更新带有标记位的引用类型。可以原子更新一个布尔类型的标记位和引用类型,构造方法是:AtomicMarkableReference(V initialRef,boolean initialMark)
原子更新字段类
如果需要原子地更新某个类的某个字段时,就需要使用原子更新字段类
AtomicIntegerFieldUpdater:原子更新整型的字段更新器
AtomicLongFieldUpdater:原子更新长整型的字段更新器
AtomicStampedReference:原子更新带有版本号的引用类型。该类将整数值与引用关联起来,可用于原子的更新数据和版本号,可以解决使用CAS进行原子更新时可能出现的ABA问题
更新步骤
(1)使用静态方法newUpdater()创建更新器,并设置想要更新的类和属性
(2)想要更新的类的属性必须使用public volatile修饰
Java中的并发工具类
等待多线程完成的CountDownLatch
CountDownLatch允许一个或多个线程等待其他线程完成操作
CountDownLatch countDownLatch = new CountDownLatch(N); 接收一个int类型的参数N作为计数器,想要等待N个点完成,就传入N
public void countDown():调用一次此方法N就会减一
public void await():调用此方法会阻塞当前线程,直到N变成零
public boolean await(long timeout, TimeUnit unit):等待指定的时间后,不再阻塞当前线程,返回true表示等待条件到达,false表示条件未到达,但超时了
public long getCount():获取当前计数值
由于countDown()方法可以用在任何地方,所以这里的N个点,可以是N个线程,也可以是一个线程里的N个执行步骤
用在多个线程时,只需要把这个CountDownLatch的引用传递到线程里即可
与join()的区别
如果调用某个thread的join()方法,那么当前线程就会被阻塞,直到thread线程执行完毕,当前线程才能继续向下执行
join()内部不断判断thread线程是否存活,如果存活则让当前线程一直wait(0),直到thread线程终止,就会调用线程的this.notifyAll()方法
CountDownLatch只需要检查计数器的值为零就可以继续向下执行
注意:计数器必须大于等于零,等于零时,调用await()方法不会阻塞当前线程。CountDownLatch不可能重新初始化或修改对象内部计数器的值。一个线程调用countDown()方法happens-before,另一个线程调用await()方法
同步屏障CyclicBarrier
它要做的事情是,让一组线程到达一个屏障(也可以叫做同步点)时被阻塞,直到最后一个线程到达屏障时,屏障才会打开,所有被屏障拦截的线程才会继续执行
public CyclicBarrier(int parties):parties为屏障拦截的线程数量
public CyclicBarrier(int parties, Runnable barrierAction):在线程到达屏障时,优先执行barrierAction
public int await():每个线程调用CyclicBarrier对象的await()方法告诉CyclicBarrier我已到达屏障,然后当前线程被阻塞
应用场景:CyclicBarrier可用于多线程计算数据,最后合并计算结果的场景
CyclicBarrier和CountDownLatch的区别
CountDownLatch的计数器只能使用过一次,而CyclicBarrier的计数器可以使用reset()方法重置
CyclicBarrier还提供了一些其他方法
public int getNumberWaiting():获取CyclicBarrier阻塞的线程数量
public boolean isBroken():判断CyclicBarrier阻塞的线程是否被中断,被中断则返回true,否则返回false
控制并发线程数的Semaphore
Semaphore(信号量)用来控制同时访问特定资源的线程数量,他通过协调各个线程,以保证合理的使用公共资源
可以将其比喻为控制流量的红绿灯,比如XXX马路要限制流量,要求一次只能通过100辆车,其他的车需要在路口等待,那么前100辆车看到的便是绿灯(acquire()),其他的车看到的是红灯,如果这100辆车中有5辆车已经离开(release()),那么就允许后面5辆驶入马路。例子中车就是线程,驶入马路就是正在执行,离开便是执行完成,看到后灯表示被阻塞,不能执行。
应用场景:Semaphore可以用于做流量限制,特别是公共资源有限的场景,比如数据库连接等
public Semaphore(int permits):permits表示许可证的数量,即最大并发数
public void acquire():从信号量获取一个许可,如果无可用许可前将一直阻塞等待,使用完后调用release()归还许可证
public void acquire(int permits):获取指定数目的许可,如果无可用许可前也将会一直阻塞等待
public boolean tryAcquire():从信号量尝试获取一个许可,如果无可用许可,直接返回false,不会阻塞
public boolean tryAcquire(int permits):尝试获取指定数目的许可,如果无可用许可直接返回false
public boolean tryAcquire(int permits, long timeout, TimeUnit unit):在指定的时间内尝试从信号量中获取许可,如果在指定的时间内获取成功,返回true,否则返回false
public void release():释放一个许可,别忘了在finally中使用,注意:多次调用该方法,会使信号量的许可数增加,达到动态扩展的效果,如:初始permits为1, 调用了两次release,最大许可会改变为2
public int availablePermits():获取当前信号量可用的许可数量
public final int getQueueLength():返回正在等待获取许可证的线程数
public final boolean hasQueuedThreads():是否有线程正在等待获取许可证
线程间交换数据的Exchanger
Exchanger(交换者)可以用于线程间的数据交换,它提供一个同步点,在这个同步点,两个线程通过调用exchange()方法交换彼此的数据。如果第一个线程先执行exchange(),它会一直等待第二个线程执行exchange()方法,当两个线程都到达同步点时,便可以交换彼此的数据
public V exchange(V x):等待另一个线程到达此交换点(除非当前线程被中断),然后将给定的对象传送给该线程,并接收该线程的对象。
public V exchange(V x, long timeout, TimeUnit unit):等待另一个线程到达此交换点(除非当前线程被中断或超出了指定的等待时间),然后将给定的对象传送给该线程,并接收该线程的对象。
Java并发编程实践
Executor框架
Executor简介
两级调度模型
在HotSpot VM中,将Java线程(java.lang.Thread)一对一映射为操作系统的线程,在上层Java多线程程序将应用拆分为若干个任务,使用用户级调度器(Executor框架)将其映射为固定数量的线程,在底层操作系统内核将这些线程映射到硬件处理器上
Executor的结构和成员
结构
任务
包括被执行的任务实现的接口Runnable和Callable,实现这两个接口的类可以被Executor执行
任务的执行
核心接口Executor,将任务的提交与执行分离
子接口
ExecutorService
实现类
ThreadPoolExecutor
ScheduledThreadPoolExecutor
异步计算的结果
接口Future,代表异步任务的执行结果
实现类
FutureTask
Executors.callable(Runnable task)方法可以将Runnable转为Callable,Runnable对象可以交给ExecutorService.execute(Runnable command)执行,或可以把Runnable或Callable对象交给ExecutorService.submit(Runnable task)或submit(Callable<T> task)执行。submit()方法回返回Future接口的实现类FutureTask类的对象,此对象也实现了Runnable接口,调用其get()方法等待任务执行完成,调用cancel(boolean mayInterruptIfRunning)取消此任务的执行
成员
ThreadPoolExecutor
ThreadPoolExecutor通常使用工厂类Executors来创建
FixedThreadPool,固定线程数的线程池,适用于需要限制线程数量的场景
public static ExecutorService newFixedThreadPool(int nThreads)
public static ExecutorService newFixedThreadPool(int nThreads, ThreadFactory threadFactory)
SingleThreadExecutor,使用单个线程的线程池,适用于需要保证顺序的执行各个任务;并且在任意时间点不会有多个线程活动的场景
public static ExecutorService newSingleThreadExecutor()
public static ExecutorService newSingleThreadExecutor(ThreadFactory threadFactory)
CachedThreadPool,是一个大小无界的线程池,会根据需要创建新线程,适用于执行很多短期异步任务的小程序
public static ExecutorService newCachedThreadPool()
public static ExecutorService newCachedThreadPool(ThreadFactory threadFactory)
ScheduledThreadPoolExecutor
ScheduledThreadPoolExecutor通常使用Executors来创建
ScheduledThreadPoolExecutor,适用于需要多个后台线程执行周期任务,但又需要限制后台线程数量的场景
public static ScheduledExecutorService newScheduledThreadPool(int corePoolSize)
public static ScheduledExecutorService newScheduledThreadPool(int corePoolSize, ThreadFactory threadFactory)
SingleThreadScheduledExecutor,适用于需要单个线程后台执行周期任务,并且需要保证顺序的执行各个任务的场景
public static ScheduledExecutorService newSingleThreadScheduledExecutor()
public static ScheduledExecutorService newSingleThreadScheduledExecutor(ThreadFactory threadFactory)
Future
Future接口表示异步计算的结果,FutureTask是其实现类,当把Runnable或Callable实现类的对象提交(submit)给Executor时,将返回Future接口实现类的对象
ExecutorService中返回Future的方法
<T> Future<T> submit(Callable<T> task)
<T> Future<T> submit(Runnable task, T result)
Future<?> submit(Runnable task)
V get():获取任务的结果,如果任务还未执行完,则阻塞当前线程
Runnable和Callable
Runnable接口和Callable接口的实现类对象可以被Executor执行,Runnable无返回结果,Callable有返回结果
Runnable转Callable
使用Executors类的静态方法
public static Callable<Object> callable(Runnable task)
public static <T> Callable<T> callable(Runnable task, T result):result指定调用future.get()时的返回值
ThreadPoolExecutor详解
ThreadPoolExecutor是Executor框架的核心类,它是线程池的实现
corePoolSize:核心线程池线程数
maximumPoolSize:线程池最大线程数
blockingQueue:暂时保存任务的工作队列
RejectedExecutionHandler:当线程池关闭或已饱和(当前线程数大于maximumPoolSize且工作队列已满)时采取的饱和策略
通过工具类Executors创建的三种ThreadPoolExecutor
FixedThreadPool详解
FixedThreadPool被称为可重用的固定线程数的线程池
corePoolSize和maximumPoolSize都被设置为创建ThreadPoolExecutor时指定的线程数
keepAliveTime参数指定为0L,则空闲的线程会立即被终止
FixedThreadPool的execute()方法
(1)如果当前线程数小于corePoolSize,则创建新的线程执行任务
(2)如果当前线程数等于corePoolSize,则将任务放到LinkedBlockingQueue中等待
(3)当线程执行完(1)中的任务后,会循环从LinkedBlockingQueue中获取任务执行
(4)FixedThreadPool使用无界队列LinkedBlockingQueue(队列容量为Integer.MAX_VALUE)作为工作队列
使用无界工作队列有什么影响?
(1)由于当前线程数等于corePoolSize时会把任务放在工作队列等待,而此队列是个无界队列,所以将导致maximumPoolSize参数无效
(2)由于(1),会导致keepAliveTime参数无效
(3)由于使用无界队列,运行中的FixedThreadPool(没有调用shutdown()或shutDownNow())不会拒绝新的任务,即RejectedExecutionHandler.rejectedExecution()方法不会被调用
SingleThreadExecutor详解
SingleThreadExecutor是一个只使用单个worker线程的线程池
corePoolSize和maximumPoolSize被设置为1,其他参数和FixedThreadPool相同,其工作队列也是LinkedBlockingQueue,带来的影响与FixedThreadPool相同
SingleThreadExecutor的execute()方法
(1)当线程池中的线程数小于corePoolSize(无运行的线程),则创建新线程去执行任务
(2)当线程池中有一个运行的线程时,则将任务放到LinkedBlockingQueue中
(3)当线程执行完(1)中的任务,则循环从LinkedBlockingQueue中获取任务执行
CachedThreadPool详解
CachedThreadPool是一个根据需要创建线程的线程池
corePoolSize为0,即corePool为空,maximumPoolSize为Integer.MAX_VALUE,即maximumPool是无界的。keepAliveTime为60L,即空闲线程最多等待新任务60秒,超时则终止空闲线程
CachedThreadPool使用没有容量的SynchronousQueue作为工作队列,意味着当主线程提交任务的速度快于maximumPool中线程处理任务的速度时,CachedThreadPool会不断创建新的线程,极端情况下会导致CPU和内存因线程过多而耗尽
CachedThreadPool的execute()方法
(1)首先执行SynchronousQueue.offer(Runnable task)方法,如果maximumPool线程池中有空闲线程在执行SynchronousQueue.poll(keepAliveTime,TimeUnit.NANOSECONDS),则主线程将任务交给该空闲线程执行,execute()方法执行完成;否则执行步骤(2)
(2)如果maximumPool中为空或没有空闲的线程(没有线程执行SynchronousQueue.poll(keepAliveTime,TimeUnit.NANOSECONDS)操作),则步骤(1)失败,CachedThreadPool创建一个新的线程去执行任务。execute()方法执行完成
(3)当步骤(2)执行完成,SynchronousQueue.poll(keepAliveTime,TimeUnit.NANOSECONDS)方法会使空闲的线程等待60秒,如果60秒内有任务则执行,超时则终止空闲的线程
由于CachedThreadPool中的线程超时会终止,所以长时间保持空闲的CachedThreadPool并不会占用任何资源
SynchronousQueue是一个没有容量的阻塞队列,每个插入操作必须等待一个移除操作,反之亦然。CachedThreadPool使用SynchronousQueue将主线程提交的任务传递给空闲线程执行
ScheduledThreadPoolExecutor详解
ScheduledThreadExecutor继承了ThreadPoolExecutor,用于在给定的延迟后执行任务或定期执行任务
运行机制
(1)当调用scheduleAtFixedRate()或scheduleWithFixedDelay()时,会向ScheduleThreadExecutor的DelayQueue中添加一条实现了RunnableScheduledFuture的ScheduledFutureTask
(2)线程池中的线程从DelayQueue中获取ScheduledFutureTask,然后执行
DelayQueue是一个无界队列,所以设置maximumPoolSize无效
ScheduledThreadPoolExecutor的实现
ScheduledThreadPoolExecutor会把待调度的任务(ScheduledFutureTask)放在DelayQueue中
ScheduledFutureTask主要包括三个成员变量
long类型的time,表示任务的执行时间
long类型的sequenceNumber,表示任务被添到线程池中的序号
long类型的period,表示执行任务的间隔周期
DelayQueue封装了一个PriorityQueue,会对ScheduledFutureTask按time进行排序,time小的排在前面(时间早的先执行),time相同则按照sequenceNumber排序,sequenceNumber小的在前(先添加的先执行)
ScheduledThreadPoolExecutor中的线程执行周期任务的过程
(1)首先从DelayQueue中获取到期的ScheduledFutureTask(DelayQueue.take()方法)。到期是指time大于等于当前时间
获取任务的过程
(1)获取Lock
(2)循环获取周期任务
如果PriorityQueue为空,则进入Condition等待
如果队头元素的time比当前时间大,则等待到time时间
获取队头元素,如果PriorityQueue不为空,则唤醒所有在Condition中等待的线程
(3)释放Lock
(2)执行ScheduledFutureTask
(3)修改ScheduledFutureTask的time属性为下次被执行时间
(4)将修改time属性之后的ScheduledFutureTask放回DelayQueue中(DelayQueue.offer()方法)
将任务放回队列的过程
(1)获取Lock
(2)添加任务
向PriorityQueue中添加任务
如果添加的任务是PriorityQueue的队头元素,则唤醒所有在Condition中等待的线程
(3)释放Lock
FutureTask详解
Future接口和实现了Future接口的FutureTask类代表异步计算的结果
FutureTask简介
FutureTask也实现了Runnable接口,所以可以交给Executor执行或由线程直接调用(FutureTask.run())
状态
未启动,创建完成且尚未执行FutureTask.run()方法
已启动,正在执行FutureTask.run()方法
此时调用FutureTask.get()会阻塞调用的线程
已完成,FutureTask.run()正常结束,或被取消(FutureTask.cancel()),或执行时抛出异常而结束
此时调用FutureTask.get()会立即返回结果或抛出异常
public boolean cancel(boolean mayInterruptIfRunning):取消子任务的执行
子任务已经结束、被取消、不能取消,则此方法执行失败返回false
子任务未启动,则会取消子任务,不再执行
子任务已启动,正在执行,如果mayInterruptIfRunning为true,则中断执行子任务的线程,并返回true,如果mayInterruptIfRunning为false,则不中断执行子任务的线程,并返回true
FutureTask的使用
可以使用ExecutorService.submit()方法返回一个FutureTask或单独使用FutureTask
当一个线程需要等待另一个线程执行完毕才能继续向下执行时,可以使用FutureTask
FutureTask原理
Java8中修改了FutureTask的实现,参考链接:https://blog.csdn.net/zhuguang10/article/details/97374772
Java中的线程池
使用线程池的好处
(1)减低资源消耗,重复利用已创建的线程,降低创建和销毁线程时的消耗
(2)提高响应速度,任务到来时,不需要等待线程的创建便可以直接执行
(3)提高线程的可管理性,线程属于稀缺资源,使用线程池可以统一的分配、调优和监控
线程池的实现原理
主要流程
ThreadPoolExecutor执行execute()方法的4种情况
如果当前运行的线程少于corePoolSize,则创建新的线程来执行任务(执行这一步需要获取全局锁)
如果运行的线程大于等于corePoolSize,则将任务加入BlockingQueue
如果无法将任务加入BlockingQueue(队列已满),并且线程数小于maximumPoolSize,则创建新的线程来执行任务(执行这一步需要获取全局锁)
如果创建新的线程会使当前运行的线程数超出maximumPoolSize,则任务将被拒绝,并调用RejectedExecutionHandler.rejectedExecution()方法
采用这个设计思路是为了避免获取全局锁
源码分析
工作线程:线程池创建线程时,会将线程封装为工作线程Worker,Worker在执行完任务后,还会循环从工作队列中获取任务来执行
线程池的使用
创建线程池
public ThreadPoolExecutor(int corePoolSize,int maximumPoolSize, long keepAliveTime,TimeUnit unit, BlockingQueue<Runnable> workQueue, ThreadFactory threadFactory, RejectedExecutionHandler handler)
corePoolSize
线程池基本大小,当提交一个任务到线程池,即使其他空闲的基本线程可以执行任务,也会创建新的线程,直到执行任务的线程数大于corePoolSize时就不再创建新的线程。调用prestartAllCoreThreads()方法线程池会提前创建并启动所有基本线程
maximumPoolSize
线程池中允许的最大线程数,如果任务队列已满,并且当前线程数小于maximumPoolSize,则创建新的线程来执行任务。注意:如果使用了无界的任务队列,则此参数无效
keepAliveTime
线程最大空闲时间
unit
线程最大空闲时间的单位
workQueue
用于保存等待执行的任务的阻塞队列
可以选择以下几个阻塞队列
ArrayBlockingQueue
LinkedBlockingQueue,吞吐量高于ArrayBlockingQueue,Executors.newFixedThreadPool()使用了此队列
SynchronousQueue,吞吐量高于LinkedBlockingQueue,Executors.newCachedThreadPool()使用了此队列
PriorityBlockingQueue
threadFactory
创建线程的工厂,可以通过线程工厂给创建出的线程设置更有意义的名字
handler
饱和策略,当线程池和队列都满了,指定如何处理新提交的任务的策略
AbortPolicy:直接抛出异常
CallerRunsPolicy:只用调用者所在的线程来运行任务
DiscardOldestPolicy:丢弃队列里最近的一个任务,并执行当前任务
DiscardPolicy:不处理,丢弃掉
也可以实现RejectedExecutionHandler接口自定义策略
向线程池提交任务
void execute(Runnable command):用于提交不需要返回值的任务,所以无法判断任务是否被线程池执行成功
Future<?> submit(Runnable task):用于提交需要返回值的任务
可以使用future.get()方法获取返回值,此方法会阻塞当前线程,直到任务完成,可以使用get(long timeout, TimeUnit unit)方法,阻塞一段时间后返回,此时任务可能还没有执行完
关闭线程池
void shutdown():将线程池的状态设置为SHUTDOWN,并中断所有没有正在执行任务的线程
List<Runnable> shutdownNow():将线程池状态设置为STOP,然后尝试停止所有正在执行任务或暂停执行任务的线程,并返回待执行任务列表
原理是遍历线程池中的线程,逐个调用线程的interrupt()方法,所以无法响应中断的任务可能永远无法终止
调用两个关闭方法中的任意一个,isShutdown()方法就会返回true,当所有任务都关闭后调用isTerminaed()方法才返回true
合理地配置线程池
可以从以下几个角度分许线程池中任务的特性
任务的性质:CPU密集型、I/O密集型和混合型
CPU密集型任务应配置尽可能少的线程,如Ncpu + 1个线程的线程池。
I/O密集型任务并不是一直在执行计算,应配置尽可能多的线程,如2 * NCpu
混合型任务如果执行时间相差不是太大可以将其拆分为CPU密集型任务和I/O密集型任务,否则没有必要拆分
Runtime.getRuntime().availableProcessors()方法可以获取当前设备CPU的个数
任务的优先级:高、中和低
优先级不同的任务可以使用PriorityBlockingQueue来处理
任务的执行时间:长、中和短
执行时间不同的任务可以交给不同规模的线程池来处理或使用优先队列让执行时间短的任务先执行
任务的依赖性:是否依赖其他系统资源,比如数据库连接
依赖数据库连接池的任务,因为线程需要等待数据库返回结果,等待时间越长,CPU空闲时间越长,所以线程数应该设置的大一些,更好的利用CPU
推荐使用有界阻塞队列,增加系统的稳定性和预警能力,可以根据需求设置的大一些
线程池的监控
使用线程池提供的参数对线程池进行监控,方便在出现问题时,快速定位问题
线程池的属性
taskCount:任务数量
completedTaskCount:已完成的任务数量,小于等于taskCount
largestPoolSize:曾经创建过的最大线程数量,该值如果等于最大线程数,说明线程池曾经满过
getPoolSize:线程池的线程数量,如果线程池不销毁,里面的线程不会自动销毁,所以此值只增不减
getActiveCount:获取活动的线程数
可以通过继承线程池重写指定的方法,对线程池进行监控
beforeExecute():任务执行前
afterExecute():任务执行后
terminated():线程池关闭前