学习笔记-2 多线程

2.并发编程

2.1 线程的基本认识

2.1.1 线程的基本介绍

进程的执行逻辑

什么是线程:线程是操作系统能够进行运算调度的最小单位。它被包含在进程之中,是进程中的实际运作单位。
什么是进程:进程是一个具有一定独立功能的程序关于某个数据集合的一次运行活动。它是操作系统动态执行的基本单元,在传统的操作系统中,进程既是基本的[分配单元],也是基本的执行单元

为什么会有线程:
1.在多核CPU中,利用多线程可以实现真正意义上的并行执行。
2.在一个应用进程中,会存在多个同时执行的任务,如果其中一个任务被阻塞,将会引起不依赖该任务的任务也被阻塞。通过对不同任务创建不同的线程去处理,可以提升程序处理的实时性。
3.线程可以认为是轻量级的进程,所以线程的创建、销毁比进程更快。

2.1.2 线程的应用场景

为什么要用多线程

  1. 异步执行
  2. 利用多CPU资源实现真正意义上的并行执行。

线程的价值:


线程的价值

线程的应用场景

  1. 使用多线程实现文件下载。
  2. 后台任务:如定时向大量(100W以上)的用户发送邮件。
  3. 异步处理:记录日志。
  4. 多步骤的任务处理,可根据步骤特征选用不同个数和特征的线程来协作处理,多任务的分割,由一个主线程分割给多个线程完成。

总结
多线程的本质是:合理的利用多核心CPU资源来实现线程的并行处理,来实现同一个进程内的多个任务的并行执行,同时基于线程本身的异步执行特性,提升任务处理的效率。

2.1.3 如何在 Java 中应用多线程

Java中使用多线程的方式:

  • 继承Thread类
public class ThreadDemo extends Thread{

    @Override
    public void run() {
        System.out.println("当前线程:"+Thread.currentThread().getName());
    }

    public static void main(String[] args) {
        ThreadDemo threadDemo=new ThreadDemo();
        //threadDemo.start(); //启动一个线程
    }
}
  • 实现Runnable接口
public class RunnableDemo implements Runnable{
    @Override
    public void run() {
        System.out.println("当前线程:"+Thread.currentThread().getName());
    }
    public static void main(String[] args) {
        RunnableDemo runnableDemo=new RunnableDemo();
        Thread thread=new Thread(runnableDemo);
        thread.start();//启动线程
    }
}
  • 实现Callable接口
public class CallableDemo implements Callable<String> {
    @Override
    public String call() throws Exception {
        System.out.println("当前线程:"+Thread.currentThread().getName());
        Thread.sleep(10000);
        return "Hello";
    }

    public static void main(String[] args) throws ExecutionException, InterruptedException {
        ExecutorService executorService= Executors.newFixedThreadPool(1);
        Future<String> future=executorService.submit(new CallableDemo());
        //future.get 是一个阻塞方法
        System.out.println(Thread.currentThread().getName()+"-"+future.get());
    }
}

2.1.4 Java 线程的生命周期

Java 线程从创建到销毁,一共有6 个状态:

  • NEW:初始状态,线程被构建,但是还没有调用start方法。
  • RUNNABLED:运行状态,JAVA线程把操作系统中的就绪和运行两种状态统一称为“运行中”。
  • BLOCKED:阻塞状态,表示线程进入等待状态,也就是线程因为某种原因放弃了CPU使用权,阻塞也分为几种情况。
  • WAITING: 等待状态
  • TIME_WAITING:超时等待状态,超时以后自动返回
  • TERMINATED:终止状态,表示当前线程执行完毕
线程的生命周期
public class ThreadStatusDemo {

    public static void main(String[] args) {
        //TIME_WAITING
        new Thread(()->{
            while(true){
                try {
                    TimeUnit.SECONDS.sleep(100); // 超时等待
                } catch (InterruptedException e) {
                    e.printStackTrace();
                }
            }
        },"Time_Wating_Demo").start();
        //WAITING
        new Thread(()->{
            while(true){
                synchronized (ThreadStatusDemo.class){
                    try {
                        ThreadStatusDemo.class.wait(); //等待
                    } catch (InterruptedException e) {
                        e.printStackTrace();
                    }
                }
            }
        },"Wating").start();

        // 优先抢占锁的进入time_waiting状态, 没抢到锁的进入block状态
        new Thread(new BlockedDemo(),"Blocked-Demo-01").start();
        new Thread(new BlockedDemo(),"Blocked-Demo-02").start();
    }
    
    static class BlockedDemo extends  Thread{
        @Override
        public void run(){
            synchronized (BlockedDemo.class){
                while(true){ //死循环
                    try {
                        TimeUnit.SECONDS.sleep(100);
                    } catch (InterruptedException e) {
                        e.printStackTrace();
                    }
                }
            }
        }
    }
}

2.2 线程的基本操作及原理

2.2.1 Thread.join的使用及原理

public class ThreadJonDemo {
    private static int x=0;
    private static int i=0;
    public static void main(String[] args) throws InterruptedException {
        Thread t1=new Thread(()->{
            //阻塞操作
            i=1;
            x=2;
        });
        Thread t2=new Thread(()->{
            i=x+2;
        });
        //两个线程的执行顺序,
        t1.start();
        t1.join(); //t1线程的执行结果对于t2可见(t1线程一定要比t2线程有限执行) --- 阻塞
        t2.start();
        Thread.sleep(1000);
        System.out.println("result:"+i);

    }

Thread.join的作用是保证线程执行结果的可见性。
可见性:可以理解为对结果对后续代码可见,即意味着,线程执行完成以后,后续代码才会开始执行。

Thread.join的原理

Thread.join的本质其实是wait/notifyall。它阻塞的被调用时, join所在的线程,在上述图片示例中为主线程。

2.2.2 Thread.sleep

Thread.sleep: 使线程暂停执行一段时间,直到等待的时间结束才恢复执行或在这段时间内被中断。
Thread.sleep的工作流程:

  • 挂起线程并修改其运行状态
  • 用sleep()提供的参数来设置一个定时器。
  • 当时间结束,定时器会触发,内核收到中断后修改线程的运行状态。
    例如线程会被标志为就绪而进入就绪队列等待调度。
    Thread.Sleep(0) 的意义: 类似Thred.yield.使线程放弃CPU的使用并重新加入竞争中。
    操作系统中,CPU竞争有很多种策略。Unix系统使用的是时间片算法,而Windows则属于抢占式的。

2.2.3 wait 和 notify的使用

一个线程修改了一个对象的值,而另个线程感知到了变化,然后进行响应的操作。
例如可以通过waitnotify来实现阻塞队列。

// 生产者
public class Producer implements Runnable{

    private Queue<String> bags;
    private int size;

    public Producer(Queue<String> bags, int size) {
        this.bags = bags;
        this.size = size;
    }


    @Override
    public void run() {
        int i=0;
        while(true){
            i++;
            synchronized (bags){
                while(bags.size()==size){
                    System.out.println("bags已经满了");
                    //TODO? 阻塞?
                    try {
                        bags.wait();
                    } catch (InterruptedException e) {
                        e.printStackTrace();
                    }
                }
                try {
                    Thread.sleep(1000);
                } catch (InterruptedException e) {
                    e.printStackTrace();
                }
                System.out.println("生产者-生产:bag"+i);
                bags.add("bag"+i);
                //TODO? 唤醒处于阻塞状态下的消费者
                bags.notifyAll();
            }
        }
    }
}
// 生产者
public class Consumer implements Runnable{
    private Queue<String> bags;
    private int size;

    public Consumer(Queue<String> bags, int size) {
        this.bags = bags;
        this.size = size;
    }


    @Override
    public void run() {
        while(true){
            synchronized (bags){
                while(bags.isEmpty()){
                    System.out.println("bags为空");
                    try {
                        bags.wait();
                    } catch (InterruptedException e) {
                        e.printStackTrace();
                    }
                }
                try {
                    Thread.sleep(1000);
                } catch (InterruptedException e) {
                    e.printStackTrace();
                }
                String bag=bags.remove();
                System.out.println("消费者消费:"+bag);
                bags.notifyAll();
            }
        }
    }
}
// 主程序入口
public class WaitNotifyDemo {

    public static void main(String[] args) {
        // 队列为共享资源
        Queue<String> queue=new LinkedList<>();
        int size=10;
        Producer producer=new Producer(queue,size);
        Consumer consumer=new Consumer(queue,size);
        Thread t1=new Thread(producer);
        Thread t2=new Thread(consumer);
        t1.start();
        t2.start();
    }
}

为什么wait/notify需要加synchronized?

  1. 从刚刚的案例来看,其实wait/notify本质上其实是一种条件的竞争,至少来说,
    wait和notify方法一定是互斥存在的,既然要实现互斥,那么synchronized就是一
    个很好的解决方法。
  2. wait和notify是用于实现多个线程之间的通信,而通信必然会存在一个通信的
    载体,比如我们小时候玩的游戏,用两个纸杯,中间用一根线连接,然后可以实现
    比较远距离的对话。而这根线就是载体,那么在这里也是一样,wait/notify是基
    于synchronized来实现通信的。也就是两者必须要在同一个频道也就是同一个锁的
    范围内。

2.3.4 Thread.interrupt和 Thread.interrupted

interrupt是类种的一个属性,在JVM中维护。
interrupt(): 中断一个线程。当其他线程通过调用当前线程的interrupt方法,表示向当前线程打个招呼,告诉他可以中断线程的执行了,至于什么时候中断,取决于当前线程自己。
interrupted:对设置中断标识的线程复位,并且返回当前的中断状态。

public static void main(String[] args) throws InterruptedException {

        Thread thread=new Thread(()->{
            int i = 0;
// isInterrupted只有在循环内才有效
            while(!Thread.currentThread().isInterrupted()){
                  i++;
                }
            }
        });
        thread.start();
        TimeUnit.SECONDS.sleep(1);
        thread.interrupt(); //中断
    }

线程处于阻塞状态下的情况下(中断才有意义)

  • thread.join
  • wait
  • Thread.sleep
public static void main(String[] args) throws InterruptedException {
        Thread thread=new Thread(()->{
            //Thread.currentThread().isInterrupted() 默认是false
           //正常的任务处理..
            try {
                // 线程阻塞时,调用interrupt()方法可以使阻塞的线程立马跑出中断异常,我们可以捕获中断异常来进行后续的逻辑处理。
                Thread.sleep(10000);
            } catch (InterruptedException e) {
                //抛出异常来相应客户端的中断请求
                e.printStackTrace();
            }
        });

        thread.start();
        Thread.sleep(5000);
        //interrupt 这个属性由false-true
        thread.interrupt(); //中断(友好)
    }

2.3 线程的安全分析

2.3.1 并发编程问题的源头- - 原子性、可见性、有序性

如何理解线程安全:当多个线程访问某个对象时,不管运行时环境采用何种调度方式或者这些线程将如何交替执行,并且在主调代码中不需要任何额外的同步或者协同,这个类都能表现出正确的行为,那么就称这个类是线程安全的。

线程问题的源头

cpu内存模型

CPU为了提高处理速度,减少IO瓶颈的优化

  • CPU增加了高速缓存,均衡与内存的速度差异
  • 操作系统增加进程、线程、以及分时复用cpu,均衡cpu与i/o设备的速度差异
  • 编译程序优化指令的执行顺序,使得能够更加合理的利用缓存

可见性问题的源头

可见性问题的源头

线程安全问题的本质:

  • 原子性
  • 可见性
  • 有序性
    CPU 高速缓存和内存的架构设计,以及指令重排序问题

2.3.2 Java 内存模型 -Java 如何解决可见性有序性问题

JVM: Java内存模型是一种抽象结构,它提供了合理的禁用缓存以及禁止重排序
的方法来解决可见性、有序性问题。

JMM和硬件模型的对应简图


image.png

2.3.3 同步关键字Synchronized

synchronized锁的范围

  • 对于普通同步方法,锁是当前实例对象。
  • 对于静态同步方法,锁是当前类的Class对象。
  • 对于同步方法块,锁是Synchonized括号里配置的对象
public class SyncDemo {

    //对象锁(同一个对象有效)  this
    public synchronized void demo(){

    }

    public void demo1(){
        //TODO
        synchronized (this){  //对象锁
        }
        //TODO
    }

    //类级别的锁  SyncDemo.class
    public synchronized static void demo3(){

    }

    public void demo4(){
        //TODO
        // 方法快的锁
        synchronized (SyncDemo.class){

        }
    }

    public static void main(String[] args) {
        SyncDemo syncDemo=new SyncDemo();
        SyncDemo syncDemo1=new SyncDemo();
        //无法实现两个线程的互斥.
        //如果访问demo3的话,那么下面两个线程会存在互斥
        new Thread(()->{  //syncDemo这个实例
            syncDemo.demo1();
        }).start();

        new Thread(()->{ //BLOCKED状态
            syncDemo1.demo1();//syncDemo1这个实例
        }).start();
    }

}

synchronized 在指令中会增加monitorenter指令

synchronized的本质

2.3.4 volatile关键字分析

volatile是什么?
volatile可以用来解决可见性和有序性问题

volatile会在指令中增加Lock指令

Lock指令的作用

  • 将当前处理器缓存行的数据写回到系统内存。
  • 这个写回内存的操作会使在其他CPU里缓存了该内存地址的数据无效。

什么情况下需要用到volatile
当存在多个线程对同一个共享变量进行修改的时候,需要增加volatile,保证数据修改的实时可见。

Volatile是如何解决有序性问题的
下图中,进入if后,value可能不等于10,这就是有序性问题。

可能产生的问题

CPU层面的内存屏障

  • Store Barrier:强制所有在store屏障指令之前的store指令,都在该store屏障指令执行之前被执行,并把store缓冲区的数据都刷到CPU缓存;
  • Load Barrier:强制所有在load屏障指令之后的load指令,都在该load屏障指令执行之后被执行,并且一直等到load缓冲区被该CPU读完才能执行之后的load指令;
  • Full Barrier:复合了load和storee屏障的功能。

CPU层面的内存屏障

CPU层面的内存屏障

JVM提供的内存屏障


JVM提供的内存屏障

volatile的总结
本质上来说:volatile实际上是通过内存屏障来防止指令重排序以及禁止cpu高速缓存来解决可见性问题。
Lock指令,
它本意上是禁止高速缓存解决可见性问题,但实际上在这里,它表示的是一种内存屏障的功能。也就是说针对当前的硬件环境,JMM层面采用Lock指令作为内存屏障来解决可见性问题.

2.3.5 final域

final关键字
final在Java中是一个保留的关键字,可以声明成员变量、方法、类以及本地变量。一旦你将引用声明作final,你将不能改变这个引用了。

final域和线程安全有什么关系?
对于 final域,编译器和处理器要遵守两个重排序规则:

  • 在构造函数内对一个final域的写入,与随后把这个被构造对象的引用赋值给一
    个引用变量,这两个操作之间不能重排序。
  • 初次读一个包含final域的对象的引用,与随后初次读这个final域,这两个操
    作之间不能重排序。
final代码

写可能执行的情况

可能执行的情况

写final域重排序规则

  • JMM禁止编译器把final域的写重排序到构造函数之外。
  • 编译器会在final域的写之后,构造函数return之前,插入一个StoreStore屏障。这个屏障禁止处理器把final域的写重排序到构造函数之外。

读的时候可能存在的执行情况

读的时候可能存在的执行情况

读域的重排序规则
在一个线程中,初次读对象引用与初次读该对象包含的final域,JMM禁止处理器重排序这两个操作,编译器会在读final域操作的前面插入一个LoadLoad屏障。

构造器内部造成逃逸

逃逸代码

溢出带来的重排序问题
溢出带来的重排序问题

2.3.6 Happens-Before 规则

什么是Happens-Before?
Happens-Before是一种可见性规则,它表达的含义是前面一个操作的结果对后续操作是可见的。

6种Happens-Before规则:

  • 程序顺序规则
  • 监视器锁规则:对一个锁的解锁 Happens-Before 于后续对这个锁的加锁。


    监视器锁规则
  • Volatile变量规则:对一个volatile域的写,happens-before于任意后续对这个volatile域的读.

  • 传递性: 如果A happens-before B,且B happens-before C,那么A happens-before C.

  • start()规则: 如果线程A执行操作ThreadB.start()(启动线程B),那么A线程的ThreadB.start()操作happens-before于线程B中的任意操作.

  • Join()规则: 如果线程A执行操作ThreadB.join()并成功返回,那么线程B中的任意操作happens-before于线程A从ThreadB.join()操作成功返回.

2.3.7 原子类 Atomic-无锁工具的典范

原子性问题的解决方案:

  • synchronizedLock
    J.U.C包下的Atomic
    Atomic类的示例:
public class AtomicDemo {

  //  public static int count=0;
    private static AtomicInteger atomicInteger=new AtomicInteger(0);
    public static void incr(){
        try {
            Thread.sleep(1);
        } catch (InterruptedException e) {
            e.printStackTrace();
        }
        //count++; //count++ (只会由一个线程来执行)
        atomicInteger.incrementAndGet();
    }
    public static void main(String[] args) throws InterruptedException {
        TimeUnit.SECONDS.sleep(1);
        for (int i = 0; i < 1000; i++) {
            new Thread(AtomicDemo::incr).start();
        }
        Thread.sleep(2000);
        System.out.println("result:"+count);

        //System.out.println("result:"+atomicInteger.get());
    }
}

Atomic实现原理:

  • Unsafe类;
  • CAS:


    image.png

2.3.8 ThreadLocal的实现原理

ThreadLocal实现线程隔离

public class ThreadLocalDemo {

    private static Integer num = 0;

    private static final ThreadLocal<Integer> local = new ThreadLocal<Integer>() {

        // 重写设置初始值的方法 
        @Override
        protected Integer initialValue() {
            return 0;
        }
    };

    public static void main(String[] args) {
        Thread[] threads = new Thread[5];
        for (int i = 0; i < 5; i++) {
            threads[i] = new Thread(() -> {
                // 希望所有线程拿到初始值0
                int num = local.get();
                num += 5;
                local.set(local.get() + 5);
                System.out.println(Thread.currentThread().getName() + "---" + Integer.toString(local.get()));
            }, "thread" + i);
        }

        for (Thread thread : threads) {
            thread.start();
        }
    }
} 

2.4.1 发布与逃逸

  • 发布

    发布的意思是使一个对象能够被当前范围之外的代码所使用.

public static HashSet<Person> persons;
public void init() {
    person = new HashSet<Person>();
}

不安全发布

不安全发布是指某个不应该发布的对象被发布了.

下面案例中, states是个私有变量,但是因为不安全发布使得其数值可以被外部对象修改,导致不安全发布。


网页捕获_1-3-2022_211432_.jpeg

发布溢出(逃逸)

一种错误的发布,当一个对象还没有构造完成时,就使它被其他线程所见.
下面实例中,就可能因为指令重排序问题导致this对象逃逸。


错误的发布
逃逸

2.4.2 安全发布对象

  • 在静态初始化函数中初始化一个对象引用.
public class StaticDemo {
    private StaticDemo(){}
    private static StaticDemo instance=new StaticDemo();
    public static StaticDemo getInstance(){
        return instance;
    }
}

静态初始化之所以安全是因为是在JVM类的初始化阶段进行初始化的,在类被加载后并在被线程使用之前,所有内存的写入操作会自动对所有线程可见。但该规则只适用于构造时的状态,不保证后续线程使用的安全性。

  • 将对象的引用保存到volatile类型的域或者AtomicReference对象中(利用volatile happen-before规则)
  • 将对象的引用保存到一个由锁保护的域中(读写都上锁)
public class VolatileSyncDemo {

    private VolatileSyncDemo(){}

    private volatile static VolatileSyncDemo instance=null;
    //DCL问题
    /**
     * instance = new VolatileSyncDemo();
     * ->
     * 1. memory=allocate()
     * 2.
     * 3. instance=memory
     *    ctorInstance(memory)
     *
     * 1.3.2 (不完整实例)
     */
    public static VolatileSyncDemo getInstance(){
        if(instance==null){

            // 把锁放入到判断里,而不是防止方法上,防止当instance不为null时,音方法加锁带来的性能损失。
            // 再锁里再进行一次判断就可以防止多线程时的多实例问题。但是会带来DCL问题,所以需要加入volatile,防止指令重排序。
            synchronized(VolatileSyncDemo.class) {
                if(instance==null) {
                    instance = new VolatileSyncDemo();
                }
            }
        }
        return instance;
    }
}

将对象的引用保存到某个正确构造对象的final类型域中(初始化安全性)

public class FinalDemo {
    private final Map<String,String> states;

    public FinalDemo(){
        states=new HashMap<>();
        states.put("mic","mic");
    }
}

对还有final域的对象,其防止指令重排序的作用能防止对象初始化引用被重排序到构造函数之外。后续使用的安全性还需要保障.

2.5 J.U.C核心之AQS

2.5.1 重入锁ReentrantLock的初步认识

什么是锁

锁是用来解决多线程并发访问共享资源所带来的数据安全性问题的手段。

对一个共享资源加锁后,如果有一个线程获得了锁,那么其他线程无法访问这个共享资源

什么是重入锁?

一个持有锁的线程,在释放锁之前,如果再次访问加了该同步锁的其他方法,这个线程不需要再次争抢锁,只需要记录重入次数。

public class LockDemo {
    static Lock lock=new ReentrantLock(); //重入锁
    public static int count=0;
    public static void incr(){ //递增
        lock.lock(); //获得锁 线程A
        try {
            Thread.sleep(1);
            decr();
            count++; //count++ (只会由一个线程来执行)
        } catch (InterruptedException e) {
            e.printStackTrace();
        }finally {
            lock.unlock(); //释放锁
        }
    }
    public static void decr(){ //递减
        lock.lock(); //有需要争抢锁 ,线程A (不需要争抢锁,记录重入次数即可)
        try {
            count -= 1;
        }catch (Exception e) {
            e.printStackTrace();
        }finally {
            lock.unlock(); //释放锁
        }
    }
    public static void main(String[] args) throws InterruptedException {
        for (int i = 0; i < 1000; i++) {
            new Thread(LockDemo::incr).start();
        }
        Thread.sleep(4000);
        System.out.println("result:"+count);
    }
}

什么是AQS

AbstractQueuedSynchronizer,类如其名,抽象的队列式的同步器,AQS定义了一套多线程访问共享资源的同步器框架,许多同步类实现都依赖于它,如常用的ReentrantLock/Semaphore/CountDownLatch...它维护了一个volatile int state(代表共享资源)和一个FIFO线程等待队列(多线程争用资源被阻塞时会进入此队列)。这里volatile是核心关键词,具体volatile的语义,在此不述。state的访问方式有三种:

  • getState()
  • setState()
  • compareAndSetState()

AQS定义两种资源共享方式:Exclusive(独占,只有一个线程能执行,如ReentrantLock)和Share(共享,多个线程可同时执行,如Semaphore/CountDownLatch)。

不同的自定义同步器争用共享资源的方式也不同。自定义同步器在实现时只需要实现共享资源state的获取与释放方式即可,至于具体线程等待队列的维护(如获取资源失败入队/唤醒出队等),AQS已经在顶层实现好了。自定义同步器实现时主要实现以下几种方法:

  • isHeldExclusively():该线程是否正在独占资源。只有用到condition才需要去实现它。
  • tryAcquire(int):独占方式。尝试获取资源,成功则返回true,失败则返回false。
  • tryRelease(int):独占方式。尝试释放资源,成功则返回true,失败则返回false。
  • tryAcquireShared(int):共享方式。尝试获取资源。负数表示失败;0表示成功,但没有剩余可用资源;正数表示成功,且有剩余资源。
  • tryReleaseShared(int):共享方式。尝试释放资源,如果释放后允许唤醒后续等待结点返回true,否则返回false。

以ReentrantLock为例,state初始化为0,表示未锁定状态。A线程lock()时,会调用tryAcquire()独占该锁并将state+1。此后,其他线程再tryAcquire()时就会失败,直到A线程unlock()到state=0(即释放锁)为止,其它线程才有机会获取该锁。当然,释放锁之前,A线程自己是可以重复获取此锁的(state会累加),这就是可重入的概念。但要注意,获取多少次就要释放多么次,这样才能保证state是能回到零态的。

再以CountDownLatch以例,任务分为N个子线程去执行,state也初始化为N(注意N要与线程个数一致)。这N个子线程是并行执行的,每个子线程执行完后countDown()一次,state会CAS减1。等到所有子线程都执行完后(即state=0),会unpark()主调用线程,然后主调用线程就会从await()函数返回,继续后余动作。

一般来说,自定义同步器要么是独占方法,要么是共享方式,他们也只需实现tryAcquire-tryRelease、tryAcquireShared-tryReleaseShared中的一种即可。但AQS也支持自定义同步器同时实现独占和共享两种方式,如ReentrantReadWriteLock。


AQS原理流程图

CountDownLatch是什么?

countdownlatch是一个同步工具类,它允许一个或多个线程一直等待,直到其他线程的操作执行完毕再执行。从命名可以解读到
countdown是倒数的意思,类似于我们倒计时的概念。

public class CountDownLatchDemo {

    public static void main(String[] args) throws InterruptedException {
        //state存储
        CountDownLatch countDownLatch=new CountDownLatch(3);

        new Thread(()->{
            countDownLatch.countDown(); //倒计时 3-1=2
            //修改state=state-1 通过cas设置到state这个字段上
        }).start();
        new Thread(()->{
            countDownLatch.countDown(); //倒计时 2-1=1
        }).start();
        new Thread(()->{
            countDownLatch.countDown(); //倒计时 1-1=0 ->触发唤醒操作
        }).start();

        countDownLatch.await(); //阻塞主线程
        System.out.println("线程执行完毕");
    }
}

什么是Semaphore

Semaphore 通常我们叫它信号量, 可以用来控制同时访问特定资源的线程数量,通过协调各个线程,以保证合理的使用资源。

acquire()  
获取一个令牌,在获取到令牌、或者被其他线程调用中断之前线程一直处于阻塞状态。
release()
释放一个令牌,唤醒一个获取令牌不成功的阻塞线程。
public class SemaphoreDemo {

    public static void main(String[] args) {
        //当前可以获得的最大许可数量是5个
        //AQS   ->state
        Semaphore semaphore=new Semaphore(5);
        for (int i = 0; i < 10; i++) {
            new Car(i,semaphore).start();
        }
    }
    static class Car extends Thread{
        private int num;
        private Semaphore semaphore;
        public Car(int num, Semaphore semaphore) {
            this.num = num;
            this.semaphore = semaphore;
        }
        @Override
        public void run() {
            try {
                semaphore.acquire(); //获得一个许可,如果获取不到许可,就会被阻塞
                System.out.println("第"+num+" 占用一个停车位");
                TimeUnit.SECONDS.sleep(2);
                System.out.println("第"+num+"  俩车走了");
                semaphore.release(); //释放许可
            } catch (InterruptedException e) {
                e.printStackTrace();
            }
        }
    }
}

什么是CyclicBarrier

CyclicBarrier 的字面意思是可循环使用(Cyclic)的屏障(Barrier)。它要做的事情是,让一组线程到达一个屏障(也可以叫同步点)时被阻塞,直到最后一个线程到达屏障时,屏障才会开门,所有被屏障拦截的线程才会继续工作。

CyclicBarrier实例:

数据加载类,调用CyclicBarrier

public class DataImportThread extends Thread{

    private String path;

    private CyclicBarrier cyclicBarrier;

    public DataImportThread(String path, CyclicBarrier cyclicBarrier) {
        this.path = path;
        this.cyclicBarrier = cyclicBarrier;
    }

    @Override
    public void run() {
        System.out.println("开始导入:"+path+" 位置的数据");
        try {
            cyclicBarrier.await(); //阻塞
            //TODO
        } catch (InterruptedException e) {
            e.printStackTrace();
        } catch (BrokenBarrierException e) {
            e.printStackTrace();
        }
    }
}

CycliBarrierDemo实例


public class CycliBarrierDemo extends Thread{
    @Override
    public void run() {
        System.out.println("开始进行数据汇总和分析");
    }

    /**
     * 1. parties ,如果因为某种原因导致没有足够多的线程来调用await,
     * 这个时候会导致所有线程都会被阻塞
     * 2. await(timeout,unit) 设置一个超市等待时间.
     * 3. reset重置计数,brokenBarrierException
     * @param args
     */

    public static void main(String[] args) {
        //parties=3
        CyclicBarrier cyclicBarrier=
                new CyclicBarrier(3,new CycliBarrierDemo());
        new DataImportThread("path1",cyclicBarrier).start();
        new DataImportThread("path2",cyclicBarrier).start();
        new DataImportThread("path3",cyclicBarrier).start();
        //TODO 希望三个线程执行结束之后,再做一个汇总处理.
    }
}
CyclicBarrier

什么是condition

Condition是一个多线程协调通信的工具类,可以让某些线程一起等待某个条件(condition),只有满足条件时,线程才会被唤醒。

// condition 等待类
public class ConditionDemoWait extends Thread{

    private Lock lock;
    private Condition condition;

    public ConditionDemoWait(Lock lock, Condition condition) {
        this.lock = lock;
        this.condition = condition;
    }

    /**
     * synchronized(class){
     *     class.wait();
     * }
     */
    @Override
    public void run() {
        System.out.println("begin - ConditionDemoWait");
        try {
            lock.lock();//ThreadA获得了锁
            condition.await(); //阻塞
            System.out.println("end - ConditionDemoWait");
        } catch (InterruptedException e) {
            e.printStackTrace();
        }finally {
            lock.unlock();
        }
    }
}
// condition 唤醒类
public class ConditionDemoNotify extends Thread{

    private Lock lock;
    private Condition condition;

    public ConditionDemoNotify(Lock lock, Condition condition) {
        this.lock = lock;
        this.condition = condition;
    }

    /**
     * synchronized(class){
     *     class.wait();
     * }
     */
    @Override
    public void run() {
        System.out.println("begin - ConditionDemoNotify");
        lock.lock(); //线程B被阻塞在这个位置,由于被线程A唤醒,所以线程B继续执行
        try {
            condition.signal();//线程B会执行这个方法
            System.out.println("end - ConditionDemoNotify");
        }finally {
            lock.unlock();
        }
    }
}
// condition 主类
public class CondtionDemo {

    public static void main(String[] args) {
        Lock lock=new ReentrantLock();
        Condition condition=lock.newCondition();
        ConditionDemoWait conditionDemoWait=new ConditionDemoWait(lock,condition);
        ConditionDemoNotify conditionDemoNotify=new ConditionDemoNotify(lock,condition);
        conditionDemoWait.start();
        conditionDemoNotify.start();

    }
}
condition原理

2.6 线程调度之线程池

2.6.1 什么是线程池

提前创建好若干个线程放在一个容器中。如果有任务需要处理,则将任务直接分配给线程池中的线程来执行,任务处理完以后这个线程不会被销毁,而是等待后续分配任务。

2.6.2 Java 中提供的线程池

  • newFixedThreadPool

    创建一个定长线程池,可控制线程最大并发数,超出的线程会在队列中等待。

public class ThreadPoolDemo implements Runnable{

    public static void main(String[] args) {

        ThreadPoolExecutor executorService=(ThreadPoolExecutor) ExecutorsSelf.newFixedThreadPool(3);
        executorService.prestartAllCoreThreads(); //可以提前预热所有核心线程
        for (int i = 0; i < 100; i++) {
            executorService.execute(new ThreadPoolDemo());
        }
        executorService.shutdown();
    }
    @Override
    public void run() {
        try {
            Thread.sleep(10);
        } catch (InterruptedException e) {
            e.printStackTrace();
        }
        System.out.println(Thread.currentThread().getName());
    }
}
  • newSingleThreadExecutor

    创建一个单线程化的线程池,它只会用唯一的工作线程来执行任务,保证所有任务按照指定顺序(FIFO, LIFO, 优先级)执行

  • newCachedThreadPool

    创建一个可缓存线程池,如果线程池长度超过处理需要,可灵活回收空闲线程,若无可回收,则新建线程。

  • newScheduledThreadPool

    创建一个定长线程池,支持定时及周期性任务执行.


    线程池提交任务的逻辑图

    image-20220308222631001.png

2.6.3 线程池的监控

线程池的监控一般通过自己实现线程池的类,并在其中进行监控。

public class ThreadPoolSelf extends ThreadPoolExecutor {

    public ThreadPoolSelf(int corePoolSize, int maximumPoolSize, long keepAliveTime, TimeUnit unit, BlockingQueue<Runnable> workQueue) {
        super(corePoolSize, maximumPoolSize, keepAliveTime, unit, workQueue);
    }

    @Override
    public void shutdown() {
        super.shutdown();
    }

    //任务执行开始
    @Override
    protected void beforeExecute(Thread t, Runnable r) {
        //TODO 通过一个属性来记录任务的开始时间
    }

    @Override
    protected void afterExecute(Runnable r, Throwable t) {
        //任务执行结束
        System.out.println("初始线程数:"+this.getPoolSize());
        System.out.println("核心线程数:"+this.getCorePoolSize());
        System.out.println("正在执行的任务数量:"+this.getActiveCount());
        System.out.println("已经执行的任务数:"+this.getCompletedTaskCount());
        System.out.println("任务总数"+this.getTaskCount());
    }
}

2.6.4 带返回值的线程

public class CallableFutureDemo implements Callable<String> {
    @Override
    public String call() throws Exception {
        System.out.println("Hello Mic");
        Thread.sleep(3000); //睡眠3s
        return "Mic";
    }

    public static void main(String[] args) throws ExecutionException, InterruptedException {
        CallableFutureDemo callableFutureDemo=new CallableFutureDemo();
        FutureTask futureTask=new FutureTask(callableFutureDemo);
        new Thread(futureTask).start();
        //get方法是属于阻塞方法
        System.out.println(futureTask.get());

        ExecutorService executorService= Executors.newFixedThreadPool(1);
        CallableFutureDemo callableFutureDemo2=new CallableFutureDemo();
        FutureTask future=(FutureTask) executorService.submit(callableFutureDemo2);
        System.out.println(future.get());

    }
}

2.7 多线程并发拓展

2.7.1 死锁发生的条件

  • 互斥,共享资源 X 和 Y 只能被一个线程占用;
  • 占有且等待,线程 T1 已经取得共享资源 X,在等待共享资源 Y 的时候,不释放共享资源 X;
  • 不可抢占,其他线程不能强行抢占线程 T1 占有的资源;
  • 循环等待,线程 T1 等待线程 T2 占有的资源,线程 T2 等待线程 T1 占有的资源,就是循环等待。

四个条件同时满足时,产生死锁。

pojo

public class Account {
    private String accountName;
    private int balance; //余额

    public Account(String accountName, int balance) {
        this.accountName = accountName;
        this.balance = balance;
    }

    public String getAccountName() {return accountName; }

    public void setAccountName(String accountName) {this.accountName = accountName; }

    public int getBalance() {        return balance;    }

    public void setBalance(int balance) {this.balance = balance; }
    public void debit(int amount){this.balance-=amount;}
    public void credbit(int amount){this.balance+=amount;}
public class TransferAccount implements Runnable{
    private Account fromAccount; //转出账户
    private Account toAccount; //转入账户
    private int amount;

    public TransferAccount(Account fromAccount, Account toAccount, int amount) {
        this.fromAccount = fromAccount;
        this.toAccount = toAccount;
        this.amount = amount;
    }

    @Override
    public void run() {
        while(true){
            synchronized (fromAccount){
                synchronized (toAccount){
                    if(fromAccount.getBalance()>=amount){
                        fromAccount.debit(amount);
                        toAccount.credbit(amount);
                    }
                }
                System.out.println(fromAccount.getAccountName()+"----"+fromAccount.getBalance());
                System.out.println(toAccount.getAccountName()+"----"+toAccount.getBalance());
            }
        }
    }
}

程序入口

public class TestMain {

    public static void main(String[] args) {
        Account fromAccount=new Account("张三",100000);
        Account toAccount=new Account("李四",200000);
        Allocator allocator=new Allocator(); //统一分配锁
        Thread a=new Thread(new TransferAccount01(fromAccount,toAccount,1,allocator));
        Thread b=new Thread(new TransferAccount01(toAccount,fromAccount,2,allocator));
        a.start();
        b.start();
    }
}

此时,该程序因为a b两个账户互相争抢资源导致死锁。

2.7.2 如何解决死锁问题

一旦发生死锁,一般没什么好的方法来解决,只能通过重启应用。所以如果要解决死锁问题,最好的方式就是提前规避。

因为死锁的四个条件中 互斥,共享资源 X 和 Y 只能被一个线程占用 是先决条件,因此只能用其他三个条件入手。

  1. 同时申请两个临界资源
    public class Allocator {

    private List<Object> list=new ArrayList<>();
    /**

    • 申请资源
    • @param from
    • @param to
    • @return
      */
    synchronized boolean apply(Object from,Object to){
        if(list.contains(from)||list.contains(to)){
            return false;
        }else {
            list.add(from);
            list.add(to);
            return true;
        }
    }
    synchronized void free(Object from,Object to){
        list.remove(from);
        list.remove(to);
    }
}
public class TransferAccount implements Runnable{
   
@Override
    public void run() {
        while (true) {
            // run方法改写成如下:
            if (allocator.apply(fromAccount, toAccount)) { //都会在这个地方去获得临界资源
                try {
                    if (fromAccount.getBalance() >= amount) {
                        fromAccount.debit(amount);
                        toAccount.credbit(amount);
                
                        System.out.println(fromAccount.getAccountName() + "----" + fromAccount.getBalance());
                        System.out.println(toAccount.getAccountName() + "----" + toAccount.getBalance());
                    }
                } finally {
                    allocator.free(fromAccount, toAccount);
                }
            }
        }
      }
}

此时因为程序会同时去获取两个临界资源,因此不会导致死锁。

  1. 当抢资源失败时,释放已占有资源,破坏不可抢占的条件。

synchronized 天然具有阻塞性,抢占到锁是就会阻塞,用lock来进行加锁。

public class TransferAccount02 implements Runnable {
    private Account fromAccount; //转出账户
    private Account toAccount; //转入账户
    private int amount;
    Lock fromLock = new ReentrantLock();
    Lock toLock = new ReentrantLock();

    public TransferAccount02(Account fromAccount, Account toAccount, int amount) {
        this.fromAccount = fromAccount;
        this.toAccount = toAccount;
        this.amount = amount;
    }

    @Override
    public void run() {
        while (true) {
            //  fromLock.tryLock() 返回布尔值,不会阻塞
            if (fromLock.tryLock()) {
                if (toLock.tryLock()) {
                    if (fromAccount.getBalance() >= amount) {
                        fromAccount.debit(amount);
                        toAccount.credbit(amount);
                    }
                    System.out.println(fromAccount.getAccountName() + "----" + fromAccount.getBalance());
                    System.out.println(toAccount.getAccountName() + "----" + toAccount.getBalance());
                }
            }
        }
    }
}
  1. 破坏循环等待,让加锁的循序一致。
@Override
    public void run() {
        Account left=null;
        Account right=null;
        if(fromAccount.hashCode()>toAccount.hashCode()){
            left=toAccount;
            right=fromAccount;
        }
        while(true){
            synchronized (left) {
                synchronized (right) {//申请不到资源,就已经阻塞
                    if (fromAccount.getBalance() >= amount) {
                        fromAccount.debit(amount);
                        toAccount.credbit(amount);
                    }
                }
                System.out.println(fromAccount.getAccountName() + "----" + fromAccount.getBalance());
                System.out.println(toAccount.getAccountName() + "----" + toAccount.getBalance());
            }
        }
    }

2.7.3 单机版 MapReduce-Fork-Join

2.7.3.1什么是Fork-Join

Fork/Join框架是Java 7提供的一个用于并行执行任务的框架,是一个把大任务分割成若干个小任务,最终汇总每个小任务结果后得到大任务结果的框架。


Fork-Join

2.7.3.2 Fork-Join 实例

public class ForkJoinDemo extends RecursiveTask<Integer> {
    // 分隔的阈值
    private final int THREADHOLD = 11;

    private int start;
    private int end;

    public ForkJoinDemo(int start, int end) {
        this.start = start;
        this.end = end;
    }

    @Override
    protected Integer compute() {
        int sum = 0;
        if (end - start > THREADHOLD) {
            int middle = (end + start) /2;
            ForkJoinDemo left = new ForkJoinDemo(start, middle);
            ForkJoinDemo right = new ForkJoinDemo(middle + 1, end);
            left.fork();
            right.fork();
            int leftResult = left.join();
            int rightResult = right.join();
            sum = leftResult + rightResult;
        } else {
            System.out.println("start:" + start + "  end:" + end);
            for (int i = start; i <= end; ++i) {
                sum += i;
            }
         }

        return sum;
    }

    public static void main(String[] args) throws ExecutionException, InterruptedException {
        ForkJoinPool forkJoinPool = new ForkJoinPool();
        ForkJoinDemo forkJoinDemo = new ForkJoinDemo(1, 100);
        Future<Integer> submit = forkJoinPool.submit(forkJoinDemo);
        System.out.println(submit.get());
    }
}
最后编辑于
©著作权归作者所有,转载或内容合作请联系作者
【社区内容提示】社区部分内容疑似由AI辅助生成,浏览时请结合常识与多方信息审慎甄别。
平台声明:文章内容(如有图片或视频亦包括在内)由作者上传并发布,文章内容仅代表作者本人观点,简书系信息发布平台,仅提供信息存储服务。

相关阅读更多精彩内容

友情链接更多精彩内容