Java并发--线程计数器

2024-06-10 18:08
文章标签 java 线程 并发 计数器

本文主要是介绍Java并发--线程计数器,希望对大家解决编程问题提供一定的参考价值,需要的开发者们随着小编来一起学习吧!

Java中经常存在以下的需求,启动多个相同或者不同的线程,主线程需要等待所有的线程执行完才继续往下执行

要实现上面的需求,基本的思路: 创建一个计数器, 来记录线程的执行

有两种实现方法

方法1:

使用锁和计数器:需要有一个对象锁,作用一:保证这个计数器的线程安全,作用二:阻塞主线程,等待所有线程执行完再来唤醒主线程继续执行

方法2:

使用Java线程包中的CountDownLatch:不需要加锁, 不需要wait notify这么复杂

方法1:

package com.yaya.thread.threadCount.count;import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;public class TestCount {int count = 0;int total = 0;Object object = new Object();void testCount() {total = 10;ExecutorService pool = Executors.newFixedThreadPool(5);for (int i = 0; i < 10; i++) {final int j = i;Runnable runnable = new Runnable() {@Overridepublic void run() {// TODO Auto-generated method stubtry {System.out.println("runnalbe:" + j);Thread.sleep(2000);} catch (InterruptedException e) {// TODO Auto-generated catch blocke.printStackTrace();} finally {synchronized (object) {count++;if (count == total) {object.notify();}}}}};pool.execute(runnable);}synchronized (object) {if (count != total) {try {object.wait();} catch (InterruptedException e) {// TODO Auto-generated catch blocke.printStackTrace();}}}System.out.println("end");pool.shutdown();}public static void main(String[] args) {TestCount testCount = new TestCount();testCount.testCount();}}

方法2: CountDownLatch

package com.yaya.thread.threadCount.countDownLatch;import java.util.ArrayList;
import java.util.List;
import java.util.concurrent.Callable;
import java.util.concurrent.CountDownLatch;import com.yaya.thread.future.ThreadPoolUtil;public class TestCountDownLatch {static CountDownLatch count = null;public void testRunnalbe(){List<String> list = new ArrayList<String>();for (int i = 1; i <= 10; i++) {list.add("list" + i);}count = new CountDownLatch(list.size());List<Runnable> runnables = new ArrayList<Runnable>();for (int i = 0; i < list.size(); i++) {final String listName = list.get(i);Runnable runnable = new TestRunnable1(listName, count);runnables.add(runnable);}try {ThreadPoolUtil.exeRunnableList(runnables);} catch (Exception e) {System.err.println(e.getMessage());}try {count.await();System.out.println("end");ThreadPoolUtil.shutDown(true);} catch (InterruptedException e) {// TODO Auto-generated catch blocke.printStackTrace();}}public void testDifferentRunnalbes(){List<String> list = new ArrayList<String>();for (int i = 1; i <= 10; i++) {list.add("list" + i);}count = new CountDownLatch(list.size()*2);List<Runnable> runnables = new ArrayList<Runnable>();for (int i = 0; i < list.size(); i++) {final String listName = list.get(i);Runnable runnable = new TestRunnable1(listName, count);runnables.add(runnable);}for (int i = 0; i < list.size(); i++) {final String listName = list.get(i);Runnable runnable = new TestRunnable2(listName, count);runnables.add(runnable);}try {ThreadPoolUtil.exeRunnableList(runnables);} catch (Exception e) {System.err.println(e.getMessage());}try {count.await();System.out.println("end");ThreadPoolUtil.shutDown(true);} catch (InterruptedException e) {// TODO Auto-generated catch blocke.printStackTrace();}}public static void main(String[] args) {TestCountDownLatch  testCountDownLatch  = new TestCountDownLatch();testCountDownLatch.testDifferentRunnalbes();}}
package com.yaya.thread.threadCount.countDownLatch;import java.util.concurrent.CountDownLatch;public class TestRunnable2 implements Runnable {String appid;CountDownLatch count;public TestRunnable2(String appid, CountDownLatch count) {super();this.appid = appid;this.count = count;}@Overridepublic void run() {System.out.println("task" + this.appid + "开始");try {Thread.sleep(3000);} catch (InterruptedException e) {e.printStackTrace();}System.out.println("task" + this.appid + "睡了3s");count.countDown();}}

package com.yaya.thread.future;import java.util.List;
import java.util.concurrent.Callable;
import java.util.concurrent.ExecutionException;
import java.util.concurrent.Future;
import java.util.concurrent.LinkedBlockingQueue;
import java.util.concurrent.ThreadFactory;
import java.util.concurrent.ThreadPoolExecutor;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.TimeoutException;
import java.util.concurrent.atomic.AtomicInteger;public class ThreadPoolUtil {// 线程池private static ThreadPoolExecutor threadPool;// 线程池核心线程数private static final int CORE_POOL_SIZE = 5;// 线程池最大线程数private static final int MAX_POOL_SIZE = 10;// 额外线程空状态生存时间private static final int KEEP_ALIVE_TIME = 10000;private static final int CANCEL_TASK_TIME = 20;private ThreadPoolUtil() {}static {threadPool = new ThreadPoolExecutor(CORE_POOL_SIZE, MAX_POOL_SIZE, KEEP_ALIVE_TIME, TimeUnit.SECONDS,new LinkedBlockingQueue<>(30), new ThreadFactory() {private final AtomicInteger integer = new AtomicInteger();@Overridepublic Thread newThread(Runnable r) {return new Thread(r, "mock thread:" + integer.getAndIncrement());}});}/*** 从线程池中抽取线程,执行指定的Runnable对象* * @param runnable*/public static void execute(Runnable runnable) {threadPool.execute(runnable);}/*** 批量执行 Runnable任务* * @param runnableList*/public static void exeRunnableList(List<Runnable> runnableList) {for (Runnable runnable : runnableList) {threadPool.execute(runnable);}}/*** 从线程池中抽取线程,执行指定的Callable对象* * @param callable* @return 返回执行完毕后的预期结果*/public static Future exeCallable(Callable<String> callable) {return threadPool.submit(callable);}/*** 批量执行 Callable任务* * @param callableList*            callable的实例列表* @return 返回指定的预期执行结果*/public static List<Future<String>> exeCallableList(List<Callable<String>> callableList) {List<Future<String>> futures = null;try {for (Callable<String> task : callableList) {threadPool.submit(task);}futures = threadPool.invokeAll(callableList);} catch (InterruptedException e) {e.printStackTrace();}return futures;}/*** 批量执行 Callable任务, 但不等待执行完* * @param callableList*            callable的实例列表* @return 返回指定的预期执行结果*/public static void exeCallableListNoReturn(List<Callable<String>> callableList) {for (Callable<String> task : callableList) {threadPool.submit(task);}}/*** 中断任务的执行* * @param isForceClose*            true:强制中断 false:等待任务执行完毕后,关闭线程池*/public static void shutDown(boolean isForceClose) {if (isForceClose) {threadPool.shutdownNow();} else {threadPool.shutdown();}}/*** 若超出CANCEL_TASK_TIME的时间,没有得到执行结果,则尝试中断线程* * @param future* @return 中断成功,则返回true 否则返回false*/public static boolean attemptCancelTask(Future future) {boolean cancel = false;try {future.get(CANCEL_TASK_TIME, TimeUnit.MINUTES);cancel = true;} catch (InterruptedException e) {e.printStackTrace();} catch (ExecutionException e) {e.printStackTrace();} catch (TimeoutException e) {cancel = future.cancel(true);e.printStackTrace();}return cancel;}}


这篇关于Java并发--线程计数器的文章就介绍到这儿,希望我们推荐的文章对编程师们有所帮助!



http://www.chinasem.cn/article/1048882

相关文章

springboot集成easypoi导出word换行处理过程

《springboot集成easypoi导出word换行处理过程》SpringBoot集成Easypoi导出Word时,换行符n失效显示为空格,解决方法包括生成段落或替换模板中n为回车,同时需确... 目录项目场景问题描述解决方案第一种:生成段落的方式第二种:替换模板的情况,换行符替换成回车总结项目场景s

SpringBoot集成redisson实现延时队列教程

《SpringBoot集成redisson实现延时队列教程》文章介绍了使用Redisson实现延迟队列的完整步骤,包括依赖导入、Redis配置、工具类封装、业务枚举定义、执行器实现、Bean创建、消费... 目录1、先给项目导入Redisson依赖2、配置redis3、创建 RedissonConfig 配

SpringBoot中@Value注入静态变量方式

《SpringBoot中@Value注入静态变量方式》SpringBoot中静态变量无法直接用@Value注入,需通过setter方法,@Value(${})从属性文件获取值,@Value(#{})用... 目录项目场景解决方案注解说明1、@Value("${}")使用示例2、@Value("#{}"php

SpringBoot分段处理List集合多线程批量插入数据方式

《SpringBoot分段处理List集合多线程批量插入数据方式》文章介绍如何处理大数据量List批量插入数据库的优化方案:通过拆分List并分配独立线程处理,结合Spring线程池与异步方法提升效率... 目录项目场景解决方案1.实体类2.Mapper3.spring容器注入线程池bejsan对象4.创建

线上Java OOM问题定位与解决方案超详细解析

《线上JavaOOM问题定位与解决方案超详细解析》OOM是JVM抛出的错误,表示内存分配失败,:本文主要介绍线上JavaOOM问题定位与解决方案的相关资料,文中通过代码介绍的非常详细,需要的朋... 目录一、OOM问题核心认知1.1 OOM定义与技术定位1.2 OOM常见类型及技术特征二、OOM问题定位工具

基于 Cursor 开发 Spring Boot 项目详细攻略

《基于Cursor开发SpringBoot项目详细攻略》Cursor是集成GPT4、Claude3.5等LLM的VSCode类AI编程工具,支持SpringBoot项目开发全流程,涵盖环境配... 目录cursor是什么?基于 Cursor 开发 Spring Boot 项目完整指南1. 环境准备2. 创建

Spring Security简介、使用与最佳实践

《SpringSecurity简介、使用与最佳实践》SpringSecurity是一个能够为基于Spring的企业应用系统提供声明式的安全访问控制解决方案的安全框架,本文给大家介绍SpringSec... 目录一、如何理解 Spring Security?—— 核心思想二、如何在 Java 项目中使用?——

SpringBoot+RustFS 实现文件切片极速上传的实例代码

《SpringBoot+RustFS实现文件切片极速上传的实例代码》本文介绍利用SpringBoot和RustFS构建高性能文件切片上传系统,实现大文件秒传、断点续传和分片上传等功能,具有一定的参考... 目录一、为什么选择 RustFS + SpringBoot?二、环境准备与部署2.1 安装 RustF

springboot中使用okhttp3的小结

《springboot中使用okhttp3的小结》OkHttp3是一个JavaHTTP客户端,可以处理各种请求类型,比如GET、POST、PUT等,并且支持高效的HTTP连接池、请求和响应缓存、以及异... 在 Spring Boot 项目中使用 OkHttp3 进行 HTTP 请求是一个高效且流行的方式。

java.sql.SQLTransientConnectionException连接超时异常原因及解决方案

《java.sql.SQLTransientConnectionException连接超时异常原因及解决方案》:本文主要介绍java.sql.SQLTransientConnectionExcep... 目录一、引言二、异常信息分析三、可能的原因3.1 连接池配置不合理3.2 数据库负载过高3.3 连接泄漏