JUC
快来分享你的内容吧~
- 偏向锁在什么条件下会升级为轻量级锁?问题描述synchronized锁升级时,偏向锁在什么条件下会升级为轻量级锁?背景信息Java版本为java11。代码import org.openjdk.jol.info.ClassLayout;public class BiasedLockThreadIdCheck { static final Object lock = new Object(); public static v...查看全文编程导航小智:Java偏向锁升级为轻量级锁的条件分析这是一个关于Java锁机制的有趣问题。我来详细解释偏向锁升级为轻量级锁的条件和原理。偏向锁升级为轻量级锁的主要条件1. 存在锁竞争:当另一个线程尝试获取已经被偏向锁持有的对象时2. 撤销偏向锁:JVM需要撤销当前偏向锁,然后才能升级为轻量级锁3. 全局安全点检查:撤销偏向锁需要在全局安全点(global safepoint)进行为什么sleep会影响锁升级?在
Day28 JMM内存模型
- **ForkJoinPool** Fork/Join是一个并行计算框架,主要是用来支持分治模型的,Fork对应的是分治模型的任务分解,Join对应的是结果集的合并。核心思想是价格大任务拆成小任务求解结果集,然后将小结果集进行合并成最终结果,适用于可以采用分治策略的计算密集型任务。 - ForkJoinPool使用 - 构造器 有四个核心参数,分别用于控制线程池的并行数、工作线程的创建、异常处理和模式指定 **int parallelism**: 指定并行级别,决定工作线程的数量,不设置默认使用Runtime.getRuntime().availableProcessors() **ForkJoinWorkerThreadFactory factory**:ForkJoinPool在创建线程时,会通过factory来创建,不指定默认使用DefaultForkJoinWorkerThreadFactory **UncaughtExceptionHandler handler**:指定异常处理器,当任务在运行中出错时,将由设定的处理器处理 **boolean asyncMode**:设置队列的工作模式。当asyncMode为true时,将使用先进先出队列,而为false时则使用后进先出的模式。 - 任务提交方式 ```java // 异步执行 void execute(ForkJoinTask<?> task) void execute(Runnable task) // 等待获取结果 T invoke(ForkJoinTask<T> task) // 提交执行获取Future结果 ForkJoinTask<T> submit(ForkJoinTask<T> task) ForkJoinTask<T> submit(Callable<T> task) ForkJoinTask<T> submit(Runnable task) ForkJoinTask<T> submit(Runnable task, T result) ``` > ForkJoinTask是一个抽象类,定义了执行任务的基本接口,可以通过继承ForkJoinTask并重写compute方法自定义实现。通常情况下仅需继承它的子类 **RecursiveAction**:用于递归执行但不需要返回结果的任务; **RecursiveTask** :用于递归执行需要返回结果的任务。 CountedCompleter :在任务完成执行后会触发执行一个自定义的钩子函数 compute(): 业务执行逻辑 fork(): 用于向当前任务所运行的线程池中提交任务 join(): 获取认为执行结果 > ```java public class Fibonacci extends RecursiveTask<Integer> { final int n; Fibonacci(int n) { this.n = n; } /** * 重写RecursiveTask的compute()方法 * @return */ protected Integer compute() { if (n <= 1) return n; Fibonacci f1 = new Fibonacci(n - 1); //提交任务 f1.fork(); Fibonacci f2 = new Fibonacci(n - 2); //合并结果 return f2.compute() + f1.join(); } public static void main(String[] args) { //构建forkjoin线程池 ForkJoinPool pool = new ForkJoinPool(); Fibonacci task = new Fibonacci(10); //提交任务并一直阻塞直到任务 执行完成返回合并结果。 int result = pool.invoke(task); System.out.println(result); } } ``` - 处理阻塞任务 使用ForkJoinPool处理阻塞型任务时需要注意: 1. 防止线程饥饿。当一个线程正在执行一个阻塞型任务时,会一直等待任务完成,如果没有其他线程可以窃取任务,该线程将会一直阻塞 2. 使用特定线程池。为了最大程度地利用ForkJoinPool的性能,可以使用专门的线程池来处理阻塞型任务,这些线程不会被ForkJoinPool的窃取机制所影响 3. 不要阻塞工作线程。如果在ForkJoinPool中使用阻塞型任务,需要确保这些任务不会阻塞工作线程,否则会导致整个线程池的性能下降。可以将阻塞型任务提交到一个专门的线程池中,或者使用CompletableFuture等异步编程工具来处理阻塞型任务。 ```java public class BlockingTaskDemo { public static void main(String[] args) { //构建一个forkjoin线程池 ForkJoinPool pool = new ForkJoinPool(); //创建一个异步任务,并将其提交到ForkJoinPool中执行 CompletableFuture<String> future = CompletableFuture.supplyAsync(() -> { try { // 模拟一个耗时的任务 TimeUnit.SECONDS.sleep(5); return "Hello, world!"; } catch (InterruptedException e) { e.printStackTrace(); return null; } }, pool); try { // 等待任务完成,并获取结果 String result = future.get(); System.out.println(result); } catch (InterruptedException e) { e.printStackTrace(); } catch (ExecutionException e) { e.printStackTrace(); } finally { //关闭ForkJoinPool,释放资源 pool.shutdown(); } } } ``` - 工作原理 - ForkJoinPool的任务会被内部存储了一个WorkQueue数组,提交给ForkJoinPool的任务会被分配到指定的WorkQueue上执行 - 每个WorkQueue内部维护了一个ForkJoinTask数组用来存储待执行的任务,以及一个独立的ForkJoinWorkerThread用来真正执行任务 - 当有某个线程空闲时,会去窃取其他繁忙线程的任务拿过来执行  **ForkJoinWorkerThread** ForkJoinWorkerThread是ForkJoinPool中的一个专门用于执行任务的线程。当一个ForkJoinWorkerThread被创建时,它会自动注册一个WorkQueue到ForkJoinPool中。这个WorkQueue是该线程专门用于存储自己的任务的队列,只能出现在WorkQueues[]的奇数位。 ForkJoinWorkerThread工作线程启动后就会扫描偷取任务执行,另外当其在 ForkJoinTask#join() 等待返回结果时如果被 ForkJoinPool 线程池发现其任务队列为空或者已经将当前任务执行完毕,也会通过工作窃取算法从其他任务队列中获取任务分配到其任务队列中并执行。  **WorkQueue** WorkQueue是一个双端队列,用于存储工作线程自己的任务。每个工作线程都会维护一个本地的WorkQueue,并且优先执行本地队列中的任务。当本地队列中的任务执行完毕后,工作线程会尝试从其他线程的WorkQueue中窃取任务。 WorkQueue 任务队列其实也分为了两种类型,一种是外部提交进来的任务所占用的队列,其在任务队列数组中的数组下标为偶数;另一种是属于工作线程私有的任务队列,保存大任务 fork 分解出来的任务,其在任务队列数组中的数组下标为奇数。  **工作窃取** 就是允许空闲线程从繁忙线程的双端队列中窃取任务。默认情况下,工作线程从它自己的双端队列的头部获取任务。当自己的任务为空时,线程会从其他繁忙线程双端队列的尾部中获取任务。最大限度地减少了线程竞争任务的可能性,提高工作效率  - 执行流程  ## JMM内存模型 是JVM定义的一套规范,规定了多线程程序中的变量如何在内存中存储和传递,约定了线程何时从主内存读取数据、何时把数据写回主内存。其核心目标是确保多线程环境下的**可见性、有序性和原子性** 1. 可见性:一个线程对变量的修改能被其他线程看到,volatile关键字就是用来保证可见性的,强制线程每次读写都直接跟主内存交互 2. 有序性:线程执行操作的顺序。JMM允许指令重排来提高性能,但通过happens-before关系保证跨线程的有序性 3. 原子性:操作不可分割,执行过程不会被打断。synchronized关键字可以保证代码块的原子性 - JMM的抽象内存模型组成 - 主内存存放共享变量,所有线程都能访问 - 每个线程都有自己的本地内存,存放共享变量的副本 - 线程对变量的操作必须在本地内存中进行,不能直接操作主内存 - 线程间变量传递必须通过主内存完成  - Happens-Before 定义了某个操作的结果对另一个操作可见,用于确定两个操作之间的执行顺序,确保多线程程序的正确性和一致性,底层主要是利用内存屏障来实现的。 Happens-Before规则包括: > 只是说呈现给开发者的规则是如此,但并不代表会严格按此执行 > 1. 程序顺序规则:在一个线程内,按照代码顺序,前面的操作→Happens-Before→后面的操作 2. 监视器锁规则:对一个锁的解锁操作Happens-Before后续对这个锁的加锁操作 3. volatile变量规则:对一个volatile变量的写操作Happens-Before后续对这个变量的读操作 4. 传递规则:如果A→Happens-Before→B,B→Happens-Before→C,那么A→Happens-Before→C 5. 线程启动规则:对线程的Thread.start()调用Happens-Before该线程的每一个操作 6. 线程终止规则:线程中的所有操作 Happens-Before 其他线程检测到线程已终止,通过Thread.join,Thread.isAlive等 7. 线程中断规则:对线程的interrupt()调用Happens-Before检测到中断事件 8. 对象终结规则:一个对象的初始化完成Happens-Before它的finalize()方法开始。  - 内存屏障 是一种CPU指令,用于禁止特定类型的指令重排序,JVM在编译时会根据关键字(如volatile)插入相应的内存屏障 - LoadLoad屏障:禁止读操作重排序,按序读 - StoreStore:禁止写操作重排序,按序写 - LoadStore:禁止读操作与后面的写操作重排序 - StoreLoad:禁止写操作与后面的读操作重排序 - Volatile的语义 - 保证可见性,写操作立即刷新到主内存,读操作从主内存获取最新值 - 禁止指令重排序,编译器和CPU不会把volatile变量的读写操作和其他操作乱序执行 但不能保证原子性,i++这种复合操作(读i→加1→写回i),对于单次读或单次写是可以保证原子性的  不能保证原子性的场景: - i++;读i→加1→写回i - i = i + 1;读i→加1→写回i - long x;x = 1L;long类型64位,在32位系统会被拆成两部分,不保证原子性,除非在64位系统或加volatile关键字(JVM规定,对long/double之外的所有基本类型的单次操作都是原子的,加上volatile之后读写也是原子的) - list.add(item);即使list是volatile,add()是方法调用,非原子;只保证引用本身的原子性和可见性,不保证内部的线程安全
Day25 线程池
## 线程池 - 简介 线程池(Thread Pool)是一种基于池化思想管理线程的工具,线程池维护多个线程,等待监督管理者分配可并发执行的任务。避免了处理任务时创建销毁线程开销的代价,另外也避免了线程数量膨胀导致的过分调度问题,保证了对内核的充分利用。 - **降低资源消耗**:通过池化技术重复利用已创建的线程,降低线程创建和销毁造成的损耗。 - **提高响应速度**:任务到达时,无需等待线程创建即可立即执行。 - **提高线程的可管理性**:线程是稀缺资源,如果无限制创建,不仅会消耗系统资源,还会因为线程的不合理分布导致资源调度失衡,降低系统的稳定性。使用线程池可以进行统一的分配、调优和监控。 - **提供更多更强大的功能**:线程池具备可拓展性,允许开发人员向其中增加更多的功能。比如延时定时线程池ScheduledThreadPoolExecutor,就允许任务延期执行或定期执行。 - 线程池使用 - ThreadPoolExecutor 推荐 - 核心参数 - **corePoolSize**:核心线程数,线程池初始化时默认是没有线程的,当任务来临时才开始创建线程去执行任务 - **maximumPoolSize**:最大线程数,在核心线程数已满,且队列已满时,如果池子里的工作线程数小于maximumPoolSize,则会创建非核心线程执行任务。一般设置为 (最大任务数-任务队列长度)* 单个任务执行时间 - **keepAliveTime**:非核心线程数的空闲时间超过keepAliveTime就会被自动终止回收掉,但在corePoolSize=maximumPoolSize时,该值无效,因为不存在非核心线程 - **unit**:keepAliveTime的时间单位 - **workQueue**:用于保存线程任务的队列,主要分为无界、有界、同步移交等队列,当池子里的工作线程数大于corePoolSize,就会将新进来的线程任务放入队列中 【见并发容器——BlockQueue】一般设置为**核心线程数/单个任务执行耗时 * 2** ```java // private Runnable getTask() => for(;;) try { // 使用阻塞队列,万一队列为空,那这里就会阻塞住直到有任务 // 如果使用非阻塞队列,这里的死循环会反复空转,浪费CPU资源 Runnable r = timed ? workQueue.poll(keepAliveTime, TimeUnit.NANOSECONDS) : workQueue.take(); if (r != null) return r; timedOut = true; } catch (InterruptedException retry) { timedOut = false; } ``` - **threadFactory**: - 创建线程的工厂接口,默认使用Executors.defaultThreadFactory() - 另外可以实现ThreadFactory接口,自定义线程工厂 - **handler**:线程池无法继续接收任务时(workQueue已满和maximumPoolSize已满)的拒绝策略 - AbortPolicy:默认拒绝策略,线程池和队列都满时不再接受新的提交,抛出RejectedExecutionException异常 - CallerRunsPolicy:让提交任务的主线程来执行任务 - DiscardOldestPolicy:丢弃在队列中存在时间最久的任务,重复执行 - DiscardPolicy:丢弃任务,不进行任何通知 - 另外可以实现RejectedExecutionHandler接口,自定义拒绝策略 - 使用示例 ```java import java.util.concurrent.ArrayBlockingQueue; import java.util.concurrent.ThreadPoolExecutor; import java.util.concurrent.TimeUnit; import java.util.concurrent.RejectedExecutionHandler; import java.util.concurrent.ThreadPoolExecutor.AbortPolicy; public class ThreadPoolExample2 { public static void main(String[] args) { // 线程池的核心线程数 int corePoolSize = 5; // 线程池的最大线程数 int maximumPoolSize = 10; // 线程池的任务队列 ArrayBlockingQueue<Runnable> workQueue = new ArrayBlockingQueue<>(100); // 线程池保持空闲的时间 long keepAliveTime = 60L; // 时间单位 TimeUnit unit = TimeUnit.SECONDS; // 线程池的拒绝策略 RejectedExecutionHandler handler = new ThreadPoolExecutor.AbortPolicy(); // 创建线程池 ThreadPoolExecutor executor = new ThreadPoolExecutor( corePoolSize, maximumPoolSize, keepAliveTime, unit, workQueue, handler ); // 提交任务到线程池 for (int i = 0; i < 15; i++) { final int taskId = i; executor.execute(() -> { System.out.println("Task " + taskId + " is running by " + Thread.currentThread().getName()); }); } // 关闭线程池 executor.shutdown(); try { // 等待所有任务完成,超时时间为60秒 if (!executor.awaitTermination(60, TimeUnit.SECONDS)) { // 如果超时后任务仍未完成,则强制关闭线程池 executor.shutdownNow(); } } catch (InterruptedException e) { // 如果等待过程中被中断,也强制关闭线程池 executor.shutdownNow(); } System.out.println("All tasks are done or interrupted."); } } ``` - Executors 不推荐 提供一些静态类创建线程池,但风险高,一些关键参数不可控 - 常用线程池类型 - **newFixedThreadPool** 定长线程池,核心线程池固定,使用无界队列,任务数大于线程数时入队,永不拒绝;会**有OOM风险。**仅适合任务量稳定可预测的小型应用 - **newCachedThreadPool** 弹性伸缩线程池,使用无上限的非核心线程池,来一个任务就建一个线程;**高并发时线程数过多容易导致系统崩溃**。仅适合执行大量耗时极短且并发可控的任务。 - **newSingleThreadExecutor** 定时线程池,支持周期性任务执行,核心线程池固定,非核心线程池无上限,使用无界延迟队列;**线程数爆炸+队列堆积→OOM,**可手动设置maximumPoolsize - **newSingleThreadScheduledExecutor** 单线程线程池,只有一个工作线程,任务按顺序执行,使用无界队列,**有OOM风险**;适合需要异步串行执行任务的场景。 - **newWorkStealingPool**(JDK8+)基于ForkJoinPool的抢占式线程池,空闲线程会从其他繁忙线程的队列拿任务;适合分治任务 - 使用示例 ```java import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; public class ThreadPoolExample { public static void main(String[] args) { // 创建一个固定大小的线程池 ExecutorService executor = Executors.newFixedThreadPool(5); // 提交任务到线程池 for (int i = 0; i < 10; i++) { final int taskId = i; executor.submit(() -> { System.out.println("Task " + taskId + " is running by " + Thread.currentThread().getName()); }); } // 关闭线程池 executor.shutdown(); while (!executor.isTerminated()) { // 等待所有任务完成 } System.out.println("All tasks are done."); } } ``` - 任务提交 - 无返回值的任务使用 public void execute(Runnable command) 方法提交; - 有返回值的任务使用: - Future submit(Runnable task) : 提交Runnable任务 - Future submit(Runnable task, T result): 提交Runnable任务并指定执行结果 - Future submit(Callable task) : 提交Callable任务 - 批量任务 ```java # 执行批量任务,返回它们的执行结果 public <T> List<Future<T>> invokeAll(Collection<? extends Callable<T>> tasks) throws InterruptedException #执行批量任务,返回指定时间内完成的执行结果,取消未完成的任务 public <T> List<Future<T>> invokeAll(Collection<? extends Callable<T>> tasks,long timeout, TimeUnit unit) # 执行批量任务,返回最先完成的执行结果 public <T> T invokeAny(Collection<? extends Callable<T>> tasks) throws InterruptedException, ExecutionException #执行批量任务,返回指定时间内最先完成的执行结果,取消未完成的任务 public <T> T invokeAny(Collection<? extends Callable<T>> tasks,long timeout, TimeUnit unit) ``` - 定时/延时任务 ```java #具备执行定时、延时、周期性任务的线程池 public class ScheduledThreadPoolExecutor extends ThreadPoolExecutor implements ScheduledExecutorService { #延时执行Runnable任务,只执行一次 public ScheduledFuture<?> schedule(Runnable command,long delay,TimeUnit unit) #延时执行Callable任务,只执行一次 public <V> ScheduledFuture<V> schedule(Callable<V> callable,long delay,TimeUnit unit) #廷时一段时间后,周期性执行Runnable任务,周期为固定时间 public ScheduledFuture<?> scheduleAtFixedRate(Runnable command,long initialDelay,long period,TimeUnit unit) #廷时一段时间后,周期性执行Runnable任务,周期为间隔时间 public ScheduledFuture<?> scheduleWithFixedDelay(Runnable command,long initialDelay,long delay,TimeUnit unit) ``` - 线程池关闭 - shutdownNow():**立即关闭**线程池,正在执行中的任务和队列中的任务都会被中断,关闭后状态为STOP,同时返回被中断的队列中的任务列表。 - shutdown():**延时关闭**线程池,正在执行中的任务和队列中的任务都能执行完成,关闭后状态为SHUTDOWN,后续进来的新任务会被执行拒绝策略。 - isTerminated():当正在执行的任务和队列中的任务全部都执行完时返回true。 - 原理分析 - 任务执行流程 submit方法底层实际调用的是execute,只是用Future包装了一层 ```java public Future<?> submit(Runnable task) { if (task == null) throw new NullPointerException(); RunnableFuture<Void> ftask = newTaskFor(task, null); execute(ftask); return ftask; } ``` execute()执行分为三步  在实际运行时,不管线程池中的线程是否空闲,只要数量小于核心线程数就会创建新线程。 **线程发生异常时,当前线程会被移出线程池**,源码中任务出现异常时,在finally块里会执行processWorkerExit()方法,处理当前线程并新增一个线程,以维持固定的核心线程数 - 线程池的流转状态 - RUNNING:**会**接收新任务并且**会**处理队列中的任务 - SHUTDOWN:**不会**接收新任务并且**会**处理队列中的任务 - STOP:**不会**接收新任务并且**不会**处理队列中的任务,并且会中断在处理的任务(注意:一个任务能不能被中断得看任务本身) - TIDYING:所有任务都终止了,线程池中也没有线程了,这样线程池的状态就会转为TIDYING,一旦达到此状态,就会调用线程池的terminated() - TERMINATED:terminated()执行完之后就会转变为TERMINATED RUNNING → SHUTDOWN:调用shutdown()时触发(GC时会调用) RUNNING/SHUTDOWN→STOP:调用shutdownNow()时触发 SHUTDOWN→TIDYING:队列为空且线程池中没有线程时自动转换 STOP→TIDYING:线程池中没有线程时转换 TIDYING→TERMINATED:调用terminated()  - 线程池参数动态配置 > 背景: 在日常项目中使用线程池来处理一些并发场景,提高任务处理的效率。但是实际使用时,无法准确地设置线程池参数,只能在运行过程中,不断去调整参数,然后重启服务。 >  - 基于Nacos实现 借助Nacos的Listener,在Bean初始化的时候,开启Nacos配置变更监听 ```java @Configuration @Data public class MyDynamicThreadPool implements InitializingBean { @Value("${threadPool.corePoolSize}") private int corePoolSize; @Value("${threadPool.maxPoolSize}") private int maxPoolSize; @Value("${threadPool.queueCapacity}") private int queueCapacity; @Value("${threadPool.keepAliveSeconds}") private int keepAliveSeconds; private static ThreadPoolTaskExecutor threadPoolTaskExecutor; @Autowired private NacosConfigManager nacosConfigManager; @Autowired private NacosConfigProperties nacosConfigProperties; @Override public void afterPropertiesSet() throws Exception { threadPoolTaskExecutor = new ThreadPoolTaskExecutor(); threadPoolTaskExecutor.setCorePoolSize(corePoolSize); threadPoolTaskExecutor.setMaxPoolSize(maxPoolSize); threadPoolTaskExecutor.setQueueCapacity(queueCapacity); threadPoolTaskExecutor.setKeepAliveSeconds(keepAliveSeconds); threadPoolTaskExecutor.setThreadNamePrefix( "Fox--"); threadPoolTaskExecutor.setRejectedExecutionHandler( new RejectedExecutionHandler() { @Override public void rejectedExecution(Runnable r, ThreadPoolExecutor executor) { System.out.println("队列已满,丢弃任务"); } }); threadPoolTaskExecutor.initialize(); nacosConfigManager.getConfigService().addListener("threadPool.yml", nacosConfigProperties.getGroup(), new Listener() { @Override public Executor getExecutor() { return null; } @Override public void receiveConfigInfo(String configInfo) { System.out.println("动态修改前-->"); print(); Yaml yaml = new Yaml(); InputStream inputStream = new ByteArrayInputStream(configInfo.getBytes()); Map<String, Object> dataMap = yaml.load(inputStream); // 将Map转换为JSONObject JSONObject pool = new JSONObject(dataMap).getJSONObject("threadPool"); threadPoolTaskExecutor.setCorePoolSize(pool.getInteger("corePoolSize")); threadPoolTaskExecutor.setMaxPoolSize(pool.getInteger("maxPoolSize")); threadPoolTaskExecutor.setQueueCapacity(pool.getInteger("keepAliveSeconds")); threadPoolTaskExecutor.setQueueCapacity(pool.getInteger("queueCapacity")); System.out.println("动态修改后-->"); print(); } }); } //执行任务 public void execute(Runnable runnable){ threadPoolTaskExecutor.execute(runnable); } public void print(){ System.out.println("核心线程数:" + threadPoolTaskExecutor.getThreadPoolExecutor().getCorePoolSize() + " " +"最大线程数:" + threadPoolTaskExecutor.getThreadPoolExecutor().getMaximumPoolSize() +" " + "阻塞队列数:" + threadPoolTaskExecutor.getThreadPoolExecutor().getQueue().size() + "/" + queueCapacity +" " + "活跃线程数:" + threadPoolTaskExecutor.getThreadPoolExecutor().getActiveCount()); } } ``` - 使用美团开源的DynamicTp-基于配置中心的轻量级动态可监控线程池 支持动态调参、通知告警、运行监控、三方包线程池管理等功能[https://dynamictp.cn](https://dynamictp.cn/)  ```yaml spring: dynamic: tp: enabled: true # 是否启用 dynamictp,默认true executors: # 动态线程池配置,都有默认值,采用默认值的可以不配置该项,减少配置量 - threadPoolName: dtpExecutor1 # 线程池名称,必填 threadPoolAliasName: 测试线程池 # 线程池别名,可选 executorType: common # 线程池类型 common、eager、ordered、scheduled、priority,默认 common corePoolSize: 5 # 核心线程数,默认1 maximumPoolSize: 8 # 最大线程数,默认cpu核数 queueCapacity: 2000 # 队列容量,默认1024 queueType: VariableLinkedBlockingQueue # 任务队列,查看源码QueueTypeEnum枚举类,默认VariableLinkedBlockingQueue rejectedHandlerType: CallerRunsPolicy # 拒绝策略,查看RejectedTypeEnum枚举类,默认AbortPolicy keepAliveTime: 10 # 空闲线程等待超时时间,默认60 threadNamePrefix: Fox # 线程名前缀,默认dtp allowCoreThreadTimeOut: false # 是否允许核心线程池超时,默认false waitForTasksToCompleteOnShutdown: true # 参考spring线程池设计,优雅关闭线程池,默认true awaitTerminationSeconds: 5 # 优雅关闭线程池时,阻塞等待线程池中任务执行时间,默认3,单位(s) preStartAllCoreThreads: false # 是否预热所有核心线程,默认false ```
Day24 ReentrantLock、并发容器
- **ReentrantLock** 可重入悲观独占锁,与其他两者的比较 | 特性 | `synchronized` | `ReentrantLock` | CAS | | --- | --- | --- | --- | | 实现 | JVM 内置(C++ Monitor) | JDK 层(AQS + CAS) | CPU 指令 | | 锁类型 | 可重入、悲观锁 | 可重入、悲观锁 | 乐观锁 | | 阻塞 | 是(重量级时) | 是(可中断) | 否(自旋) | | 功能 | 自动释放 | 支持超时、公平锁(默认非公平) | 仅原子操作 | | 性能(低竞争) | 极优(偏向/轻量级) | 略逊 | 最优 | | 性能(高竞争) | 差(重量级) | 较好(CLH 队列) | 差(自旋) | - Lock接口的常用API > 加锁后一定要及时释放!!! > - void lock() 获取锁,调用该方法当前线程会获取锁,当锁获得后,该方法返回 - void lockInterruptibly() throws InterruptedException可中断的获取锁,和lock()方法不同之处在于该方法会响应中断,即在锁的获取中可以中断当前线程 - boolean tryLock()尝试非阻塞的获取锁,调用该方法后立即返回。如果能够获取到返回true,否则返回false - boolean tryLock(long time, TimeUnit unit) throws InterruptedException超时获取锁,当前线程在以下三种情况下会被返回:当前线程在超时时间内获取了锁当前线程在超时时间内被中断超时时间结束,返回false - void unlock() 释放锁 - Condition newCondition()获取等待通知组件,该组件和当前的锁绑定,当前线程只有获取了锁,才能调用该组件的await()方法,而调用后,当前线程将释放锁 - ReentrantLock使用 - 公平锁与非公平锁 - 公平锁:线程在获取锁时,按照等待的先后顺序获取锁。 - 非公平锁:线程在获取锁时,不按照等待的先后顺序获取锁,而是随机获取锁。ReentrantLock默认是非公平锁 ```java ReentrantLock lock = new ReentrantLock(); //参数默认false,不公平锁 ReentrantLock lock = new ReentrantLock(true); //公平锁 ``` ```java class Counter { private final ReentrantLock lock = new ReentrantLock(); // 创建 ReentrantLock 对象 public void recursiveCall(int num) { lock.lock(); // 获取锁 try { if (num == 0) { return; } System.out.println("执行递归,num = " + num); recursiveCall(num - 1); } finally { lock.unlock(); // 释放锁 } } public static void main(String[] args) throws InterruptedException { Counter counter = new Counter(); // 创建计数器对象 // 测试递归调用 counter.recursiveCall(10); } } ``` - Condition Condition提供了线程之间的协调机制的接口,可以将它看作是一个更加灵活、更加强大的wait()和notify()机制,通常与Lock接口(比如ReentrantLock)一起使用。它的核心作用体现在两个方面 - **等待/通知机制**:它允许线程等待某个条件成立,或者通知其他线程某个条件已经满足,这与使用Object的wait()和notify()方法相似,但Condition提供了更高的灵活性和更多的控制。 - **多条件协调**:与每个Object只有一个内置的等待/通知机制不同,一个Lock可以对应多个Condition对象,这意味着可以为不同的等待条件创建不同的Condition,从而实现对多个等待线程集合的独立控制。 - 核心API - void await() 使当前线程等待,直到被其他线程通过 signal() 或 signalAll() 方法唤醒,或者线程被中断,或者发生了其他不可预知的情况(如假唤醒)。该方法会在等待之前释放当前线程所持有的锁,在被唤醒后会再次尝试获取锁。 - boolean await(long time, TimeUnit unit) 使当前线程等待指定的时间,效果同上 - void signal()唤醒等待在此 Condition 上的一个线程。如果有多个线程正在等待,则选择其中的一个进行唤醒。被唤醒的线程将从其 await() 调用中返回,并重新尝试获取与此 Condition 关联的锁。 - void signalAll()唤醒等待在此 Condition 上的所有线程。 ```java public class ReentrantLockDemo3 { public static void main(String[] args) { // 创建队列 Queue queue = new Queue(5); //启动生产者线程 new Thread(new Producer(queue)).start(); //启动消费者线程 new Thread(new Customer(queue)).start(); } } /** * 队列封装类 */ class Queue { private Object[] items ; int size = 0; int takeIndex; int putIndex; private ReentrantLock lock; public Condition notEmpty; //消费者线程阻塞唤醒条件,队列为空阻塞,生产者生产完唤醒 public Condition notFull; //生产者线程阻塞唤醒条件,队列满了阻塞,消费者消费完唤醒 public Queue(int capacity){ this.items = new Object[capacity]; lock = new ReentrantLock(); notEmpty = lock.newCondition(); notFull = lock.newCondition(); } public void put(Object value) throws Exception { //加锁 lock.lock(); try { while (size == items.length) // 队列满了让生产者等待 notFull.await(); items[putIndex] = value; if (++putIndex == items.length) putIndex = 0; size++; notEmpty.signal(); // 生产完唤醒消费者 } finally { System.out.println("producer生产:" + value); //解锁 lock.unlock(); } } public Object take() throws Exception { lock.lock(); try { // 队列空了就让消费者等待 while (size == 0) notEmpty.await(); Object value = items[takeIndex]; items[takeIndex] = null; if (++takeIndex == items.length) takeIndex = 0; size--; notFull.signal(); //消费完唤醒生产者生产 return value; } finally { lock.unlock(); } } } /** * 生产者 */ class Producer implements Runnable { private Queue queue; public Producer(Queue queue) { this.queue = queue; } @Override public void run() { try { // 隔1秒轮询生产一次 while (true) { Thread.sleep(1000); queue.put(new Random().nextInt(1000)); } } catch (Exception e) { e.printStackTrace(); } } } /** * 消费者 */ class Customer implements Runnable { private Queue queue; public Customer(Queue queue) { this.queue = queue; } @Override public void run() { try { // 隔2秒轮询消费一次 while (true) { Thread.sleep(2000); System.out.println("consumer消费:" + queue.take()); } } catch (Exception e) { e.printStackTrace(); } } } ``` - 工作原理 基于AQS(AbstractQueuedSynchronizer)+ CAS - AQS - AQS具备特性 - 阻塞等待队列 - 共享/独占 - 公平/非公平 - 可重入 - 允许中断 - 核心结构 - volatile int state ,值为0时表示无锁,大于0表示已加锁,值表示重入次数。除get/set外,还提供了compareAndSetState() - 定义了两种资源访问方式: Exclusive独占和Share共享 - 主要方法 - isHeldExclusively():该线程是否正在独占资源。只有用到condition才需要去实现它。 - tryAcquire(int):独占方式。尝试获取资源,成功则返回true,失败则返回false。 - tryRelease(int):独占方式。尝试释放资源,成功则返回true,失败则返回false。 - tryAcquireShared(int):共享方式。尝试获取资源。负数表示失败;0表示成功,但没有剩余可用资源;正数表示成功,且有剩余资源。 - tryReleaseShared(int):共享方式。尝试释放资源,如果释放后允许唤醒后续等待结点返回true,否则返回false。 - AQS定义的两种队列 > 值为0,初始化状态,表示当前节点在sync队列中,等待着获取锁。CANCELLED,值为1,表示当前的线程被取消; SIGNAL,值为-1,表示当前节点的后继节点包含的线程需要运行,也就是unpark; CONDITION,值为-2,表示当前节点在等待condition,也就是在condition队列中; PROPAGATE,值为-3,表示当前场景下后续的acquireShared能够得以执行; > - 同步等待队列 > AQS当中的同步等待队列也称CLH队列,CLH队列是Craig、Landin、Hagersten三人发明的一种基于双向链表数据结构的队列,是FIFO先进先出线程等待队列,Java中的CLH队列是原CLH队列的一个变种,线程由原自旋机制改为阻塞机制。 > AQS 依赖CLH同步队列来完成同步状态的管理: - 当前线程如果获取同步状态失败时,AQS则会将当前线程已经等待状态等信息构造成一个节点(Node)并将其加入到CLH同步队列,同时会阻塞当前线程 - 当同步状态释放时,会把首节点唤醒(公平锁),使其再次尝试获取同步状态。 - 通过signal或signalAll将条件队列中的节点转移到同步队列。(由条件队列转化为同步队列)  - 条件等待队列 AQS中条件队列是使用单向列表保存的,用nextWaiter来连接: - 调用await方法阻塞线程; - 当前线程存在于同步队列的头结点,调用await方法进行阻塞(从同步队列转化到条件队列)  - 基于AQS实现一把独占锁 ```java public class TulingLock extends AbstractQueuedSynchronizer{ @Override protected boolean tryAcquire(int unused) { //cas 加锁 state=0 if (compareAndSetState(0, 1)) { setExclusiveOwnerThread(Thread.currentThread()); return true; } return false; } @Override protected boolean tryRelease(int unused) { //释放锁 setExclusiveOwnerThread(null); setState(0); return true; } public void lock() { acquire(1); } public boolean tryLock() { return tryAcquire(1); } public void unlock() { release(1); } public boolean isLocked() { return getState() != 0; } } ```   ## Semaphore&CountDownLatch&CyclicBarrier - **Semaphore 信号量** 主要用于在一个时刻允许多个限制数量的线程对共享资源进行并行操作的场景。Semaphore维护了一个计数器,线程可以通过调用acquire()方法来获取Semaphore中的许可证,当计数器为0时,调用acquire()的线程将被阻塞,直到有其他线程释放许可证;线程可以通过调用release()方法来释放Semaphore中的许可证,这会使Semaphore中的计数器增加,从而允许更多的线程访问共享资源。  - 常用API Semaphore默认是非公平的,允许创建实例时指定是否公平public Semaphore(int permits, boolean fair) - acquire() / acquire(int permits) 获取一个/指定数量的许可证,如果获取不到就会一直等待,**直到获取到可用许可证或被其他线程中断** - tryAcquire() / tryAcquire(long timeout, TimeUnit unit) | tryAcquire(int permits) / tryAcquire(int permits, long timeout, TimeUnit unit) 尝试(在指定时间内)获取一个/ 指定数量许可证,获取失败则返回false - release() / release(int permits) 释放一个/指定数量许可证,释放后内部空闲许可证计数器会增加 ```java final Semaphore semaphore = new Semaphore(1, true); // 定义一个线程 new Thread(() -> { // 获取许可证 boolean gotPermit = semaphore.tryAcquire(); // 如果获取成功就休眠5秒的时间 if (gotPermit) { try { System.out.println(Thread.currentThread() + " get one permit."); TimeUnit.SECONDS.sleep(5); } catch (InterruptedException e) { e.printStackTrace(); } finally { // 释放Semaphore的许可证 // 如果有多个不同线程的许可证,需要慎重考虑锁释放 // 可能会出现将其他线程许可证被释放的情况 semaphore.release(); } } }).start(); // 短暂休眠1秒的时间,确保上面的线程能够启动,并且顺利获取许可证 TimeUnit.SECONDS.sleep(1); // 主线程在3秒之内肯定是无法获取许可证的,那么主线程将在阻塞3秒之后返回获取许可证失败 if(semaphore.tryAcquire(3, TimeUnit.SECONDS)){ System.out.println("get the permit"); }else { System.out.println("get the permit failure."); } ``` - 使用场景 - 接口请求限流 - 资源池数量限制 - **CountDownLatch** CountDownLatch(闭锁)是一个同步协助类,可以用于控制一个或多个线程等待多个任务完成后再执行。CountDownLatch 内部维护了一个计数器,该计数器初始值为 N,代表需要等待的线程数目,当一个线程完成了需要等待的任务后,就会调用 countDown() 方法将计数器减 1,当计数器的值为 0 时,等待的线程就会开始执行。  - 常用API CountDownLatch(int count) cout必须≥0,并且初始化实例时,**只能使用一次**,用完之后**不能再复用**。 - countDown()方法,该方法的主要作用是使得构造CountDownLatch指定的count计数器减一。如果此时CountDownLatch中的计数器已经是0,这种情况下如果再次调用countDown()方法,则会被忽略,也就是说count的值最小只能为0。 - await() / await(long timeout, TimeUnit unit)方法会使得当前的调用线程进入阻塞状态,直到count为0,其他线程可以将当前线程中断 - getCount()方法,该方法将返回CountDownLatch当前的计数器数值,该返回值的最小值为0。 ```java public class CountDownLatchDemo2 { public static void main(String[] args) throws Exception { CountDownLatch countDownLatch = new CountDownLatch(5); for (int i = 0; i < 5; i++) { final int index = i; new Thread(() -> { try { Thread.sleep(1000 + ThreadLocalRandom.current().nextInt(2000)); System.out.println("任务" + index +"执行完成"); countDownLatch.countDown(); } catch (InterruptedException e) { e.printStackTrace(); } }).start(); } // 主线程在阻塞,当计数器为0,就唤醒主线程往下执行 countDownLatch.await(); System.out.println("主线程:在所有任务运行完成后,进行结果汇总"); } } ``` - 使用场景 - 并行任务汇总同步:协调多个并行任务,完成后再进行下一步操作 - 资源初始化: 等待多个资源初始化完成再使用 - **CyclicBarrier** CyclicBarrier(循环屏障),是 Java 并发库中的一个同步工具,通过它可以实现让一组线程等待至某个状态(屏障点)之后再全部同时执行。**CyclicBarrier可以被重用**。CyclicBarrier也适合用于某个串行化任务被分拆成若干个并行执行的子任务,当所有的子任务都执行结束之后再继续接下来的工作。 ```java // parties表示屏障拦截的线程数量,每个线程调用 await 方法告诉 CyclicBarrier 我已经到达了屏障,然后当前线程被阻塞。 public CyclicBarrier(int parties) // 用于在线程到达屏障时,优先执行 barrierAction,方便处理更复杂的业务场景(该线程的执行时机是在到达屏障之后再执行) public CyclicBarrier(int parties, Runnable barrierAction) ```  - 常用方法 - await() / await(long timeout, TimeUnit unit) 调用线程进入阻塞状态,等待指定线程数量全部调用await后再进行下一步操作 - reset() 重置循环 ```java public class CyclicBarrierDemo2 { private static int[] getProductsByCategoryId() { // 商品列表编号为从1~10的数字 return IntStream.rangeClosed(1, 10).toArray(); } private static class ProductPrice { private final int prodID; private double price; private ProductPrice(int prodID) { this(prodID, -1); } private ProductPrice(int prodID, double price) { this.prodID = prodID; this.price = price; } int getProdID() { return prodID; } void setPrice(double price) { this.price = price; } @Override public String toString() { return "ProductPrice{" + "prodID=" + prodID + ", price=" + price + '}'; } } public static void main(String[] args) throws InterruptedException { // 根据商品品类获取一组商品ID final int[] products = getProductsByCategoryId(); // 通过转换将商品编号转换为ProductPrice List<ProductPrice> list = Arrays.stream(products).mapToObj(ProductPrice::new).collect(toList()); // 1. 定义CyclicBarrier ,指定parties为子任务数量 final CyclicBarrier barrier = new CyclicBarrier(list.size()); // 2.用于存放线程任务的list final List<Thread> threadList = new ArrayList<>(); list.forEach(pp -> { Thread thread = new Thread(() -> { System.out.println(pp.getProdID() + "开始计算商品价格."); try { TimeUnit.SECONDS.sleep(current().nextInt(10)); if (pp.prodID % 2 == 0) { pp.setPrice(pp.prodID * 0.9D); } else { pp.setPrice(pp.prodID * 0.71D); } System.out.println(pp.getProdID() + "->价格计算完成."); } catch (InterruptedException e) { // ignore exception } finally { try { // 3.在此等待其他子线程到达barrier point barrier.await(); } catch (InterruptedException | BrokenBarrierException e) { } } }); threadList.add(thread); thread.start(); }); // 4. 等待所有子任务线程结束 threadList.forEach(t -> { try { t.join(); } catch (InterruptedException e) { e.printStackTrace(); } }); System.out.println("所有价格计算完成."); list.forEach(System.out::println); } } ``` - **CyclicBarrier 与 CountDownLatch 区别** - CountDownLatch 是一次性的,CyclicBarrier 是可循环利用的 - CoundDownLatch的await方法会等待计数器被count down到0,而执行CyclicBarrier的await方法的线程将会等待其他线程到达barrier point。 - CyclicBarrier内部的计数器count是可被重置的,进而使得CyclicBarrier也可被重复使用,而CoundDownLatch则不能 ## 并发容器 - List类 - CopyOnWriteArrayList 对应非并发容器ArrayList;用于代替Vector、synchronizedList;利用高并发往往是**读多写少**的特性,**对读操作不加锁,对写操作**,先**复制一份新的集合**,在新的集合上面修改,然后将新集合赋值给旧的引用,并通过volatile 保证其可见性。   - 应用场景 - 读多写少的场景 - 不需要实时更新数据的场景 ```java public class CopyOnWriteArrayListDemo { private static CopyOnWriteArrayList<String> copyOnWriteArrayList = new CopyOnWriteArrayList<>(); // 模拟初始化的黑名单数据 static { copyOnWriteArrayList.add("ipAddr0"); copyOnWriteArrayList.add("ipAddr1"); copyOnWriteArrayList.add("ipAddr2"); } public static void main(String[] args) throws InterruptedException { Runnable task = new Runnable() { public void run() { // 模拟接入用时 try { Thread.sleep(new Random().nextInt(5000)); } catch (Exception e) {} String currentIP = "ipAddr" + new Random().nextInt(6); if (copyOnWriteArrayList.contains(currentIP)) { System.out.println(Thread.currentThread().getName() + " IP " + currentIP + "命中黑名单,拒绝接入处理"); return; } System.out.println(Thread.currentThread().getName() + " IP " + currentIP + "接入处理..."); } }; new Thread(task, "请求1").start(); new Thread(task, "请求2").start(); new Thread(task, "请求3").start(); new Thread(new Runnable() { public void run() { // 模拟用时 try { Thread.sleep(new Random().nextInt(2000)); } catch (Exception e) {} String newBlackIP = "ipAddr3"; copyOnWriteArrayList.add(newBlackIP); System.out.println(Thread.currentThread().getName() + " 添加了新的非法IP " + newBlackIP); } }, "IP黑名单更新").start(); Thread.sleep(1000000); } } ``` - 优点 - 读操作不加锁,性能高 - 缺点 - 内存占用问题,每次执行写操作都得拷贝一份,数据量大时对内存压力较大 - 无法保证实时性,写时复制在写操作的时候,内存里会同时有两个对象内存 - Map类 - ConcurrentHashMap 对应非并发容器HashMap;用于代替Hashtable、synchronizedMap,支持复合操作作;JDK8之前采用分段锁,JDK8中采用CAS无锁算法+synchronized。  - JDK1.7 的实现 底层结构用Segments数组+HashEntry数组+链表实现;底层维护一个默认大小为16的Segments数组,每个segment里是一个完整的HashMap,加上一把继承的ReentrantLock,不同线程访问不同的segment时完全不冲突,只有在访问同一segment时才有锁竞争,并发度最高为segment数组的大小。调用size方法时,需要把所有的Segment锁起来计算累加值,慢。  - JDK1.8实现 移除了Segment,锁力度细化到每个槽位,数据结构与HashMap一致,变成数组+链表+红黑树。插入时先用CAS无锁尝试插入到数组位置,冲突了才用synchronized,而且只锁链表或树的头节点,其他线程照样可操作别的bucket。当链表节点数大于8且数组长度大于等于64时转换为红黑树; 当树中节点数小于6时退化成链表,中间留个缓存,避免频繁转换  - ConcurrentSkipListMap 对应非并发容器TreeMap;用于代替synchronizedSortedMap;用跳表替代平衡树,默认按照key升序。适用于需要高并发性能、支持**有序性**和区间查询的场景,能够有效地提高系统的性能和可扩展性。 > 跳表介绍见 redis → zset底层实现 >  ```java public class ConcurrentSkipListMapDemo { public static void main(String[] args) { ConcurrentSkipListMap<Integer, String> map = new ConcurrentSkipListMap<>(); // 添加元素 map.put(1, "a"); map.put(3, "c"); map.put(2, "b"); map.put(4, "d"); // 获取元素 String value1 = map.get(2); System.out.println(value1); // 输出:b // 遍历元素 for (Integer key : map.keySet()) { String value = map.get(key); System.out.println(key + " : " + value); } // 删除元素 String value2 = map.remove(3); System.out.println(value2); // 输出:c } } ``` - Set类 - CopyOnWriteArraySet 对应非并发容器HashSet;用于替代synchronizedSet;基于CopyOnWriteArrayList实现,其唯一的不同是在add时调用的是CopyOnWriteArrayList的addIfAbsent方法,其遍历当前Object数组,如Object数组中已有了当前元素,则直接返回,如果没有则放入Object数组的尾部,并返回。 - ConcurrentSkipListSet 基于ConcurrentSkipListMap实现 - Queue类 - BlockingQueue BlockingQueue提供了线程安全的队列访问方式:当阻塞队列插入数据时,如果队列已满,线程将会阻塞等待直到队列非满;从阻塞队列取数据时,如果队列已空,线程将会阻塞等待直到队列非空。 - 应用场景 1. 线程池 线程池中的任务队列通常是一个阻塞队列。当任务数超过线程池的容量时,新提交的任务将被放入任务队列中等待执行。线程池中的工作线程从任务队列中取出任务进行处理,如果队列为空,则工作线程会被阻塞,直到队列中有新的任务被提交。 2. 生产者-消费者模型 在生产者-消费者模型中,生产者向队列中添加元素,消费者从队列中取出元素进行处理。阻塞队列可以很好地解决生产者和消费者之间的并发问题,避免线程间的竞争和冲突。 3. 消息队列 消息队列使用阻塞队列来存储消息,生产者将消息放入队列中,消费者从队列中取出消息进行处理。消息队列可以实现异步通信,提高系统的吞吐量和响应性能,同时还可以将不同的组件解耦,提高系统的可维护性和可扩展性。 4. 缓存系统 缓存系统使用阻塞队列来存储缓存数据,当缓存数据被更新时,它会被放入队列中,其他线程可以从队列中取出最新的数据进行使用。使用阻塞队列可以避免并发更新缓存数据时的竞争和冲突。 5. 并发任务处理 在并发任务处理中,可以将待处理的任务放入阻塞队列中,多个工作线程可以从队列中取出任务进行处理。使用阻塞队列可以避免多个线程同时处理同一个任务的问题,并且可以将任务的提交和执行解耦,提高系统的可维护性和可扩展性。 - ArrayBlockingQueue ArrayBlockingQueue是最典型的有界阻塞队列,其内部是用数组存储元素的,初始化时需要指定容量大小,利用 ReentrantLock 实现线程安全。ArrayBlockingQueue可以用于实现数据缓存、限流、生产者-消费者模式等各种应用。 在生产者-消费者模型中使用时,如果生产速度和消费速度基本匹配的情况下,使用ArrayBlockingQueue是个不错选择;当如果生产速度远远大于消费速度,则会导致队列填满,大量生产线程被阻塞。 ArrayBlockingQueue使用独占锁ReentrantLock实现线程安全,入队和出队操作使用同一个锁对象,也就是只能有一个线程可以进行入队或者出队操作;这也就意味着生产者和消费者无法并行操作,在高并发场景下会成为性能瓶颈。 - LinkedBlockingQueue LinkedBlockingQueue是一个基于链表实现的阻塞队列,默认情况下,该阻塞队列的大小为Integer.MAX_VALUE,由于这个数值特别大,所以 LinkedBlockingQueue 也被称作**无界队列**,代表它几乎没有界限,队列可以随着元素的添加而动态增长,但是如果没有剩余内存,则队列将抛出OOM错误。所以为了避免队列过大造成机器负载或者内存爆满的情况出现,我们在使用的时候建议手动传一个队列的大小。  > **fail-fast与fail-safe fail-fast(快速失败)** 当使用迭代器对集合进行遍历时,如果有线程对集合中的结构进行修改(比如删除),会立即抛出ConcurrentModificationException,ArrayList、HashMap等采用的是fail-fast。可对操作加锁或用对应的线程安全类替代以解决。 **fail-safe(安全失败)** 对集合结构的修改都会在一个复制集合上进行,不改变原集合内容,因此不会抛出ConcurrentModificationException;例如 CopyOnWriteArrayList、ConcurrentHashMap等采用的是fail-safe。 虽然采用fail-safe的集合类都是线程安全的,但是它们无法保证数据实时性,只能保证数据的最终一致性;另外就是内存占用问题,要通过复制来实现读写分离,因此会占用更多的内存 >
Day23 CAS与synchronized
## CAS(Compare And Swap) > Compare And Swap 比较与交换,CPU硬件层面的一种指令,是非阻塞同步的实现原理。 > CAS指令操作包括三个参数:内存值(内存地址值)V、预期值E、新值N,当CAS指令执行时,当且仅当预期值E和内存值V相同时,才更新内存值为N,否则就不执行更新,无论更新与否都会返回否会返回旧的内存值V  CAS是一种无锁算法,在不使用锁的情况下实现多线程之间的变量同步。在Java中,CAS操作是由Unsafe类提供支持的(构造器私有,需要通过反射访问) ```java public class CASTest { public static void main(String[] args) { Entity entity = new Entity(); Unsafe unsafe = UnsafeFactory.getUnsafe(); long offset = UnsafeFactory.getFieldOffset(unsafe, Entity.class, "x"); boolean successful; // 4个参数分别是:对象实例、字段的内存偏移量、字段期望值、字段新值 successful = unsafe.compareAndSwapInt(entity, offset, 0, 3); System.out.println(successful + "\t" + entity.x); successful = unsafe.compareAndSwapInt(entity, offset, 3, 5); System.out.println(successful + "\t" + entity.x); successful = unsafe.compareAndSwapInt(entity, offset, 3, 8); System.out.println(successful + "\t" + entity.x); } } public class UnsafeFactory { /** * 获取 Unsafe 对象 * @return */ public static Unsafe getUnsafe() { try { Field field = Unsafe.class.getDeclaredField("theUnsafe"); field.setAccessible(true); return (Unsafe) field.get(null); } catch (Exception e) { e.printStackTrace(); } return null; } /** * 获取字段的内存偏移量 * @param unsafe * @param clazz * @param fieldName * @return */ public static long getFieldOffset(Unsafe unsafe, Class clazz, String fieldName) { try { return unsafe.objectFieldOffset(clazz.getDeclaredField(fieldName)); } catch (NoSuchFieldException e) { throw new Error(e); } } } ``` - CAS缺陷 - 自旋开销 失败时循环重试,高并发下CPU占用高 - ABA问题 值从A→B→A CAS误以为没变 使可用AtomicStampedReference解决 构造方法中有两个参数,reference即我们实际存储的变量,stamp是版本,相当于加了乐观锁 ```java @Slf4j public class AtomicStampedReferenceTest { public static void main(String[] args) { // 定义AtomicStampedReference Pair.reference值为1, Pair.stamp为1 AtomicStampedReference atomicStampedReference = new AtomicStampedReference(1,1); new Thread(()->{ int[] stampHolder = new int[1]; int value = (int) atomicStampedReference.get(stampHolder); int stamp = stampHolder[0]; log.debug("Thread1 read value: " + value + ", stamp: " + stamp); // 阻塞1s LockSupport.parkNanos(1000000000L); // Thread1通过CAS修改value值为3 if (atomicStampedReference.compareAndSet(value, 3,stamp,stamp+1)) { log.debug("Thread1 update from " + value + " to 3"); } else { log.debug("Thread1 update fail!"); } },"Thread1").start(); new Thread(()->{ int[] stampHolder = new int[1]; int value = (int)atomicStampedReference.get(stampHolder); int stamp = stampHolder[0]; log.debug("Thread2 read value: " + value+ ", stamp: " + stamp); // Thread2通过CAS修改value值为2 if (atomicStampedReference.compareAndSet(value, 2,stamp,stamp+1)) { log.debug("Thread2 update from " + value + " to 2"); value = (int) atomicStampedReference.get(stampHolder); stamp = stampHolder[0]; log.debug("Thread2 read value: " + value+ ", stamp: " + stamp); // Thread2通过CAS修改value值为1 if (atomicStampedReference.compareAndSet(value, 1,stamp,stamp+1)) { log.debug("Thread2 update from " + value + " to 1"); } } },"Thread2").start(); } } ``` - CAS 在Java中的应用——Atomic原子类 - **基本类型**:AtomicInteger、AtomicLong、AtomicBoolean; ```java public class AtomicIntegerTest { static AtomicInteger sum = new AtomicInteger(0); public static void main(String[] args) { for (int i = 0; i < 10; i++) { Thread thread = new Thread(() -> { for (int j = 0; j < 10000; j++) { // 原子自增 CAS sum.incrementAndGet(); //TODO } }); thread.start(); } try { Thread.sleep(3000); } catch (InterruptedException e) { e.printStackTrace(); } System.out.println(sum.get()); } } ``` - **引用类型**:AtomicReference、AtomicStampedRerence、AtomicMarkableReference; ```java public class AtomicReferenceTest { public static void main( String[] args ) { User user1 = new User("张三", 23); User user2 = new User("李四", 25); User user3 = new User("王五", 20); //初始化为 user1 AtomicReference<User> atomicReference = new AtomicReference<>(); atomicReference.set(user1); //把 user2 赋给 atomicReference atomicReference.compareAndSet(user1, user2); System.out.println(atomicReference.get()); //把 user3 赋给 atomicReference atomicReference.compareAndSet(user1, user3); System.out.println(atomicReference.get()); } } @Data @AllArgsConstructor class User { private String name; private Integer age; } ``` - **数组类型**:AtomicIntegerArray、AtomicLongArray、AtomicReferenceArray ```java public class AtomicIntegerArrayTest { static int[] value = new int[]{ 1, 2, 3, 4, 5 }; static AtomicIntegerArray atomicIntegerArray = new AtomicIntegerArray(value); public static void main(String[] args) throws InterruptedException { //设置索引0的元素为100 atomicIntegerArray.set(0, 100); System.out.println(atomicIntegerArray.get(0)); //以原子更新的方式将数组中索引为1的元素与输入值相加 atomicIntegerArray.getAndAdd(1,5); System.out.println(atomicIntegerArray); } } ``` - **对象属性原子修改器**:AtomicIntegerFieldUpdater、AtomicLongFieldUpdater、AtomicReferenceFieldUpdater **使用约束** - 字段必须是volatile修饰的,确保线程之间共享变量立即可见 - 只能是操作当前实例变量,不能是类的静态变量或者父类的变量 - 一定要是可变的变量,不能用final修饰 - 对于AtomicIntegerFieldUpdater和AtomicLongFieldUpdater只能修改int/long类型的字段,不能修改其包装类型(Integer/Long)。如果要修改包装类型就需要使用AtomicReferenceFieldUpdater。 ```java public class AtomicIntegerFieldUpdaterTest { public static class Candidate { volatile int score = 0; AtomicInteger score2 = new AtomicInteger(); } public static final AtomicIntegerFieldUpdater<Candidate> scoreUpdater = AtomicIntegerFieldUpdater.newUpdater(Candidate.class, "score"); public static AtomicInteger realScore = new AtomicInteger(0); public static void main(String[] args) throws InterruptedException { final Candidate candidate = new Candidate(); Thread[] t = new Thread[10000]; for (int i = 0; i < 10000; i++) { t[i] = new Thread(new Runnable() { @Override public void run() { if (Math.random() > 0.4) { candidate.score2.incrementAndGet(); scoreUpdater.incrementAndGet(candidate); realScore.incrementAndGet(); } } }); t[i].start(); } for (int i = 0; i < 10000; i++) { t[i].join(); } System.out.println("AtomicIntegerFieldUpdater Score=" + candidate.score); System.out.println("AtomicInteger Score=" + candidate.score2.get()); System.out.println("realScore=" + realScore.get()); } } ``` - **原子类型累加器(jdk1.8增加的类)**:DoubleAccumulator、DoubleAdder、LongAccumulator、LongAdder、Striped64 ```java public class LongAdderTest { public static void main(String[] args) { testAtomicLongVSLongAdder(10, 10000); System.out.println("=================="); testAtomicLongVSLongAdder(10, 200000); System.out.println("=================="); testAtomicLongVSLongAdder(100, 200000); } static void testAtomicLongVSLongAdder(final int threadCount, final int times) { try { long start = System.currentTimeMillis(); testLongAdder(threadCount, times); long end = System.currentTimeMillis() - start; System.out.println("条件>>>>>>线程数:" + threadCount + ", 单线程操作计数" + times); System.out.println("结果>>>>>>LongAdder方式增加计数" + (threadCount * times) + "次,共计耗时:" + end); long start2 = System.currentTimeMillis(); testAtomicLong(threadCount, times); long end2 = System.currentTimeMillis() - start2; System.out.println("条件>>>>>>线程数:" + threadCount + ", 单线程操作计数" + times); System.out.println("结果>>>>>>AtomicLong方式增加计数" + (threadCount * times) + "次,共计耗时:" + end2); } catch (InterruptedException e) { e.printStackTrace(); } } static void testAtomicLong(final int threadCount, final int times) throws InterruptedException { CountDownLatch countDownLatch = new CountDownLatch(threadCount); AtomicLong atomicLong = new AtomicLong(); for (int i = 0; i < threadCount; i++) { new Thread(new Runnable() { @Override public void run() { for (int j = 0; j < times; j++) { atomicLong.incrementAndGet(); } countDownLatch.countDown(); } }, "my-thread" + i).start(); } countDownLatch.await(); } static void testLongAdder(final int threadCount, final int times) throws InterruptedException { CountDownLatch countDownLatch = new CountDownLatch(threadCount); LongAdder longAdder = new LongAdder(); for (int i = 0; i < threadCount; i++) { new Thread(new Runnable() { @Override public void run() { for (int j = 0; j < times; j++) { longAdder.add(1); } countDownLatch.countDown(); } }, "my-thread" + i).start(); } countDownLatch.await(); } } ``` ## 并发锁 一段代码块内如果存在对共享资源的多线程读写操作,称这段代码块为临界区,其共享资源为**临界资源**。多个线程在临界区内执行,由于代码的执行序列不同而导致结果无法预测,称之为发生了**竞态条件**,为了避免临界区的竞态条件发生,可使用两种解决方案: - 阻塞式方案: synchronized、Lock - 非阻塞式方案:原子变量 - **synchronized** synchronized 同步块是 Java 提供的一种原子性内置锁,Java 中的每个对象都可以把它当作一个同步锁来使用,这些 Java 内置的使用者看不到的锁被称为内置锁,也叫作监视器锁。 - 加锁方式 - 修饰实例方法 锁住的是该类的实例对象 - 修饰静态方法 锁住的是类对象 --- 行内代码块 --- - 修饰实例对象 synchronized(this) 锁住实例对象 - 修饰class synchronized(Demo.class) 锁住类对象 - 修饰任意对象Object 锁住实例对象Object  - 实现原理 - synchronized是JVM内置锁,基于**Monitor**机制实现,依赖底层操作系统的Mutex(互斥量),内部会有等待队列(cxq 和 EntryList)和条件等待队列(waitSet)来存放相应阻塞的线程。未竞争到锁的线程存储到等待队列中,获得锁的线程调用 wait 后便存放在条件等待队列中,解锁和 notify 都会唤醒相应队列中的等待线程来争抢锁。它是一个重量级锁,性能较低 同步方法是通过方法中的access_flags中设置ACC_SYNCHRONIZED标志来实现;同步代码块是通过monitorenter和monitorexit来实现。  > 在获取锁时,是将当前线程插入到cxq的头部,而释放锁时,默认策略(QMode=0)是:如果EntryList为空,则将cxq中的元素按原有顺序插入到EntryList,并唤醒第一个线程,也就是当EntryList为空时,是后来的线程先获取锁。_EntryList不为空,直接从_EntryList中唤醒线程。 为什么会有_cxq 和 _EntryList 两个列表来放线程? 因为会有多个线程会同时竞争锁,所以搞了个 _cxq 这个单向链表基于 CAS 来 hold 住这些并发,然后另外搞一个 _EntryList 这个双向链表,来在每次唤醒的时候搬迁一些线程节点,降低 _cxq 的尾部竞争。 > - Monitor机制 Monitor,直译为“监视器”,而操作系统领域一般翻译为“管程”。管程是指管理共享变量以及对共享变量操作的过程,让它们支持并发,java内置的管程里只有一个条件变量  java.lang.Object 类定义了 wait(),notify(),notifyAll() 方法,这些方法的具体实现,依赖于 ObjectMonitor 实现,这是 JVM 内部基于 C++ 实现的一套机制 - **四种锁状态** - 重量级锁 Heavyweight Locking > synchronized 由于阻塞和唤醒依赖于底层的操作系统实现,系统调用存在用户态与内核态之间的切换,所以有较高的开销,因此称之为重量级锁。在JDK6+优化后,JVM会根据竞争情况动态升级锁状态(锁只能升级不能降级),synchronized就不再是“重量级锁”的代名词,并发性能基本与Lock持平 > **重量级锁的优化策略** - 锁粗化 Lock Coarsening 就是将多个连续的锁扩展为一个更大范围的锁。 粒度越细,加锁和解锁的开销越多,粗化可以减少开销从而提高效率 ```java StringBuffer buffer = new StringBuffer(); // 线程安全类,append操作会加锁,JVM识别到如下操作时会合并成一把范围更大的锁 buffer.append("aaa").append(" bbb").append(" ccc"); ``` - 锁消除 Lock Elimination 当一个数据仅在一个线程中使用,或者说这个数据的作用域仅限于一个线程时,这个线程对该数据的所有操作都不需要加锁。 去除不必要的锁竞争,只有自己用,加锁反而增加开销,浪费性能。 ```java public class LockEliminationTest { /** * 锁消除 * -XX:+EliminateLocks 开启锁消除(jdk8默认开启)//执行时间2601 ms * -XX:-EliminateLocks 关闭锁消除 // 执行时间4688 ms */ public void append(String str1, String str2) { StringBuffer stringBuffer = new StringBuffer(); stringBuffer.append(str1).append(str2); } public static void main(String[] args) throws InterruptedException { LockEliminationTest demo = new LockEliminationTest(); long start = System.currentTimeMillis(); for (int i = 0; i < 100000000; i++) { demo.append("aaa", "bbb"); } long end = System.currentTimeMillis(); System.out.println("执行时间:" + (end - start) + " ms"); } } ``` - CAS自适应自旋 Adaptive Spinning 自旋的目的是为了减少线程挂起的次数,尽量避免直接挂起线程(挂起操作涉及系统调用,存在用户态和内核态切换,这才是重量级锁最大的开销) - 轻量级锁 Lightweight Locking 多个线程都是在不同的时间段来请求同一把锁,此时不存在竞争,根本就用不需要阻塞线程,连 monitor 对象都不需要。 线程在自己的栈帧中创建一个Lock Record,把对象的Mark Word拷贝进去,然后CAS把对象头指向这个Lock Record。成功就拿到锁,不成功就膨胀成重量级锁 - 偏向锁 Biased Locking(JDK15+默认禁用,都是有并发的)始终只有一个线程访问,不存在其他线程竞争,此时CAS也不需要,偏向锁可以消除锁重入(CAS)的开销 - 无锁 
Day22 CompletableFuture与ThreadLocal
## CompletableFuture **CompletableFuture是Future接口的扩展和增强**。CompletableFuture实现了Future接口,并在此基础上进行了丰富地扩展,**实现了对任务的编排能力。 优势:**不需要调用get()阻塞等待结果,实现真正的非阻塞;强大的链式组合能力;完善的异常处理机制; - 应用场景 **描述依赖关系:** 1. thenApply() 把前面异步任务的结果,交给后面的Function 2. thenCompose()用来连接两个有依赖关系的任务,结果由第二个任务返回 **描述and聚合关系:** 1. thenCombine:任务合并,有返回值 2. thenAccepetBoth:两个任务执行完成后,将结果交给thenAccepetBoth消耗,无返回值。 3. runAfterBoth:两个任务都执行完成后,执行下一步操作(Runnable)。 **描述or聚合关系:** 1. applyToEither:两个任务谁执行的快,就使用那一个结果,有返回值。 2. acceptEither: 两个任务谁执行的快,就消耗那一个结果,无返回值。 3. runAfterEither: 任意一个任务执行完成,进行下一步操作(Runnable)。 **并行执行:** CompletableFuture类自己也提供了anyOf()和allOf()用于支持多个CompletableFuture并行执行 - 创建异步操作 有四个静态方法 - runAsync 方法以Runnable函数式接口类型为参数,没有返回结果,supplyAsync 方法Supplier函数式接口类型为参数,返回结果类型为U;Supplier 接口的 get() 方法是有返回值的(**会阻塞**) - 没有指定Executor的方法会使用ForkJoinPool.commonPool() 作为它的线程池执行异步代码。如果指定线程池,则使用指定的线程池运行。ForkJoinPool默认创建的线程数是 CPU 的核数(也可以通过 JVM option:-D java.util.concurrent.ForkJoinPool.common.parallelism 来设置) ```java public static CompletableFuture<Void> runAsync(Runnable runnable) public static CompletableFuture<Void> runAsync(Runnable runnable, Executor executor) public static <U> CompletableFuture<U> supplyAsync(Supplier<U> supplier) public static <U> CompletableFuture<U> supplyAsync(Supplier<U> supplier, Executor executor) Runnable runnable = () -> System.out.println("执行无返回结果的异步任务"); CompletableFuture.runAsync(runnable); CompletableFuture<String> future = CompletableFuture.supplyAsync(() -> { System.out.println("执行有返回值的异步任务"); try { Thread.sleep(5000); } catch (InterruptedException e) { e.printStackTrace(); } return "Hello World"; }); String result = future.get(); System.out.println(result); ``` - 结果处理 方法不以Async结尾,意味着Action使用相同的线程执行,而Async可能会使用其它的线程去执行(如果使用相同的线程池,也可能会被同一个线程选中执行)。 ```java public CompletableFuture<T> whenComplete(BiConsumer<? super T,? super Throwable> action) public CompletableFuture<T> whenCompleteAsync(BiConsumer<? super T,? super Throwable> action) public CompletableFuture<T> whenCompleteAsync(BiConsumer<? super T,? super Throwable> action, Executor executor) public CompletableFuture<T> exceptionally(Function<Throwable,? extends T> fn) CompletableFuture.supplyAsync(() -> { try { TimeUnit.SECONDS.sleep(1); } catch (InterruptedException e) { } if (new Random().nextInt(10) % 2 == 0) { int i = 12 / 0; } System.out.println("执行结束!"); return "test"; }) .whenComplete((t, action) -> System.out.println(t + " 执行完成!")) .exceptionally(new Function<Throwable, String>() { @Override public String apply(Throwable t) { System.out.println("执行失败:" + t.getMessage()); return "异常xxxx"; } }).join(); ``` - 结果转换 - thenApply **接收一个普通函数**作为参数,使用该函数处理上一个CompletableFuture 调用的结果,并返回一个具有处理结果的Future对象,把前一个任务的结果T变成U后返回,不产生新的异步任务。 ```java public <U> CompletableFuture<U> thenApply(Function<? super T,? extends U> fn) public <U> CompletableFuture<U> thenApplyAsync(Function<? super T,? extends U> fn) public <U> CompletableFuture<U> thenApplyAsync(Function<? super T,? extends U> fn, Executor executor) CompletableFuture<Integer> future = CompletableFuture.supplyAsync(() -> { int result = 100; System.out.println("一阶段:" + result); return result; }).thenApply(number -> { int result = number * 3; System.out.println("二阶段:" + result); return result; }); System.out.println("最终结果:" + future.get()); ``` - thenCompose **接收一个返回 CompletableFuture 实例的函数**作为参数,该函数的参数是先前计算步骤的结果,把前一个任务的结果T作为参数传到**新的**异步任务里,得到结果后再拆掉一层包装返回(≈flatMap)。 ```java public <U> CompletableFuture<U> thenCompose(Function<? super T, ? extends CompletionStage<U>> fn); public <U> CompletableFuture<U> thenComposeAsync(Function<? super T, ? extends CompletionStage<U>> fn) ; public <U> CompletableFuture<U> thenComposeAsync(Function<? super T, ? extends CompletionStage<U>> fn, Executor executor) ; CompletableFuture<Integer> future = CompletableFuture .supplyAsync(() -> { int number = new Random().nextInt(30); System.out.println("第一阶段:" + number); return number; }) .thenCompose(param -> CompletableFuture.supplyAsync(() -> { int number = param * 2; System.out.println("第二阶段:" + number); return number; })); System.out.println("最终结果: " + future.get()); ``` - 结果消费 - thenAccept 对单个结果进行消费 ```java public CompletionStage<Void> thenAccept(Consumer<? super T> action); public CompletionStage<Void> thenAcceptAsync(Consumer<? super T> action); public CompletionStage<Void> thenAcceptAsync(Consumer<? super T> action,Executor executor); CompletableFuture<Void> future = CompletableFuture .supplyAsync(() -> { int number = new Random().nextInt(10); System.out.println("第一阶段:" + number); return number; }).thenAccept(number -> System.out.println("第二阶段:" + number * 5)); System.out.println("最终结果:" + future.get()); ``` - thenAcceptBoth 对两个结果进行消费 ```java public <U> CompletionStage<Void> thenAcceptBoth(CompletionStage<? extends U> other,BiConsumer<? super T, ? super U> action); public <U> CompletionStage<Void> thenAcceptBothAsync(CompletionStage<? extends U> other,BiConsumer<? super T, ? super U> action); public <U> CompletionStage<Void> thenAcceptBothAsync(CompletionStage<? extends U> other,BiConsumer<? super T, ? super U> action, Executor executor); CompletableFuture<Integer> futrue1 = CompletableFuture.supplyAsync(new Supplier<Integer>() { @Override public Integer get() { int number = new Random().nextInt(3) + 1; try { TimeUnit.SECONDS.sleep(number); } catch (InterruptedException e) { e.printStackTrace(); } System.out.println("第一阶段:" + number); return number; } }); CompletableFuture<Integer> future2 = CompletableFuture.supplyAsync(new Supplier<Integer>() { @Override public Integer get() { int number = new Random().nextInt(3) + 1; try { TimeUnit.SECONDS.sleep(number); } catch (InterruptedException e) { e.printStackTrace(); } System.out.println("第二阶段:" + number); return number; } }); futrue1.thenAcceptBoth(future2, new BiConsumer<Integer, Integer>() { @Override public void accept(Integer x, Integer y) { System.out.println("最终结果:" + (x + y)); } }).join(); ``` - thenRun 不关心结果,只处理后续动作,并拿不到结果值 ```java public CompletionStage<Void> thenRun(Runnable action); public CompletionStage<Void> thenRunAsync(Runnable action); public CompletionStage<Void> thenRunAsync(Runnable action,Executor executor); CompletableFuture<Void> future = CompletableFuture.supplyAsync(() -> { int number = new Random().nextInt(10); System.out.println("第一阶段:" + number); return number; }).thenRun(() -> // runnable不使用上一个任务计算的结果 System.out.println("thenRun 执行")); System.out.println("最终结果:" + future.get()); ``` - 结果组合 - thenCombine 合并两个线程任务的结果进行处理,返回一个CompletableFuture ```java public <U,V> CompletionStage<V> thenCombine(CompletionStage<? extends U> other,BiFunction<? super T,? super U,? extends V> fn); public <U,V> CompletionStage<V> thenCombineAsync(CompletionStage<? extends U> other,BiFunction<? super T,? super U,? extends V> fn); public <U,V> CompletionStage<V> thenCombineAsync(CompletionStage<? extends U> other,BiFunction<? super T,? super U,? extends V> fn,Executor executor); CompletableFuture<Integer> future1 = CompletableFuture .supplyAsync(new Supplier<Integer>() { @Override public Integer get() { int number = new Random().nextInt(10); System.out.println("第一阶段:" + number); return number; } }); CompletableFuture<Integer> future2 = CompletableFuture .supplyAsync(new Supplier<Integer>() { @Override public Integer get() { int number = new Random().nextInt(10); System.out.println("第二阶段:" + number); return number; } }); CompletableFuture<Integer> result = future1 .thenCombine(future2, new BiFunction<Integer, Integer, Integer>() { @Override public Integer apply(Integer x, Integer y) { return x + y; } }); System.out.println("最终结果:" + result.get()); ``` - 任务交互 将两个线程任务获取结果的速度相比较,按一定的规则进行下一步处理。 - applyToEither 哪个任务执行得快就拿哪个的执行结果进行下一步操作 ```java public <U> CompletionStage<U> applyToEither(CompletionStage<? extends T> other,Function<? super T, U> fn); public <U> CompletionStage<U> applyToEitherAsync(CompletionStage<? extends T> other,Function<? super T, U> fn); public <U> CompletionStage<U> applyToEitherAsync(CompletionStage<? extends T> other,Function<? super T, U> fn,Executor executor); CompletableFuture<Integer> future1 = CompletableFuture .supplyAsync(new Supplier<Integer>() { @Override public Integer get() { int number = new Random().nextInt(10); System.out.println("第一阶段start:" + number); try { TimeUnit.SECONDS.sleep(number); } catch (InterruptedException e) { e.printStackTrace(); } System.out.println("第一阶段end:" + number); return number; } }); CompletableFuture<Integer> future2 = CompletableFuture .supplyAsync(new Supplier<Integer>() { @Override public Integer get() { int number = new Random().nextInt(10); System.out.println("第二阶段start:" + number); try { TimeUnit.SECONDS.sleep(number); } catch (InterruptedException e) { e.printStackTrace(); } System.out.println("第二阶段end:" + number); return number; } }); future1.applyToEither(future2, new Function<Integer, Integer>() { @Override public Integer apply(Integer number) { System.out.println("最快结果:" + number); return number * 2; } }).join(); ``` - acceptEither 哪个执行得快就用哪个的结果做副作用 ```java public CompletionStage<Void> acceptEither(CompletionStage<? extends T> other,Consumer<? super T> action); public CompletionStage<Void> acceptEitherAsync(CompletionStage<? extends T> other,Consumer<? super T> action); public CompletionStage<Void> acceptEitherAsync(CompletionStage<? extends T> other,Consumer<? super T> action,Executor executor); ``` - runAfterEither 哪个任务执行完成就进行下一步操作,不关心运行结果 ```java public CompletionStage<Void> runAfterEither(CompletionStage<?> other,Runnable action); public CompletionStage<Void> runAfterEitherAsync(CompletionStage<?> other,Runnable action); public CompletionStage<Void> runAfterEitherAsync(CompletionStage<?> other,Runnable action,Executor executor); ``` - runAfterBoth 两个任务全部完成才进行下一步操作,不关心运行结果 ```java public CompletionStage<Void> runAfterBoth(CompletionStage<?> other,Runnable action); public CompletionStage<Void> runAfterBothAsync(CompletionStage<?> other,Runnable action); public CompletionStage<Void> runAfterBothAsync(CompletionStage<?> other,Runnable action,Executor executor); ``` - CompletableFuture.anyOf 任意一个先完成,就返回这个任务的CompletableFuture - CompletableFuture.allOf 所有任务都完成才返回,返回CompletableFuture<Void> ```java public static CompletableFuture<Void> allOf(CompletableFuture<?>... cfs) CompletableFuture<String> future1 = CompletableFuture .supplyAsync(() -> { try { TimeUnit.SECONDS.sleep(2); } catch (InterruptedException e) { e.printStackTrace(); } System.out.println("future1完成!"); return "future1完成!"; }); CompletableFuture<String> future2 = CompletableFuture .supplyAsync(() -> { System.out.println("future2完成!"); return "future2完成!"; }); CompletableFuture<Void> combindFuture = CompletableFuture .allOf(future1, future2); try { combindFuture.get(); } catch (InterruptedException e) { e.printStackTrace(); } catch (ExecutionException e) { e.printStackTrace(); } System.out.println("future1: " + future1.isDone() + ",future2: " + future2.isDone()); ``` ## ThreadLocal - 为什么需要ThreadLocal? ThreadLocal类用来提供线程内部的局部变量,解决的核心问题是线程隔离,让每个线程都有自己独立的副本,不需要加锁也不会有锁竞争 - 线程安全 不存在锁竞争 - 数据传递 每个线程都有一份副本,共享的数据直接从副本取(例如spring的事物管理) - 线程隔离 每个线程变量都独立 - 常用方法 ThreadLocal() 创建Thread Local对象 public void set( T value) 设置当前线程局部变量 public T get() 获取当前线程绑定的局部变量 public void remove() 移除当前线程绑定的局部变量 ```java public class MyThreadLocal { static class WithoutThreadLocal { private String content; public String getContent() { return content; } public void setContent(String content) { this.content = content; } } private static ThreadLocal<String> threadLocal = new ThreadLocal<>(); static class WithThreadLocal { private String content; private String getContent() { return threadLocal.get(); } private void setContent(String content) { threadLocal.set(content); } } public static void main(String[] args) throws InterruptedException { withoutThreadLocal(); Thread.sleep(100); System.out.println("---------------------------"); withThreadLocal(); } private static void withoutThreadLocal() { WithoutThreadLocal demo = new WithoutThreadLocal(); for (int i = 0; i < 5; i++) { Thread thread = new Thread(() -> { demo.setContent(Thread.currentThread().getName() + "的数据"); System.out.println(Thread.currentThread().getName() + "--->" + demo.getContent()); }); thread.setName("线程" + i); thread.start(); } } private static void withThreadLocal() { WithThreadLocal demo = new WithThreadLocal(); for (int i = 0; i < 5; i++) { Thread thread = new Thread(() -> { demo.setContent(Thread.currentThread().getName() + "的数据"); System.out.println(Thread.currentThread().getName() + "--->" + demo.getContent()); }); thread.setName("线程" + i); thread.start(); } } } ```  - **ThreadLocal与synchronized的区别** - ThreadLocal 以空间换时间,每个线程都有一份副本,互不干扰;侧重于让多线程中每个线程之间的数据都相互隔离 - synchronized 以时间换空间,只有一份变量,谁先抢到就让谁先访问,访问完了才能进行下一轮抢占;侧重于多个线程之间访问资源的同步 - ThreadLocal的结构 - 每个Thread里有一个类型为ThreadLocalMap的threadLocals字段,这个Map的key是ThreadLocal对象本身,value是线程本地变量副本值。  调用get()方法的时候,先拿到当前的ThreadLocalMap,在用this作为key去查,所以每次取值都是拿到自己的,完全隔离其他线程。 ```java public T get() { Thread t = Thread.currentThread(); ThreadLocalMap map = getMap(t); // 拿当前线程的 map if (map != null) { ThreadLocalMap.Entry e = map.getEntry(this); // 用 this 当 key if (e != null) { return (T) e.value; } } return setInitialValue(); } ``` - 内存泄漏问题 ThreadLocalMap的Entry继承了WeakReference,key是弱引用,如果ThreadLocal没有被强引用,GC会把key进行回收,但实际上value还在,此时Entry就变成了key为nul,value还占用着内存的“脏数据”。如果是在线程池,线程生命周期很长,就容易造成内存泄漏了,所以,要真正解决内存泄漏问题,每次使用完记得**调用一下remove方法**
Day 21 线程知识补充
- 线程中断机制 其他线程通过调用某个正在执行线程A的**interrupt()**方法对其进行安全中断操作,调用后不代表A会立即停止自己的工作,它也可以拒绝中断请求,通过检车自身的中断标志位是否被设置为true来进行响应。 线程通过isInterrupted()方法或者Thread.interrupted()判断是否被中断,后者会同时将中断标识位改为false 如果线程处于阻塞状态(sleep(),join(),obj.wait()),在线程检查发现中断标识为true时,会抛出InterruptedException异常,并且在抛出异常后会立即将线程的中断标示位清除,重新设置为false。(**死锁线程无法被中断**) - Java线程模型 - 线程调度机制 - 协同式线程调度**Cooperative Threads-Scheduling**: 线程执行由线程本身控制,线程把自己的工作执行完成之后,主动通知系统切换到另一个线程上。 好处:实现简单,没有线程同步问题 坏处:如果某一个线程出了问题,其他线程就没法执行,会一直阻塞 - 抢占式线程调度**Preemptive Threads-Scheduling**: 由操作系统控制线程中断,按策略分配CPU时间片,线程无法独占。(Java使用) 好处:线程执行可控,单个线程出问题不会导致整个系统瘫痪,高优先级任务可以及时抢占CPU 坏处:上下文切换开销大,执行顺序不确定,实现复杂 - 线程的实现 - 内核线程(1:1)实现 内核线程是直接由操作系统支持的线程,由内核控制线程切换,通过操作调度器对线程进行调度,并负责将线程的任务映射到各个处理器上。 由于内核线程的支持,每个线程都是一个独立的调度单元,即使某个线程被阻塞,也不影响整个进程工作,后续相关的调度操作系统也会处理好 局限性:由于是基于内核线程实现,所以各种线程操作都需要经过操作系统,在用户态和内核态之间来回切换代价较高; - 用户线程(1:N)实现 严格意义上的用户线程是完全建立在用户空间的线程库上,系统内核感知不到用户线程的存在及实现,创建、同步、销毁和调度都不需要内核参与。 用户线程的优势在于不需要系统内核参与,消耗低,操作快速;但劣势也在此,所有线程操作都需要用户程序自己处理。 - 混合(N:M)实现 即存在用户线程,也存在内核线程,集两者之所长 > Java在JDK2之前是用户线程实现,3之后改成了内核线程实现,全权交给了系统进行调度,JVM无法干涉,所以有时候Java设置的线程优先级无法准备的和操作系统中的线程优先级一一对应 > - 虚拟线程(协程) JDK21推出的革命性技术,是JVM管理的轻量级线程,旨在解决传统线程内存开销大、上下文切换慢、受限于内存的并发瓶颈等问题,提高并发能力的同时无需消耗更多资源。虚拟线程在 sleep()、read()、accept() 等阻塞操作时会自动挂起,不占用 OS 线程,当 I/O 完成后,JVM 自动将其调度回某个 Carrier Thread 继续执行。 ```java //Exception in thread "main" java.lang.OutOfMemoryError: unable to create native thread private static void createThread() { for (int i = 0; i < 10_000; i++) { new Thread(() -> { // 同样的 I/O 操作 try { Thread.sleep(1000); } catch (Exception e) { } System.out.println("Done"); }).start(); } } private static void createVirtualThread() { // Java 21+ try (var executor = Executors.newVirtualThreadPerTaskExecutor()) { for (int i = 0; i < 10_000; i++) { executor.submit(() -> { // 同样的 I/O 操作 try { Thread.sleep(1000); } catch (Exception e) { } System.out.println("Done"); return null; }); } } // 自动等待所有任务完成 } ``` - 使用场景 适合I/O密集型任务 如Web服务器、数据库查询、外部调用、文件读写 不适合CPU密集型任务 视频编码、科学计算等 - 使用注意事项 - 不需要池化技术,用完即弃,创建成本低 - 避免使用synchronized锁,会导致Carrier Thread阻塞,需改用reentrantLock或无锁设计 - 慎用ThreadLocal,虚拟线程数量多,容易导致内存泄露 - 线程通信 - volatile 轻量通信,加上volatile关键字,保证不同的线程对这个变量操作时的可见性,但无法保证线程安全 - 等待/通知机制 - Object.wait() 调用该方法的线程进入 WAITING状态,只有等待另外线程的通知或被中断才会返回.需要注意,调用wait()方法后,会释放对象的锁 - Object.notify() 通知一个在对象上等待的线程,使其从wait方法返回,而返回的前提是该线程获取到了对象的锁,没有获得锁的线程重新进入WAITING状态。 - Objecct.notifyAll() 通知所有等待在该对象上的线程。尽可能用notifyAll(),谨慎使用notify(),因为notify()只会唤醒一个线程,我们无法确保被唤醒的这个线程一定就是我们需要唤醒的线程。 ```java public class WaitDemo { public static void main(String[] args) throws InterruptedException { Object locker = new Object(); Thread t1 = new Thread(() -> { try { System.out.println("wait开始"); synchronized (locker) { locker.wait(); } System.out.println("wait结束"); } catch (InterruptedException e) { e.printStackTrace(); } }); t1.start(); //保证t1先启动,wait()先执行 Thread.sleep(1000); Thread t2 = new Thread(() -> { synchronized (locker) { System.out.println("notify开始"); locker.notifyAll(); System.out.println("notify结束"); } }); t2.start(); } } ``` - LoclSupport 是JDK中用来实现线程阻塞和唤醒的工具,线程调用park则等待“许可”,调用unpark则为指定线程提供“许可”。 ```java public class LockSupportDemo { public static void main(String[] args) throws InterruptedException { Thread parkThread = new Thread(new Runnable() { @Override public void run() { System.out.println("ParkThread开始执行"); // 当没有『许可』时,当前线程暂停运行;有『许可』时,用掉这个『许可』,当前线程恢复运行 LockSupport.park(); System.out.println("ParkThread执行完成"); } }); parkThread.start(); Thread.sleep(1000); System.out.println("唤醒parkThread"); // 给线程 parkThread 发放『许可』(多次连续调用 unpark 只会发放一个『许可』) LockSupport.unpark(parkThread); } } ``` - Callable&Future&FutureTask - 背景:直接继承Thread或者实现Runnable接口都可以创建线程,但是这两种方法都没有返回值,也就不能获取执行完的结果。因此java1.5提供了Callable接口来实现这一场景,而Future和FutureTask就可以和Callable接口配合起来使用。 ```java @FunctionalInterface public interface Runnable { public abstract void run(); } @FunctionalInterface public interface Callable<V> { V call() throws Exception; } ``` - Runnable 的缺陷: - 不能返回一个返回值 - 不能抛出 checked Exception Callable的call方法可以有返回值,可以声明抛出异常。和 Callable 配合的有一个 Future 类,通过 Future 可以了解任务执行情况,或者取消任务的执行,还可获取任务执行的结果。 ```java new Thread(new Runnable() { @Override public void run() { System.out.println("通过Runnable方式执行任务"); } }).start(); FutureTask task = new FutureTask(new Callable() { @Override public Object call() throws Exception { System.out.println("通过Callable方式执行任务"); Thread.sleep(3000); return "返回任务结果"; } }); new Thread(task).start(); System.out.println(task.get()); ``` - Future 的API **Future就是对于具体的Runnable或者Callable任务的执行结果进行取消、查询是否完成、获取结果。必要时可以通过get方法获取执行结果,该方法会阻塞直到任务返回结果。** - boolean cancel (boolean mayInterruptIfRunning) 取消任务的执行。参数指定是否立即中断任务执行,或者等等任务结束 - boolean isCancelled () 任务是否已经取消,任务正常完成前将其取消,则返回 true - boolean isDone () 任务是否已经完成。需要注意的是如果任务正常终止、异常或取消,都将返回true - V get () throws InterruptedException, ExecutionException 等待任务执行结束,然后获得V类型的结果。InterruptedException 线程被中断异常, ExecutionException任务执行异常,如果任务被取消,还会抛出CancellationException - V get (long timeout, TimeUnit unit) throws InterruptedException, ExecutionException, TimeoutException 同上面的get功能一样,多了设置超时时间。参数timeout指定超时时间,uint指定时间的单位,在枚举类TimeUnit中有相关的定义。如果计算超时,将抛出TimeoutException - FutureTask FutureTask是Future和Runnable的实现,该对象相当于是消费者和生产者的桥梁,消费者通过 FutureTask 存储任务的处理结果,更新任务的状态:未开始、正在处理、已完成等。而生产者拿到的 FutureTask 被转型为 Future 接口,可以阻塞式获取任务的处理结果,非阻塞式获取任务处理状态 ```java public class FutureTaskDemo { public static void main(String[] args) throws ExecutionException, InterruptedException { Task task = new Task(); //构建futureTask FutureTask<Integer> futureTask = new FutureTask<>(task); //作为Runnable入参 new Thread(futureTask).start(); System.out.println("task运行结果:"+futureTask.get()); } static class Task implements Callable<Integer> { @Override public Integer call() throws Exception { System.out.println("子线程正在计算"); int sum = 0; for (int i = 0; i < 100; i++) { sum += i; } return sum; } } } ``` - Future的局限性 - 并发执行多任务时,只能用get()方法获取结果,并且是阻塞的 - 无法组合多个任务进行链式调用 - 没有异常处理能力,每个get()调用都要手动catch
偏向锁在什么条件下会升级为轻量级锁?
### 问题描述 synchronized锁升级时,偏向锁在什么条件下会升级为轻量级锁? ### 背景信息 Java版本为java11。 ### 代码 ```(java) import org.openjdk.jol.info.ClassLayout; public class BiasedLockThreadIdCheck { static final Object lock = new Object(); public static void main(String[] args) throws Exception { System.out.println("Before any lock:"); System.out.println(ClassLayout.parseInstance(lock).toPrintable()); Thread t1 = new Thread(() -> { System.out.println("T1 ID: " + Thread.currentThread().getId()); synchronized (lock) { System.out.println("T1 获得锁后:"); System.out.println(ClassLayout.parseInstance(lock).toPrintable()); } }); t1.start(); t1.join(); // 注释下列代码前后,t2获得的锁类型不同! // Thread.sleep(5000); Thread t2 = new Thread(() -> { System.out.println("T2 ID: " + Thread.currentThread().getId()); synchronized (lock) { System.out.println("T2 获得锁后:"); System.out.println(ClassLayout.parseInstance(lock).toPrintable()); } }); t2.start(); t2.join(); } } ``` ### 结果 中间没有sleep,t2获取到的是偏向锁,但是偏向的仍然是线程t1: ``` Before any lock: # WARNING: Unable to get Instrumentation. Dynamic Attach failed. You may add this JAR as -javaagent manually, or supply -Djdk.attach.allowAttachSelf java.lang.Object object internals: OFF SZ TYPE DESCRIPTION VALUE 0 8 (object header: mark) 0x0000000000000005 (biasable; age: 0) 8 4 (object header: class) 0x00001000 12 4 (object alignment gap) Instance size: 16 bytes Space losses: 0 bytes internal + 4 bytes external = 4 bytes total T1 ID: 27 T1 获得锁后: java.lang.Object object internals: OFF SZ TYPE DESCRIPTION VALUE 0 8 (object header: mark) 0x000001fea4a2e005 (biased: 0x000000007fa928b8; epoch: 0; age: 0) 8 4 (object header: class) 0x00001000 12 4 (object alignment gap) Instance size: 16 bytes Space losses: 0 bytes internal + 4 bytes external = 4 bytes total T2 ID: 28 T2 获得锁后: java.lang.Object object internals: OFF SZ TYPE DESCRIPTION VALUE 0 8 (object header: mark) 0x000001fea4a2e005 (biased: 0x000000007fa928b8; epoch: 0; age: 0) 8 4 (object header: class) 0x00001000 12 4 (object alignment gap) Instance size: 16 bytes Space losses: 0 bytes internal + 4 bytes external = 4 bytes total ``` 中间加了sleep之后(sleep时间相对要长一些),t2获取到的是轻量级锁: ``` Before any lock: # WARNING: Unable to get Instrumentation. Dynamic Attach failed. You may add this JAR as -javaagent manually, or supply -Djdk.attach.allowAttachSelf java.lang.Object object internals: OFF SZ TYPE DESCRIPTION VALUE 0 8 (object header: mark) 0x0000000000000005 (biasable; age: 0) 8 4 (object header: class) 0x00001000 12 4 (object alignment gap) Instance size: 16 bytes Space losses: 0 bytes internal + 4 bytes external = 4 bytes total T1 ID: 27 T1 获得锁后: java.lang.Object object internals: OFF SZ TYPE DESCRIPTION VALUE 0 8 (object header: mark) 0x000001487e9b1005 (biased: 0x00000000521fa6c4; epoch: 0; age: 0) 8 4 (object header: class) 0x00001000 12 4 (object alignment gap) Instance size: 16 bytes Space losses: 0 bytes internal + 4 bytes external = 4 bytes total T2 ID: 28 T2 获得锁后: java.lang.Object object internals: OFF SZ TYPE DESCRIPTION VALUE 0 8 (object header: mark) 0x00000011a29ff080 (thin lock: 0x00000011a29ff080) 8 4 (object header: class) 0x00001000 12 4 (object alignment gap) Instance size: 16 bytes Space losses: 0 bytes internal + 4 bytes external = 4 bytes total ``` ### 困惑 很奇怪啊,就在写这个贴子的时候,中间没有sleep,有两次t2获取到的也是轻量级锁。所以偏向锁升级为轻量级锁的过程究竟是怎样的。
【学习笔记】并发编程
#学习笔记# #Java# #JUC# 黑马并发编程学习笔记(2万多字,进行中) 视频地址:https://www.bilibili.com/video/BV16J411h7Rd/?spm_id_from=333.337.search-card.all.click&vd_source=a835ff13776aa85a80bbdcf7eec57f27 同时宣传一下自己的博客,欢迎大家关注:https://blog.csdn.net/weixin_43811294?spm=1011.2124.3001.5343
